Day 14 / 共 20 天 · 阶段4 RAG 知识库

检索 retrieval:一句提问,怎么从上万片里捞出最相关的几片

Day13 把文档切片、向量化、建好了索引。今天是在线的另一半:用户一句提问进来,怎么快、准地找出最相关的几个片段喂给大模型。主角是 core/rag/datasource/retrieval_service.py。弄清三件事:①Dify 支持哪四种检索方式(语义/全文/关键词/混合);②为什么要用线程池"多路并发召回";③混合检索为什么必须靠 rerank(重排)来融合,重排又分"模型重排"和"权重重排"两种。检索质量直接决定 RAG 答得准不准——这是 RAG 效果的第二个命门(第一个是 Day13 的切分)。

📍 你在 20 天里的位置(阶段4:RAG 知识库)
S2 模型运行时 S3 工作流引擎 D13 RAG 与索引 D14 检索 retrieval D15 知识库数据流 S5 工具/Agent S6 收官
💡 先用两个类比兜住今天 类比一:多路检索 + rerank 像公司招聘的"多渠道初筛 + HR 终面"。先让猎头(向量检索)、内推(全文检索)、招聘网站(关键词检索)各自并行捞一批候选人(召回快、宁滥勿缺);再由资深 HR(rerank 模型)把所有候选人放一起精挑细选、重新排序,选出真正最合适的前几名。初筛图快图全,终面图准。类比二:向量检索像按"气味"找东西(语义相近就行,换个说法也能找到),关键词检索像按"名字"找东西(必须字面命中,但精确)。混合检索就是两种鼻子一起用,再合并。
L01

痛点:光有向量检索还不够

🤔 痛点Day13 建好了向量库,是不是"把问题转成向量,找最近的 top-k"就完事了?现实没这么简单:①纯向量检索对"专有名词、编号、型号"不敏感——用户搜"错误码 E1024",向量可能觉得"E1023"也很像,但你要的是精确命中;②纯关键词检索又不懂同义——搜"如何请假"匹配不到写着"休假申请流程"的段落;③把两种结果简单拼在一起,分数体系不一样没法排序。怎么既要语义的"懂意思",又要关键词的"抠字眼",还能给出一个统一、可信的最终排序?
💡 本质:多路并发召回 + 统一重排Dify 的解法两步走:召回(recall)——同时用多种方式各捞一批候选(要快、要全,宁可多召回);重排(rerank)——把所有候选放进同一个"裁判"重新打分排序,取真正的 top-k。召回用线程池并发,因为向量查询、全文查询是彼此独立的 IO,串行等太慢。重排则是把"不同来源、不同分数体系"的结果拉到同一把尺子上——这是混合检索能work的关键。
L02

四种检索方式:一个枚举讲清

所有检索方式钉在 RetrievalMethodcore/rag/retrieval/retrieval_methods.py:4):

# core/rag/retrieval/retrieval_methods.py:4
class RetrievalMethod(StrEnum):
    SEMANTIC_SEARCH = "semantic_search"      # 语义:向量最近邻(懂意思)
    FULL_TEXT_SEARCH = "full_text_search"    # 全文:倒排索引(抠字眼)
    HYBRID_SEARCH = "hybrid_search"          # 混合:语义 + 全文,再 rerank
    KEYWORD_SEARCH = "keyword_search"        # 关键词:经济模式的字面匹配

    @staticmethod
    def is_support_semantic_search(retrieval_method: str) -> bool:
        return retrieval_method in {RetrievalMethod.SEMANTIC_SEARCH, RetrievalMethod.HYBRID_SEARCH}

    @staticmethod
    def is_support_fulltext_search(retrieval_method: str) -> bool:
        return retrieval_method in {RetrievalMethod.FULL_TEXT_SEARCH, RetrievalMethod.HYBRID_SEARCH}
两个 is_support_ 方法★注意 HYBRID_SEARCH 同时出现在两个集合里——这就是"混合"的定义:既支持语义、又支持全文。后面 _retrieve 就靠这两个判断决定"要不要提交向量检索任务、要不要提交全文检索任务"。
StrEnum又是字符串枚举——检索方式这个配置从数据库/前端传进来就是字符串,枚举保证只能是这四个之一。
大白话你在知识库设置里选的"检索设置",就对应这四个值。选"混合检索",Dify 就会同时跑语义和全文两路,再融合;选"向量检索"就只跑语义那一路。is_support_* 这两个小函数,就是"这次要开哪几路"的开关。
L03

retrieve:对外总入口,先并发编排

对外入口是 RetrievalService.retrievecore/rag/datasource/retrieval_service.py:114)。它本身也用线程池,把"文本查询"和"多个附件查询"并发发出去:

# core/rag/datasource/retrieval_service.py:114
def retrieve(cls, retrieval_method, dataset_id, query, top_k=4, score_threshold=0.0,
             reranking_model=None, reranking_mode="reranking_model", weights=None,
             document_ids_filter=None, attachment_ids=None):
    if not query and not attachment_ids:
        return []                                     # ① 空查询直接返回
    dataset = cls._get_dataset(dataset_id)
    if not dataset:
        return []
    all_documents: list[Document] = []
    exceptions: list[str] = []
    with ThreadPoolExecutor(max_workers=dify_config.RETRIEVAL_SERVICE_EXECUTORS) as executor:  # ② 线程池
        futures = []
        retrieval_service = RetrievalService()
        if query:
            futures.append(executor.submit(_propagate_otel_context(retrieval_service._retrieve),
                flask_app=current_app._get_current_object(), retrieval_method=retrieval_method,
                dataset=dataset, query=query, top_k=top_k, ..., all_documents=all_documents, exceptions=exceptions))
        if attachment_ids:
            for attachment_id in attachment_ids:      # 多个附件各起一个任务
                futures.append(executor.submit(..., attachment_id=attachment_id, ...))
        if futures:
            for _ in concurrent.futures.as_completed(futures, timeout=3600):
                if exceptions:                        # ③ 任何子任务出错就取消其余、尽早失败
                    for f in futures:
                        f.cancel()
                    break
    if exceptions:
        raise ValueError(";\n".join(exceptions))
    return all_documents
top_k=4默认只要最相关的 4 片。这个数字是"喂给模型的上下文预算"和"召回覆盖率"的平衡——太多会撑爆 Prompt、稀释重点。
_propagate_otel_context包一层是为了把 OpenTelemetry 的链路追踪上下文带进子线程——否则子线程里的 span 就和主链路断了,排查性能问题时看不到检索这一段。
all_documents 作参数传进去★注意结果不是 return 出来的,而是把同一个 list 传给每个子任务往里 append。多线程共享收集容器——所以 _retrieve 的签名里有 all_documentsexceptions 两个"输出参数"。
as_completed + cancel + break尽早失败:谁先出错,立刻取消其余任务并跳出,不傻等 timeout。检索是用户在等的实时操作,快速失败比拖着强。
L04

_retrieve:三路并发召回

真正干活的 _retrievecore/rag/datasource/retrieval_service.py:781)里,又开一层线程池,按检索方式把"关键词/语义/全文"三路并发提交:

# core/rag/datasource/retrieval_service.py:781
def _retrieve(self, flask_app, retrieval_method, dataset, all_documents, exceptions,
              query=None, top_k=4, score_threshold=0.0, reranking_model=None, ...):
    with flask_app.app_context():
        all_documents_item: list[Document] = []
        with ThreadPoolExecutor(max_workers=dify_config.RETRIEVAL_SERVICE_EXECUTORS) as executor:
            futures = []
            if retrieval_method == RetrievalMethod.KEYWORD_SEARCH and query:
                futures.append(executor.submit(self.keyword_search, ...))          # 关键词路
            if RetrievalMethod.is_support_semantic_search(retrieval_method):        # ★ 语义/混合 → 开
                if query:
                    futures.append(executor.submit(self.embedding_search, ..., query_type=QueryType.TEXT_QUERY))
                if attachment_id:
                    futures.append(executor.submit(self.embedding_search, ..., query_type=QueryType.IMAGE_QUERY))
            if RetrievalMethod.is_support_fulltext_search(retrieval_method) and query:   # ★ 全文/混合 → 开
                futures.append(executor.submit(self.full_text_index_search, ...))
            if futures:
                for future in concurrent.futures.as_completed(futures, timeout=300):
                    if future.exception():
                        for f in futures: f.cancel()
                        break
is_support_semantic_search★L02 那个开关在这儿用上了:只要方式是"语义"或"混合",就提交向量检索任务 embedding_search
is_support_fulltext_search同理:方式是"全文"或"混合",就提交 full_text_index_search。所以混合检索会同时提交语义 + 全文两个任务,并行跑。
embedding_search 分文本/图片向量检索还分文本查询和图片查询(多模态知识库)——同一个函数用 query_type 区分。
all_documents_item三路的结果先汇到这个局部 list,等下(L05)如果是混合检索,还要对它做去重 + rerank,再并回外层 all_documents
⚠️ 坑:召回阶段的 score_threshold 对混合检索要放宽源码里对混合检索特意把向量的早期阈值设成 0.0(见 retrieval_service.py:343 附近注释),并说明:混合检索的最终排序由 rerank/融合决定,如果在召回阶段就用用户的分数阈值卡,会用"未重排的原始分数"误杀掉本该高质量的片段(对应 issue #35233)。这提醒你:召回阶段要宽(宁多勿漏),过滤留到重排之后再做
L05

混合检索:去重 + rerank 融合

三路召回汇齐后,混合检索要做"去重 + 重排"这一步(core/rag/datasource/retrieval_service.py:881):

# core/rag/datasource/retrieval_service.py:881
if retrieval_method == RetrievalMethod.HYBRID_SEARCH:
    all_documents_item = self._deduplicate_documents(all_documents_item)     # ① 去重(同一片可能被多路命中)
    data_post_processor = DataPostProcessor(
        str(dataset.tenant_id), reranking_mode, reranking_model, weights, False)  # ② 造重排器
    if query:
        rerank_query, query_type = query, QueryType.TEXT_QUERY
    elif attachment_id:
        rerank_query, query_type = attachment_id, QueryType.IMAGE_QUERY
    else:
        return
    all_documents_item = data_post_processor.invoke(                         # ③ ★重排:统一打分排序取 top_k
        query=rerank_query, documents=all_documents_item,
        score_threshold=score_threshold, top_n=top_k, query_type=query_type)
    if not data_post_processor.rerank_runner and score_threshold:            # ④ 没重排器才在这兜底过滤阈值
        all_documents_item = self._filter_documents_by_vector_score_threshold(
            all_documents_item, score_threshold)
all_documents.extend(all_documents_item)                                     # 并回总结果
_deduplicate_documents★同一个片段可能被向量路和全文路召回。不去重的话,最终结果里会出现重复片段,既占 top-k 名额又误导模型。按 doc_id 去重。
DataPostProcessor.invoke★核心:把去重后的所有候选送进重排器,用同一把尺子重新打分,排序后取 top_n=top_k。这一步把"不同来源、分数不可比"的问题彻底解决——因为最终分数全来自同一个 rerank。
阈值过滤的位置★呼应 L04 的坑:只有在没有重排器时,才在这里用 score_threshold 过滤(拿向量分数)。有重排器时,过滤已经在 invoke 里用重排分数做了——顺序对了才不误杀。
L06

两种重排:模型重排 vs 权重重排

DataPostProcessor 在构造时就按 reranking_mode 选好了哪种重排器(core/rag/data_post_processor/data_post_processor.py:65):

# core/rag/data_post_processor/data_post_processor.py:65
def _get_rerank_runner(self, reranking_mode, tenant_id, reranking_model=None, weights=None):
    if reranking_mode == RerankMode.WEIGHTED_SCORE and weights:        # ① 权重重排
        return RerankRunnerFactory.create_rerank_runner(
            runner_type=reranking_mode, tenant_id=tenant_id,
            weights=Weights(vector_setting=VectorSetting(...), keyword_setting=KeywordSetting(...)))
    elif reranking_mode == RerankMode.RERANKING_MODEL:                 # ② 模型重排
        rerank_model_instance = self._get_rerank_model_instance(tenant_id, reranking_model)
        if rerank_model_instance is None:
            return None
        return RerankRunnerFactory.create_rerank_runner(
            runner_type=reranking_mode, rerank_model_instance=rerank_model_instance)
    return None

权重重排不调模型,纯靠"关键词分 × 权重 + 向量分 × 权重"算加权分。看 WeightRerankRunner.runcore/rag/rerank/weight_rerank.py:59)的核心公式:

# core/rag/rerank/weight_rerank.py:59
query_scores = self._calculate_keyword_score(query, documents)            # 关键词相似度(TF-IDF)
query_vector_scores = self._calculate_cosine(self.tenant_id, query, documents, self.weights.vector_setting)  # 向量余弦相似度
rerank_documents = []
for document, query_score, query_vector_score in zip(documents, query_scores, query_vector_scores):
    score = (self.weights.vector_setting.vector_weight * query_vector_score      # ★ 加权融合
             + self.weights.keyword_setting.keyword_weight * query_score)
    if score_threshold and score < score_threshold:
        continue
    document.metadata["score"] = score
    rerank_documents.append(document)
模型重排(RERANKING_MODEL)★用一个专门的 rerank 模型(如 bge-reranker、Cohere rerank)给"问题 + 每个候选片段"打相关性分。最准,但要额外调一次模型(花钱、增延迟)。
权重重排(WEIGHTED_SCORE)★不调模型,本地算:关键词分(TF-IDF)和向量分(余弦相似度)各乘一个权重相加。免费、快,效果不如模型重排但可调。你在界面拖的那个"语义 0.7 / 关键词 0.3"滑块,就是这里的 vector_weight / keyword_weight
_calculate_cosine余弦相似度:把问题和片段的向量做点积除以模长,值越接近 1 越相似。这是"向量接近 = 语义接近"的具体数学实现。
混合检索:多路并发召回 → 去重 → 重排 → top-k 语义(向量)检索 全文检索 关键词检索 并发召回(宁多勿漏) 去重 dedup rerank 重排 模型 / 权重 统一打分 top-k 给模型 权重重排:score = 向量分×w1 + 关键词分×w2 模型重排:rerank 模型直接给"问题+片段"相关性打分(最准)
图注:三路并发召回图快图全,去重防重复,rerank 把所有候选拉到同一把尺子上排序,取 top-k。
L07

串起来 + 今日小结

重排器统一入口 DataPostProcessor.invokecore/rag/data_post_processor/data_post_processor.py:49)——重排完还能接一个 reorder(把最相关的放两端,缓解"中间遗忘"):

# core/rag/data_post_processor/data_post_processor.py:49
def invoke(self, query, documents, score_threshold=None, top_n=None, query_type=QueryType.TEXT_QUERY):
    if self.rerank_runner:
        documents = self.rerank_runner.run(query, documents, score_threshold, top_n, query_type)   # 重排
    if self.reorder_runner:
        documents = self.reorder_runner.run(documents)     # 可选:重新摆放(头尾放最相关的)
    return documents
📝 真实值:一次混合检索的完整旅程 知识库设了"混合检索 + 权重重排(语义 0.7 / 关键词 0.3),top_k=4" → 用户问「新员工试用期多久」→ retrieve() 起线程 → _retrieveis_support_semanticis_support_fulltext 都为真,同时提交向量检索(召回讲"试用期/转正/考核期"的片段)和全文检索(召回字面含"试用期"的片段)→ 两路各回 10 片,合并成 20 片 → _deduplicate 去掉重叠的 5 片剩 15 片 → DataPostProcessor.invoke 用权重公式给 15 片重新打分:score = 0.7×余弦 + 0.3×TF-IDF → 排序取前 4 → 返回给上层,拼进 Prompt:"根据以下资料回答……" → 大模型答"试用期为 3 个月"。用户感觉"它真的读了我们手册",背后是并发召回 + 去重 + 加权重排这条链。

👶 小白:既然模型重排最准,为什么还要留个权重重排?

👨‍🏫 老师:三个字——成本、延迟、依赖。模型重排每次检索都要额外调一次 rerank 模型:要么花钱调云端 API、要么自己部署一个模型服务,还多一截延迟。对于量大、对成本敏感、或者没有 rerank 模型可用的场景,权重重排"本地算加权分"就够用了,而且那个"语义/关键词"权重滑块还能手动调,可解释性强。准确度、成本、可控性之间做权衡——这正是工程和刷榜的区别。选哪种,看你的业务能接受什么。

🧠 今天你应该能回答

  • Dify 有哪四种检索方式?(语义/全文/关键词/混合)
  • 混合检索为什么同时开语义和全文?(is_support_* 两个开关都命中)
  • 为什么召回要用线程池并发?(各路检索是独立 IO,串行太慢)
  • 混合检索为什么必须 rerank?(多路结果分数不可比,要统一尺子排序)
  • 为什么要先去重?(同一片可能被多路召回,重复占名额)
  • 模型重排 vs 权重重排怎么选?(准 vs 成本/延迟/可控,看业务)

✋ 10 分钟动手

cd /Users/bitmart/work/codes/github/AI_WORK/dify/api

# 1. 四种方式
cat core/rag/retrieval/retrieval_methods.py

# 2. 并发召回
sed -n '114,190p' core/rag/datasource/retrieval_service.py   # retrieve 入口
sed -n '781,875p' core/rag/datasource/retrieval_service.py   # _retrieve 三路并发
sed -n '881,909p' core/rag/datasource/retrieval_service.py   # 混合:去重 + rerank

# 3. 两种重排
sed -n '49,97p' core/rag/data_post_processor/data_post_processor.py  # invoke + 选 runner
sed -n '59,73p' core/rag/rerank/weight_rerank.py                     # 加权融合公式
明日预告 · Day 15:这两天我们从"文档"讲到"片段"再讲到"检索命中"。明天收口——知识库数据流Dataset / Document / DocumentSegment 这几张表在数据库里怎么串起来,以及 core/rag 里那个统一的 Document 数据结构,把整个 RAG 的"数据骨架"拼完整。
← Day 13 RAG 与索引 Day 15 · 知识库数据流 →