Skip to content

阶段五:查询流程 — 融合与输出(对应课程 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 state

5.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 state

5.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('✅ 端到端查询通过')
"

OPC 超级个体实战指南