Skip to content

阶段四:查询流程 — 三路检索(对应课程 day11~16)

当前状态:导入流程完成(8 节点),数据已入库 → 开始建设查询流程

知识标记总览

知识点出现次数扮演的角色
LangGraph 条件边🟢 第 2 次(Ch18→Ch19)三路并行检索的路由与合并
Milvus 向量检索🔵 第 3 次稠密+稀疏混合检索
HyDE 检索策略🔴 第 1 次假设性文档增强检索
Neo4J 图谱查询🔴 第 2 次(上阶段已学语法)通过 Cypher 查询关联实体
MCP 检索🟢 第 3 次(Ch16→Ch17→Ch19)外部实时数据检索
MongoDB 对话历史🔴 第 1 次多轮对话上下文管理

4.1 查询流程全景

节点职责知识标记
item_name_confirm_node从用户问题中提取商品名称,与 Milvus item_name 集合匹配🟢 LLM 第 2 次
vector_search_node将用户查询向量化后在 Milvus chunks 集合中混合检索🔵 Milvus 第 3 次
kg_search_nodeCypher 查询关联实体🟢🔴 Neo4J 第 2 次
mcp_search_nodeMCP 工具调用获取实时数据🟢 MCP 第 3 次

4.2 核心代码骨架

① 查询流程 State 定义

python
# → knowledge/processor/query_process/state.py
from typing import TypedDict, List

class QueryGraphState(TypedDict, total=False):
    # 输入
    query: str                    # 用户原始提问
    item_name: str                # 识别出的商品名
    history: list                 # 历史对话(MongoDB 加载)
    
    # 三路检索结果
    vector_results: list          # Milvus 返回的文档片段
    kg_results: list              # Neo4J 返回的关联实体
    mcp_results: list             # MCP 返回的外部数据
    
    # 融合后
    fused_results: list           # RRF 融合后的结果
    reranked_results: list        # Reranker 重排后的结果
    
    # 输出
    answer: str                   # LLM 最终生成的回答
    answer_sources: list          # 引用来源

② 商品名确认节点

python
class ItemNameConfirmNode(BaseNode):
    """从用户问题中确认商品名称"""
    name = "item_name_confirm"
    
    def process(self, state: QueryGraphState) -> QueryGraphState:
        query = state['query']
        milvus_util = MilvusUtil(self.config)
        embed_util = BGEM3EmbeddingUtil(self.config)
        
        # 1. 将用户查询向量化 → 在 item_name 集合中搜索
        query_vec = embed_util.encode(query)
        results = milvus_util.search("item_name", query_vec['dense_vector'], limit=3)
        
        # 2. 取最匹配的商品名
        if results:
            state['item_name'] = results[0]['entity']['text']
        return state

③ Milvus 混合检索引擎

python
class VectorSearchNode(BaseNode):
    """Milvus 稠密+稀疏混合检索"""
    name = "vector_search"
    
    def process(self, state: QueryGraphState) -> QueryGraphState:
        query = state['query']
        milvus_util = MilvusUtil(self.config)
        embed_util = BGEM3EmbeddingUtil(self.config)
        
        # 1. 计算双向量
        query_vectors = embed_util.encode(query)
        
        # 2. 混合检索(稠密 + 稀疏 + RRF)
        from pymilvus import AnnSearchRequest, RRFRanker
        
        dense_req = AnnSearchRequest(
            data=[query_vectors['dense_vector']],
            anns_field="vector", param={}, limit=10
        )
        sparse_req = AnnSearchRequest(
            data=[query_vectors['sparse_vector']],
            anns_field="sparse_vector", param={}, limit=10
        )
        
        results = milvus_util.client.hybrid_search(
            collection_name="chunks",
            reqs=[dense_req, sparse_req],
            ranker=RRFRanker(),
            limit=10,
            output_fields=["text", "title_path"]
        )
        
        state['vector_results'] = [
            {"text": r['entity']['text'], "title": r['entity']['title_path'], "source": "vector"}
            for r in results[0]
        ]
        return state

④ Neo4J 图谱查询节点

python
class KgSearchNode(BaseNode):
    """Neo4J 知识图谱查询"""
    name = "kg_search"
    
    def process(self, state: QueryGraphState) -> QueryGraphState:
        item_name = state.get('item_name', '')
        if not item_name:
            return state
        
        neo4j_util = Neo4jUtil(self.config)
        
        # Cypher 查询:找到目标商品及所有关联实体
        cypher = """
        MATCH (p:Product {name: $name})
        OPTIONAL MATCH (p)-[r]-(related)
        RETURN p.name as product, 
               type(r) as relation, 
               related.name as related_name,
               labels(related) as related_type
        """
        results = neo4j_util.query(cypher, {"name": item_name})
        
        state['kg_results'] = [
            {
                "product": r['product'],
                "relation": r['relation'],
                "related_name": r['related_name'],
                "related_type": r['related_type'],
                "source": "kg"
            }
            for r in results
        ]
        return state

⑤ 查询流程 main_graph

python
# → knowledge/processor/query_process/main_graph.py
from langgraph.graph import StateGraph
from langgraph.constants import START, END
from langgraph.checkpoint.memory import InMemorySaver

def create_query_graph():
    builder = StateGraph(QueryGraphState)
    
    # 添加节点
    builder.add_node("item_name_confirm", ItemNameConfirmNode())
    builder.add_node("vector_search", VectorSearchNode())
    builder.add_node("kg_search", KgSearchNode())
    builder.add_node("mcp_search", McpSearchNode())
    builder.add_node("rrf_node", RRFNode())
    builder.add_node("rerank_node", RerankNode())
    builder.add_node("answer_output", AnswerOutputNode())
    
    # 连接边
    builder.add_edge(START, "item_name_confirm")
    
    # 三路并行
    builder.add_edge("item_name_confirm", "vector_search")
    builder.add_edge("item_name_confirm", "kg_search")
    builder.add_edge("item_name_confirm", "mcp_search")
    
    # 汇聚
    builder.add_edge("vector_search", "rrf_node")
    builder.add_edge("kg_search", "rrf_node")
    builder.add_edge("mcp_search", "rrf_node")
    builder.add_edge("rrf_node", "rerank_node")
    builder.add_edge("rerank_node", "answer_output")
    builder.add_edge("answer_output", END)
    
    return builder.compile(checkpointer=InMemorySaver())

4.3 设计决策

决策选项理由
三路并行检索并行 / 串行三路互不依赖,并行降低延迟
商品名确认用 Milvus 搜索而非 LLMMilvus / LLM避免每次查询都调 LLM,减轻延迟和成本
HyDE 检索使用 / 不用HyDE 先生成假设文档再检索,对短查询提升显著

📂 对应的原始代码快照day11/day16/ 之间的逐日增量


4.4 验证命令

bash
# 验证三路检索正常
python -c "
from knowledge.processor.query_process.main_graph import create_query_graph
graph = create_query_graph()
result = graph.invoke({'query': '万用表使用方法'})
print('✅ 向量结果:', len(result.get('vector_results', [])))
print('✅ 图谱结果:', len(result.get('kg_results', [])))
"

OPC 超级个体实战指南