检索 retrieval:一句提问,怎么从上万片里捞出最相关的几片
Day13 把文档切片、向量化、建好了索引。今天是在线的另一半:用户一句提问进来,怎么快、准地找出最相关的几个片段喂给大模型。主角是 core/rag/datasource/retrieval_service.py。弄清三件事:①Dify 支持哪四种检索方式(语义/全文/关键词/混合);②为什么要用线程池"多路并发召回";③混合检索为什么必须靠 rerank(重排)来融合,重排又分"模型重排"和"权重重排"两种。检索质量直接决定 RAG 答得准不准——这是 RAG 效果的第二个命门(第一个是 Day13 的切分)。
痛点:光有向量检索还不够
四种检索方式:一个枚举讲清
所有检索方式钉在 RetrievalMethod(core/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又是字符串枚举——检索方式这个配置从数据库/前端传进来就是字符串,枚举保证只能是这四个之一。is_support_* 这两个小函数,就是"这次要开哪几路"的开关。retrieve:对外总入口,先并发编排
对外入口是 RetrievalService.retrieve(core/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_documents、exceptions 两个"输出参数"。as_completed + cancel + break★尽早失败:谁先出错,立刻取消其余任务并跳出,不傻等 timeout。检索是用户在等的实时操作,快速失败比拖着强。_retrieve:三路并发召回
真正干活的 _retrieve(core/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。0.0(见 retrieval_service.py:343 附近注释),并说明:混合检索的最终排序由 rerank/融合决定,如果在召回阶段就用用户的分数阈值卡,会用"未重排的原始分数"误杀掉本该高质量的片段(对应 issue #35233)。这提醒你:召回阶段要宽(宁多勿漏),过滤留到重排之后再做。混合检索:去重 + 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 里用重排分数做了——顺序对了才不误杀。两种重排:模型重排 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.run(core/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 越相似。这是"向量接近 = 语义接近"的具体数学实现。串起来 + 今日小结
重排器统一入口 DataPostProcessor.invoke(core/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
retrieve() 起线程 → _retrieve 里 is_support_semantic 和 is_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 # 加权融合公式
Dataset / Document / DocumentSegment 这几张表在数据库里怎么串起来,以及 core/rag 里那个统一的 Document 数据结构,把整个 RAG 的"数据骨架"拼完整。