阶段四:查询流程 — 三路检索(对应课程 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_node | Cypher 查询关联实体 | 🟢🔴 Neo4J 第 2 次 |
mcp_search_node | MCP 工具调用获取实时数据 | 🟢 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 搜索而非 LLM | Milvus / 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', [])))
"