阶段五:查询流程 — 融合与输出(对应课程 day17~18)
当前状态:三路检索已就绪 → 新增 RRF 融合 + Reranker 精排 + 答案生成
知识标记总览:
| 知识点 | 出现次数 | 扮演的角色 |
|---|---|---|
| RRF 融合算法 | 🔴 第 1 次(全新) | 将三路排序结果合并为一个有序列表 |
| Reranker 重排序 | 🔴 第 1 次(全新) | 用交叉熵模型对候选结果二次精排 |
| LLM 答案生成 | 🟢 第 2 次 | 基于融合后的上下文生成最终回答 |
| SSE 流式输出 | 🟢 第 2 次(Ch17→Ch19) | 逐 token 推送给前端 |
5.1 阶段起始状态
5.2 🔴 RRF 融合算法(全新知识)
原理
RRF(Reciprocal Rank Fusion)是一种无监督的结果融合算法——不依赖训练数据,仅通过排序位置来融合多个检索结果。
对于每篇文档 d,在 RRF 中的得分:
N
score(d) = Σ 1 / (k + rank_i(d))
i=1
其中:
N = 检索路数(本项目中 N=3:向量/图谱/MCP)
rank_i(d) = 文档 d 在第 i 路的排名(从 1 开始)
k = 平滑参数(通常 k=60)直觉理解:每路检索给每篇文档打一个"排名分"(第 1 名 1/61,第 2 名 1/62,以此类推),然后把三路分数相加。一篇文档如果在多路检索中都排在前面,总分就高。
python
def rrf_fuse(result_lists: list[list[dict]], k: int = 60) -> list[dict]:
"""
RRF 融合多路检索结果
Args:
result_lists: 每路检索的结果列表 [[doc1, doc2, ...], [doc1, doc2, ...], ...]
k: 平滑参数,越大各路之间的排名差异越不明显
Returns:
按 RRF 得分降序排列的文档列表
"""
scores = {}
for rank_list in result_lists:
for rank, doc in enumerate(rank_list):
doc_id = doc.get('id') or doc.get('text', '')[:50] # 唯一标识
scores[doc_id] = scores.get(doc_id, 0) + 1 / (k + rank + 1)
# 按得分降序排列
sorted_docs = sorted(scores.items(), key=lambda x: x[1], reverse=True)
return sorted_docs🔑 关键参数
k的作用:k 值越大,名次差异带来的分数差异越小(平滑化)。当 k=0 时只有第 1 名有分。当 k→∞ 时所有名次分数趋同。实践中 k=60 是经验值。
RRFNode 实现
python
class RRFNode(BaseNode):
"""RRF 融合三路检索结果"""
name = "rrf"
def process(self, state: QueryGraphState) -> QueryGraphState:
all_results = [
state.get('vector_results', []),
state.get('kg_results', []),
state.get('mcp_results', []),
]
# 过滤空结果
all_results = [r for r in all_results if r]
if not all_results:
return state
# 构建 {id: {text, source, rank}} 映射
doc_map = {}
result_id_lists = []
for rank_list in all_results:
ids = []
for rank, doc in enumerate(rank_list):
doc_id = doc.get('text', '')[:80]
doc_map[doc_id] = doc
ids.append(doc_id)
result_id_lists.append(ids)
# RRF 融合
k = 60
scores = {}
for ids in result_id_lists:
for rank, doc_id in enumerate(ids):
scores[doc_id] = scores.get(doc_id, 0) + 1 / (k + rank + 1)
# 排序
sorted_ids = sorted(scores.items(), key=lambda x: x[1], reverse=True)
state['fused_results'] = [
{**doc_map[doc_id], "rrf_score": round(score, 4)}
for doc_id, score in sorted_ids
]
return state5.3 🔴 Reranker 重排序(全新知识)
原理
RRF 融合后的结果仍然基于排序位置,不关心内容与问题的相关性。Reranker 使用**交叉编码器(Cross-Encoder)**对每条候选计算 (问题, 文档) 的相关性分数,重新排序。
差异对比:
Milvus 向量检索:将问题和文档分别编码 → 余弦相似度(双向量的距离)
Reranker:将(问题, 文档)拼在一起输入模型 → 直接输出相关性分数(更精确但更慢)| 对比 | Milvus 向量检索 | Reranker |
|---|---|---|
| 精度 | ⭐⭐ 中等(语义近似) | ⭐⭐⭐ 高(精确匹配) |
| 速度 | ⚡ 快(毫秒级) | 🐢 慢(每对几十毫秒) |
| 适用阶段 | 粗筛(从数万条中找到 Top-K) | 精排(从 Top-K 中排序) |
| 推荐策略 | 先用向量搜索粗筛 → 再用 Reranker 精排 |
RerankNode 实现
python
class RerankNode(BaseNode):
"""使用 bge-reranker 对 RRF 结果重排序"""
name = "rerank"
def process(self, state: QueryGraphState) -> QueryGraphState:
fused = state.get('fused_results', [])
if not fused:
return state
query = state['query']
# 使用 bge-reranker-v2 对 (query, text) 对打分
from FlagEmbedding import FlagReranker
reranker = FlagReranker("BAAI/bge-reranker-v2-m3")
pairs = [(query, doc.get('text', '')) for doc in fused]
scores = reranker.compute_score(pairs) # 每个 pair 返回一个分数
# 按重排序分数降序
scored = list(zip(fused, scores))
scored.sort(key=lambda x: x[1], reverse=True)
state['reranked_results'] = [
{**doc, "rerank_score": round(score, 4)}
for doc, score in scored
]
return state💡 Reranker 选型建议:bge-reranker-v2-m3 是开源最优选择。商业场景可使用 Cohere Rerank API。部署方式:Xinference 托管或 HuggingFace TEI。
5.4 🟢 答案生成节点(第 2 次)
python
class AnswerOutputNode(BaseNode):
"""基于重排序后的上下文生成最终答案"""
name = "answer_output"
def process(self, state: QueryGraphState) -> QueryGraphState:
query = state['query']
contexts = state.get('reranked_results', [])[:5] # 取 Top-5
# 组装上下文
context_text = "\n\n".join([
f"[{doc.get('source', 'unknown')}] {doc.get('text', '')}"
for doc in contexts
])
llm_client = LLMClient(self.config)
# 调用 LLM 生成答案
response = llm_client.chat([
{"role": "system", "content":
"你是一个电商知识助手。基于提供的参考资料回答问题。"
"如果参考资料中不包含答案,明确说'资料中未找到相关信息'。"
"引用来源时用 [vector]/[kg]/[mcp] 标注。"},
{"role": "user", "content": f"参考资料:\n{context_text}\n\n问题:{query}"}
])
state['answer'] = response
state['answer_sources'] = [
{"text": c['text'][:100], "score": c.get('rerank_score', 0), "source": c.get('source', '')}
for c in contexts
]
return state5.5 查询流程最终架构图
📂 对应的原始代码快照:
day17/→day18/的逐日增量
5.6 验证命令
bash
# 验证完整查询链路
python -c "
from knowledge.processor.query_process.main_graph import create_query_graph
graph = create_query_graph()
result = graph.invoke({'query': '万用表怎么用'})
print('答案:', result.get('answer', '')[:200])
print('来源:', len(result.get('answer_sources', [])))
print('✅ 端到端查询通过')
"