知识库查询 —— 知识图谱查询节点
本文档详细介绍知识图谱查询节点(query_kg)的设计与实现。该节点是整个查询流程中最复杂的节点,涉及 LLM 实体抽取、Milvus 向量对齐、Neo4j 图谱遍历、切片关联与文本回填等多个环节。
学习理念:知识图谱查询是多路并行检索中最复杂的节点——它协调 LLM(实体抽取)、Milvus(实体对齐)、Neo4j(种子节点 + 一跳扩展)、Milvus(文本回填)四个服务。核心价值:向量检索搜的是"文本相似",图谱查询搜的是"实体关系"——两者互补,RRF 融合后效果远胜单路。
海外对标:QueryKgNode 的"LLM 抽取 → 向量对齐 → 图遍历 → 文本回填"四阶段架构对标微软 GraphRAG 项目的 Global Search(实体抽取 → 社区发现 → 答案生成),以及 Neo4j 官方推荐的 "LLM + Graph + Vector" 混合检索模式。种子节点 + 一跳扩展的图遍历策略对标 Google Knowledge Graph API 的 Entity Search + Related Entities。
本节 AI 替代率:~65% | 人工干预率:~35%
| 角色 | 能力范围 |
|---|---|
| 🤖 AI 擅长 | LLM 实体抽取代码、向量对齐代码、Cypher 查询模板、文本回填代码、三元组转文本 |
| 👤 人类需理解 | 实体对齐阈值调优(KG_ENTITY_ALIGN_MIN_SCORE=0.7)、种子节点权重策略(seed=2.0 vs neighbor=1.0)、图遍历限流参数(per_entity/max_total 的取值) |
阅读指引
| 颜色 | 章节 | AI 替代率 | 人工干预 | 说明 |
|---|---|---|---|---|
| 🟡 | §1 任务目标 | ~95% | ~5% | 学习目标明确 |
| 🟢 | §2 核心概念 | ~90% | ~10% | 实体抽取 / 对齐 / 种子 / 一跳扩展 |
| 🟡 | §3 整体流程 | ~85% | ~15% | 理解 5 阶段数据流转 |
| 🔴 | §4 分步实现 | ~60% | ~40% | 最复杂的 8 步流程 |
| 🔴 | §4.4 主代码 | ~55% | ~45% | QueryKgNode 完整架构,8 个步骤的编排 |
| 🟢 | §5 测试运行 | ~95% | ~5% | 看预期输出的 5 步执行链 |
| 🟡 | §6 总结 | ~90% | ~10% | 设计要点回顾 |
技术栈健康度标签体系
| 技术 | 健康度 | 建议 |
|---|---|---|
| Neo4j | 🟢 稳定 | 图数据库第一品牌。Cypher 的精确/模糊匹配 + 双向一跳扩展是标准操作。 |
| Milvus 实体对齐 | 🔥 巅峰 | 用向量相似度做实体对齐。generate_hybrid_embeddings + execute_hybrid_search 标准流程。 |
| LLM 实体抽取 | ⏳ 成长期 | JSON Mode 下的结构化抽取。Prompt 设计直接影响抽取质量。2025-2026 年逐步被 Function Calling 替代。 |
| 权重排序 (w_seed/w_neighbor) | 🟢 稳定 | 种子节点权重 2.0 > 邻居权重 1.0 的设计简单有效。 |
体系说明:🟢🟡🟠🔴 标识学习优先级 / AI 替代率;🔥🟢⏳⚠️💀 标识技术栈健康度。
中英文对照表
| English | 中文 | 本质 |
|---|---|---|
| Entity Extraction | 实体抽取 | 从自然语言中识别出可能存在于知识图谱中的实体名称 |
| Entity Alignment | 实体对齐 | 通过向量相似度将抽取实体匹配到图谱中的标准实体名 |
| Seed Node | 种子节点 | 在 Neo4j 中匹配到的起始实体节点 |
| One-Hop Expansion | 一跳扩展 | 从种子节点出发遍历直接关联的邻居节点和关系 |
| MENTIONED_IN | 被提及于 | 实体与原文切片之间的关联关系 |
| Triple | 三元组 | (头实体, 关系, 尾实体) 的知识表示单元 |
| Weighted Sorting | 带权排序 | 种子节点权重大于邻居节点,多个实体关联的切片排名更高 |
| Text Backfill | 文本回填 | 根据切片 ID 从 Milvus 获取完整的切片内容 |
💡 程序员比喻
- QueryKgNode 就像
git log --all --oneline --graph——从 commit message(用户问题)中提取作者/文件(实体),再沿 parent 链(关系)扩展到相关 commit(三元组),最后git show看完整 diff(切片内容)。- 实体提取 → 对齐 → 图遍历 → 回填 就像
docker pull的镜像层下载——先查 manifest(实体抽取),再找 layer digest(对齐),下载依赖层(一跳扩展),最后解压到本地(文本回填)。- 种子权重 2.0 vs 邻居权重 1.0 就像
git log --author="me"vsgit log --committer="me"——直接匹配的(作者)权重更高,间接关联的(提交者)权重减半。
1. 任务目标
1.1 本章目标
通过本章学习,你将掌握:
- 理解知识图谱在 RAG 中的作用:掌握结构化知识与向量检索的结合方式
- 学会 LLM 实体抽取:从自然语言问题中提取关键实体
- 掌握实体对齐技术:将抽取的实体与图谱中的标准实体名对齐
- 理解图遍历策略:种子节点查找与一跳扩展的实现方式
- 学会多数据源协作:LLM + Milvus + Neo4j 的联合查询模式
- 实现完整的图谱查询节点:通过
if __name__ == "__main__"验证节点功能
1.2 涉及文件
knowledge/processor/query_process/
├── nodes/
│ └── query_kg.py # 知识图谱查询节点(本章重点)
├── prompt.py # 提示词模板(含 ENTITY_EXTRACT_SYSTEM_PROMPT)
└── ...
knowledge/tools/
├── llm_utils.py # LLM 客户端工具
├── embedding_utils.py # 向量嵌入工具
├── milvus_utils.py # Milvus 向量数据库工具
└── neo4j_utils.py # Neo4j 图数据库工具1.3 节点在流程中的位置

1.4 数据存储架构

2. 核心概念扫盲
2.1 知识图谱在 RAG 中的作用
| 检索方式 | 优点 | 缺点 | 适用场景 |
|---|---|---|---|
| 向量检索 | 语义理解强 | 缺乏结构化关系 | 内容相似度匹配 |
| 知识图谱 | 关系明确、可解释 | 覆盖范围有限 | 实体关系查询 |
| 混合使用 | 兼具两者优点 | 实现复杂度高 | 复杂问答场景 |
知识图谱的独特价值:
用户问题: "万用表的表笔怎么接?"
向量检索召回:
- "万用表使用说明..."
- "测量电压时..."
知识图谱额外提供:
- 万用表 -[HAS_PART]-> 表笔
- 表笔 -[CONNECTS_TO]-> 正极端口
- 表笔 -[CONNECTS_TO]-> 负极端口
- 红表笔 -[USED_FOR]-> 正极连接
- 黑表笔 -[USED_FOR]-> 负极连接
↓ 结合后
答案可以精确描述:表笔的组成、连接方式、颜色区分等结构化信息2.2 实体抽取
从自然语言问题中识别出可能存在于知识图谱中的实体:
用户问题: "万用表的电池怎么更换?"
↓
LLM 实体抽取
↓
抽取结果: ["电池", "更换"]实体抽取提示词:
🟡 【P1 看注释就行】 Step1 预处理——清理 item_names 兼容 None/str/list,从问题中移除商品名避免干扰实体抽取。
ENTITY_EXTRACT_SYSTEM_PROMPT = """你是一个知识图谱问答系统的"实体识别"模块。
请从用户问题中抽取用于查询图数据库(Neo4j)的实体名称
(优先抽取设备/部件/操作/工具/条件/警告等名词短语)。
要求:
1) 输出必须是严格 JSON
2) 只输出一个字段:entities
3) entities 为字符串数组,每个元素是实体 name
4) 只抽取与问题相关、且可能在图中作为节点 name 存在的实体;不要输出句子
输出示例:
{"entities":["电池安装","螺丝刀","表笔"]}"""2.3 实体对齐
将 LLM 抽取的实体名与知识图谱中的标准实体名对齐:
LLM 抽取: "电池更换"
↓
Milvus 向量搜索
↓
图谱中匹配: "电池安装" (相似度 0.92)
"电池仓盖" (相似度 0.85)
↓
对齐结果: "电池安装"为什么需要对齐?
- LLM 抽取的名称可能与图谱中的名称不完全一致
- 向量相似度匹配可以容忍用词差异
- 提高后续 Neo4j 查询的命中率
2.4 种子节点与一跳扩展

2.5 MENTIONED_IN 关系
知识图谱中的实体通过 MENTIONED_IN 关系与原文切片关联:
(:Entity {name: "电池安装"}) -[:MENTIONED_IN]-> (:Chunk {id: 12345})通过这个关系,可以找到与实体相关的原文切片,获取详细的上下文信息。
2.6 带权重的节点排序
🔥 【P0 必须要学】 Step2 LLM 实体抽取——使用
ENTITY_EXTRACT_SYSTEM_PROMPT约束 LLM 输出 JSON 数组。关键设计:enable_thinking=False(禁止推理链,直接输出 JSON),set去重。
# 种子节点权重更高
w_seed = 2.0 # 直接匹配的实体
# 邻居节点权重较低
w_neighbor = 1.0 # 一跳扩展的实体
# 切片得分 = 关联节点权重之和
# 多个高权重实体关联的切片排名更靠前3. 知识图谱查询业务处理流程(总)
3.1 整体流程图

3.2 数据流转详解

4. 知识图谱查询业务处理流程(分)
4.1 目标
实现一个能够:
- 从用户问题中智能抽取关键实体
- 将抽取实体与知识图谱中的标准实体对齐
- 在图数据库中找到相关节点和关系
- 将图谱信息转化为可用于答案生成的文本
- 回填原文切片,提供详细上下文
4.2 需求分析
4.2.1 功能需求
- 实体抽取:使用 LLM 从自然语言问题中提取实体
- 实体对齐:通过向量相似度将实体对齐到知识图谱
- 种子查找:在 Neo4j 中定位对应节点(精确 + 模糊匹配)
- 关系扩展:获取种子节点的一跳邻居和关系
- 切片关联:通过 MENTIONED_IN 关系找到相关文档切片
- 文本回填:从 Milvus 获取切片的完整文本内容
4.2.2 技术依赖
| 组件 | 用途 | 数据 |
|---|---|---|
| LLM | 实体抽取 | 用户问题 → 实体列表 |
| Milvus (entity) | 实体对齐 | 实体名 → 标准实体名 |
| Neo4j | 图谱查询 | 实体 → 种子节点 → 三元组 → 切片 ID |
| Milvus (chunks) | 文本回填 | 切片 ID → 切片内容 |
4.2.3 配置参数
🔥 【P0 必须要学】 Step3 Milvus 实体对齐——核心逻辑:每个实体生成混合向量 → 在
entity_name_collection中搜索 →_pick_best_hit选择最佳匹配(KG_ENTITY_ALIGN_MIN_SCORE=0.7阈值)。失败的实体(no_hit)不影响其他实体的对齐。
# 实体对齐
KG_ENTITY_ALIGN_MIN_SCORE = 0.7 # 对齐最低分数阈值(环境变量)
# 种子节点
per_entity = 3 # 每个实体最多匹配的种子数
max_total = 30 # 种子节点总数上限
# 一跳扩展
per_seed = 50 # 每个种子最多扩展的三元组数
max_total = 200 # 三元组总数上限
# 切片获取
max_chunks = 200 # 关联切片总数上限
# 节点权重
w_seed = 2.0 # 种子节点权重
w_neighbor = 1.0 # 邻居节点权重4.3 实现流程
4.3.1 实现流程图

4.3.2 具体实现步骤
Step 1: 预处理输入
目的: 清理输入数据,为后续处理做准备。
实现逻辑:
- 获取查询文本(优先 rewritten_query)
- 清理商品名称列表(处理 None/str/list)
- 从问题中移除商品名,避免干扰实体抽取
代码片段:
🟡 【P1 看注释就行】 Step4 Neo4j 种子查找——先精确匹配
n.name = $name,失败则模糊匹配CONTAINS toLower($name)。per_entity=3+max_total=30防止结果爆炸。
def process(self, state: QueryGraphState) -> QueryGraphState:
from knowledge.tools.llm_utils import get_llm_client
# 1. 获取查询文本
question = state.get("rewritten_query") or state.get("original_query", "")
# 2. 清理商品名称
item_names = self._clean_item_names(state.get("item_names"))
# 3. 从问题中移除商品名(避免干扰实体抽取)
for name in item_names:
question = question.replace(name, "")清理函数详解:
🟡 【P1 看注释就行】 Cypher 查询——精确匹配 + 模糊匹配的双阶段查询策略。注意
toLower()确保大小写不敏感。
@staticmethod
def _clean_item_names(item_names: Any) -> List[str]:
"""清理商品名称:兼容 None/str/list,返回去重列表。"""
if not item_names:
return []
if isinstance(item_names, str):
return [item_names.strip()] if item_names.strip() else []
seen: Set[str] = set()
return [
s for x in item_names
if (s := str(x).strip()) and s not in seen and not seen.add(s)
]Step 2: LLM 实体抽取
目的: 从用户问题中提取可能存在于知识图谱中的实体名称。
实现逻辑:
- 使用专门的实体抽取提示词
- 调用 LLM(JSON 模式)
- 解析 JSON 响应,提取 entities 字段
- 去重和清理
代码片段:
🟡 【P1 看注释就行】 Step5 一跳扩展——双向遍历(出边 + 入边),排除
MENTIONED_IN关系。per_seed=50+max_total=200限流。
def _extract_entities(self, question: str, llm_client) -> List[str]:
"""LLM 抽取实体并解析 JSON。"""
try:
# 1. 调用 LLM
resp = llm_client.invoke(
[SystemMessage(content=ENTITY_EXTRACT_SYSTEM_PROMPT),
HumanMessage(content=f"用户问题:{question}")],
extra_body={"enable_thinking": False},
)
# 2. 解析 JSON
data = json.loads((resp.content or "").strip())
# 3. 提取并去重
entities = list({
e.strip() for e in data.get("entities", []) if e.strip()
})
self.logger.info(f"抽取到 {len(entities)} 个实体: {entities}")
return entities
except Exception as e:
self.logger.error(f"实体抽取失败: {e}")
return []LLM 输入输出示例:
输入问题: "电池怎么更换?"
LLM 输出:
{"entities": ["电池", "更换"]}
解析结果: ["电池", "更换"]Step 3: Milvus 实体对齐
目的: 将 LLM 抽取的实体与知识图谱中的标准实体名对齐。
实现逻辑:
- 为每个实体生成混合向量
- 在 entity_collection 中执行混合搜索
- 根据相似度分数选择最佳匹配
- 低于阈值的视为未命中
- 记录对齐详情供调试
代码片段:
🟡 【P1 看注释就行】 Cypher 扩展——
CALL { MATCH ... UNION MATCH ... }的双向查询模式。注意排除MENTIONED_IN的 WHERE 条件。
def _align_entities(
self, entities: List[str], item_names: List[str], top_k: int = 5,
) -> Dict[str, Any]:
"""向量检索对齐实体名称。"""
from knowledge.tools.embedding_utils import generate_hybrid_embeddings
from knowledge.tools.milvus_utils import (
get_milvus_client, build_hybrid_search_requests, execute_hybrid_search,
)
collection = "kb_graph_entity_names"
client = get_milvus_client()
min_score = self._get_float_env("KG_ENTITY_ALIGN_MIN_SCORE")
expr = self._build_filter_expr(item_names) # 按商品过滤
# 1. 生成混合向量
try:
emb = generate_hybrid_embeddings(entities)
except Exception as e:
self.logger.error(f"Embedding 生成失败: {e}")
return {"aligned_entities": entities, "alignments": []}
alignments: List[Dict] = []
aligned: List[str] = []
seen: Set[str] = set()
# 2. 遍历每个实体进行对齐
for idx, entity in enumerate(entities):
dense = emb["dense"][idx]
sparse = emb["sparse"][idx]
try:
# 2.1 构建搜索请求
reqs = build_hybrid_search_requests(
dense_vector=dense,
sparse_vector=sparse,
filter_expr=expr,
top_k=top_k,
)
# 2.2 执行混合搜索
res = execute_hybrid_search(
client=client,
collection_name=collection,
search_requests=reqs,
ranker_weights=(0.5, 0.5),
output_fields=["entity_name", "item_name"],
)
# 2.3 选择最佳命中
best = self._pick_best_hit(res[0] if res else [], min_score)
if best:
name = best["entity"]["entity_name"]
if name not in seen:
seen.add(name)
aligned.append(name)
alignments.append({
"original": entity,
"aligned": name,
"score": best["distance"],
})
else:
alignments.append({
"original": entity,
"aligned": None,
"reason": "no_hit"
})
except Exception as e:
alignments.append({
"original": entity,
"aligned": None,
"reason": f"error:{e}"
})
return {"aligned_entities": aligned, "alignments": alignments}对齐结果示例:
🔥 【P0 必须要学】 Step6 切片关联 + 权重排序——
_build_weighted_nodes的核心设计:种子节点权重 2.0,邻居节点权重 1.0。UNWIND $nodes+sum(n.w)计算切片得分,ORDER BY score DESC排序。这个权重策略决定了哪些切片排在前面。
{
"aligned_entities": ["电池安装"],
"alignments": [
{"original": "电池", "aligned": "电池安装", "score": 0.92},
{"original": "更换", "aligned": None, "reason": "no_hit"}
]
}Step 4: Neo4j 种子节点查找
目的: 在 Neo4j 图数据库中找到对齐实体对应的节点。
实现逻辑:
- 使用事务读取(execute_read)
- 先尝试精确匹配(name = $name)
- 精确匹配失败则模糊匹配(CONTAINS)
- 限制每个实体的种子数和总数
代码片段:
🟡 【P1 看注释就行】 Step7 Milvus 文本回填——去重 chunk_id →
fetch_chunks_by_ids批量查询 → 按原顺序回填。注意row_map的映射构建。
def _find_seed_nodes(
self, entities: List[str], item_names: List[str],
per_entity: int = 3, max_total: int = 30,
) -> List[NodeRef]:
"""在 Neo4j 中查找种子节点(精确 → 模糊)。"""
if not entities or not item_names:
return []
try:
with self._neo4j_session() as session:
seeds: List[NodeRef] = []
seen: Set[Tuple[str, str]] = set()
for name in entities:
# 事务读取
rows = session.execute_read(
self._tx_find_seeds, name, item_names, per_entity
)
for s in rows:
key = (s["item_name"], s["name"])
if key not in seen:
seen.add(key)
seeds.append(s)
if len(seeds) >= max_total:
break
if len(seeds) >= max_total:
break
return seeds
except Exception as e:
self.logger.error(f"Neo4j 种子查询异常: {e}")
return []Cypher 查询详解:
🟡 【P1 看注释就行】 Step8 三元组转文本——
f"[{item_name}] {head} -({rel})-> {tail}"格式。set去重确保无重复。
@staticmethod
def _tx_find_seeds(tx, name: str, item_names: List[str], limit: int):
"""事务: 精确匹配 → 模糊匹配。"""
# 精确匹配
seeds = tx.run("""
MATCH (n:Entity)
WHERE n.name = $name AND n.item_name IN $item_names
RETURN n.name AS name, n.item_name AS item_name
LIMIT $limit
""", name=name, item_names=item_names, limit=limit).data()
if seeds:
return seeds
# 模糊匹配(精确匹配无结果时)
return tx.run("""
MATCH (n:Entity)
WHERE n.name IS NOT NULL
AND toLower(n.name) CONTAINS toLower($name)
AND n.item_name IN $item_names
RETURN n.name AS name, n.item_name AS item_name
LIMIT $limit
""", name=name, item_names=item_names, limit=limit).data()Step 5: 一跳扩展
目的: 获取种子节点的一跳邻居关系,形成三元组。
实现逻辑:
- 双向遍历(出边 + 入边)
- 排除 MENTIONED_IN 关系(这是切片引用关系)
- 限制每个种子的三元组数和总数
- 去重处理
代码片段:
🔥 【P0 必须要学】 QueryKgNode 完整架构——8 步流程:预处理 → LLM抽取 → Milvus对齐 → Neo4j种子 → 一跳扩展 → 切片关联(权重排序) → 文本回填 → 三元组转文本。特别注意
process返回 6 个调试字段(kg_entities/kg_aligned_entities/kg_seed_nodes/kg_triples/kg_chunks/kg_alignments),便于链路易链追溯。
def _expand_one_hop(
self, seed_nodes: List[NodeRef],
per_seed: int = 50, max_total: int = 200,
) -> List[Triple]:
"""扩展种子节点的一跳关系。"""
if not seed_nodes:
return []
try:
with self._neo4j_session() as session:
triples: List[Triple] = []
seen: Set[Tuple[str, ...]] = set()
for s in seed_nodes:
rows = session.execute_read(
self._tx_expand_triples,
s["name"], s["item_name"], per_seed
)
for tr in rows:
key = (tr["item_name"], tr["head"], tr["rel"], tr["tail"])
if key not in seen:
seen.add(key)
triples.append(tr)
if len(triples) >= max_total:
break
if len(triples) >= max_total:
break
return triples
except Exception as e:
self.logger.error(f"Neo4j 扩展异常: {e}")
return []Cypher 查询详解:
🟢 【P2 后面可以查】 测试代码——5 步执行链输出,看预期输出中的三元组和切片回填即可。
@staticmethod
def _tx_expand_triples(tx, seed_name: str, item_name: str, limit: int):
"""事务: 双向一跳扩展。"""
rows = tx.run("""
MATCH (seed:Entity {name: $seed, item_name: $item_name})
CALL {
-- 出边: seed → neighbor
WITH seed
MATCH (seed)-[r]->(nbr:Entity)
WHERE type(r) <> 'MENTIONED_IN'
AND nbr.item_name = $item_name
RETURN seed.name AS head, type(r) AS rel, nbr.name AS tail
UNION
-- 入边: neighbor → seed
WITH seed
MATCH (nbr:Entity)-[r]->(seed)
WHERE type(r) <> 'MENTIONED_IN'
AND nbr.item_name = $item_name
RETURN nbr.name AS head, type(r) AS rel, seed.name AS tail
}
RETURN head, rel, tail LIMIT $limit
""", seed=seed_name, item_name=item_name, limit=limit).data()
return [
{"head": r["head"], "rel": r["rel"],
"tail": r["tail"], "item_name": item_name}
for r in rows
]Step 6: 获取关联切片
目的: 通过 MENTIONED_IN 关系找到与实体相关的文档切片。
实现逻辑:
- 构建带权重的节点列表(种子权重 > 邻居权重)
- 查询 MENTIONED_IN 关系
- 按权重和计数排序
- 返回切片引用列表
代码片段:
def _get_chunk_refs(
self, seed_nodes: List[NodeRef], triples: List[Triple],
max_chunks: int = 200,
) -> List[Dict[str, Any]]:
"""通过 MENTIONED_IN 关系获取关联切片 ID(带权重排序)。"""
# 1. 构建带权重的节点列表
nodes = self._build_weighted_nodes(seed_nodes, triples)
if not nodes:
return []
try:
with self._neo4j_session() as session:
# 2. 查询关联切片
rows = session.run("""
UNWIND $nodes AS n
MATCH (e:Entity {name: n.name, item_name: n.item_name})
-[:MENTIONED_IN]->(c:Chunk {item_name: n.item_name})
WITH c, sum(n.w) AS score, count(DISTINCT e) AS cnt
RETURN c.id AS chunk_id, c.item_name AS item_name,
score, cnt
ORDER BY score DESC, cnt DESC, chunk_id ASC
LIMIT $limit
""", nodes=nodes, limit=max_chunks).data()
return [
{
"id": None,
"distance": float(r.get("score", 0)),
"entity": {
"chunk_id": str(r["chunk_id"]),
"item_name": str(r["item_name"])
}
}
for r in rows
]
except Exception as e:
self.logger.error(f"切片引用查询异常: {e}")
return []权重计算详解:
@staticmethod
def _build_weighted_nodes(
seed_nodes: List[NodeRef], triples: List[Triple],
w_seed: float = 2.0, w_neighbor: float = 1.0,
) -> List[Dict[str, Any]]:
"""构建带权重的节点列表(种子权重 > 邻居权重)。"""
weights: Dict[Tuple[str, str], float] = {}
# 种子节点: 权重 2.0
for s in seed_nodes or []:
key = (s["item_name"], s["name"])
weights[key] = max(weights.get(key, 0), w_seed)
# 邻居节点(三元组的头尾): 权重 1.0
for tr in triples or []:
it = tr["item_name"]
for n in (tr["head"], tr["tail"]):
key = (it, n)
weights[key] = max(weights.get(key, 0), w_neighbor)
return [
{"item_name": it, "name": n, "w": w}
for (it, n), w in weights.items()
]Step 7: Milvus 文本回填
目的: 根据切片 ID 从 Milvus 获取完整的切片内容。
实现逻辑:
- 提取去重的 chunk_id 列表
- 调用 fetch_chunks_by_ids 批量查询
- 按原顺序回填内容
- 合并 item_name 信息
代码片段:
def _fetch_chunk_texts(self, hits: List[Dict[str, Any]]) -> List[Dict[str, Any]]:
"""从 Milvus 回填切片文本。"""
from knowledge.tools.milvus_utils import fetch_chunks_by_ids, get_milvus_client
if not hits:
return []
collection = "chunks_test"
# 1. 提取去重 chunk_id
chunk_ids = list({
int(str(h["entity"]["chunk_id"]))
for h in hits
if h.get("entity", {}).get("chunk_id") is not None
})
if not chunk_ids:
return []
# 2. 从 Milvus 批量查询
try:
rows = fetch_chunks_by_ids(
client=get_milvus_client(),
collection_name=collection,
chunk_ids=chunk_ids,
output_fields=["chunk_id", "content", "title", "item_name"],
)
except Exception as e:
self.logger.error(f"Milvus 回填异常: {e}")
rows = []
# 3. 构建映射表
row_map = {
str(r["chunk_id"]): r
for r in (rows or [])
if r.get("chunk_id") is not None
}
# 4. 按原顺序回填
result = []
for h in hits:
ent = h.get("entity", {})
row = row_map.get(str(ent.get("chunk_id")))
if row:
merged = dict(row)
# 合并 item_name
if ent.get("item_name") and not merged.get("item_name"):
merged["item_name"] = ent["item_name"]
result.append(merged)
self.logger.info(f"回填完成: {len(result)} 条切片")
return resultStep 8: 三元组转文本
目的: 将结构化的三元组转换为自然语言描述,用于答案生成。
实现逻辑:
- 遍历三元组列表
- 格式化为 "[商品名] 头实体 -(关系)-> 尾实体" 格式
- 去重处理
代码片段:
@staticmethod
def _triples_to_docs(triples: List[Triple]) -> List[str]:
"""三元组 → 去重文本描述。"""
seen: Set[str] = set()
docs: List[str] = []
for tr in triples:
h, r, t = tr.get("head", ""), tr.get("rel", ""), tr.get("tail", "")
if not all([h, r, t]):
continue
it = tr.get("item_name", "")
doc = f"[{it}] {h} -({r})-> {t}" if it else f"{h} -({r})-> {t}"
if doc not in seen:
seen.add(doc)
docs.append(doc)
return docs转换示例:
输入三元组:
[
{head: "电池安装", rel: "REQUIRES", tail: "螺丝刀", item_name: "万用表RS-12"},
{head: "电池仓", rel: "CONTAINS", tail: "电池", item_name: "万用表RS-12"}
]
输出文本:
[
"[万用表RS-12] 电池安装 -(REQUIRES)-> 螺丝刀",
"[万用表RS-12] 电池仓 -(CONTAINS)-> 电池"
]4.4 代码实现
完整代码请参见 knowledge/processor/query_process/nodes/query_kg.py,此处展示核心结构:
# knowledge/processor/query_process/nodes/query_kg.py
"""知识图谱查询节点
从用户查询中抽取实体,经 Milvus 对齐后在 Neo4j 中检索相关子图,
最终回填切片文本内容。
"""
class QueryKgNode(BaseNode):
"""知识图谱查询节点。"""
name = "query_kg"
def process(self, state: QueryGraphState) -> QueryGraphState:
# 1. 预处理
question = state.get("rewritten_query") or state.get("original_query", "")
item_names = self._clean_item_names(state.get("item_names"))
for name in item_names:
question = question.replace(name, "")
# 2. LLM 实体抽取
entities = self._extract_entities(question, get_llm_client(json_mode=True))
# 3. Milvus 实体对齐
align_result = self._align_entities(entities, item_names) if entities else {}
aligned = align_result.get("aligned_entities", entities)
# 4. Neo4j 种子节点查找
seed_nodes = self._find_seed_nodes(aligned, item_names)
# 5. 一跳扩展
triples = self._expand_one_hop(seed_nodes)
# 6. 获取关联切片
kg_chunk_hits = self._get_chunk_refs(seed_nodes, triples)
# 7. Milvus 文本回填
kg_chunks = self._fetch_chunk_texts(kg_chunk_hits)
# 8. 返回结果
return {
"kg_chunks": kg_chunks,
"kg_triples": self._triples_to_docs(triples),
"kg_seed_nodes": seed_nodes,
"kg_entities": entities,
"kg_aligned_entities": aligned,
"kg_alignments": align_result.get("alignments", []),
}
# ... 其他方法实现5. 测试运行
5.1 运行知识图谱查询节点测试
# 进入项目目录
cd knowledge
# 激活虚拟环境
source .venv/bin/activate # Linux/Mac
# 或
.venv\Scripts\activate # Windows
# 运行测试
python -m knowledge.processor.query_process.nodes.query_kg5.2 测试代码
if __name__ == "__main__":
from dotenv import load_dotenv
import json
load_dotenv()
try:
from knowledge.processor.query_process.base import setup_logging
setup_logging()
except ImportError:
import logging
logging.basicConfig(level=logging.INFO)
print("=" * 60)
print("知识图谱查询节点测试")
print("=" * 60)
# 构造测试状态
test_state = {
"original_query": "RS-12数字万用表怎么测电压?",
"rewritten_query": "RS-12数字万用表怎么测电压?",
"item_names": ["RS-12数字万用表"]
}
print("【输入状态】:")
print(json.dumps(test_state, ensure_ascii=False, indent=2))
print("-" * 60)
try:
result = node_query_kg(test_state)
print("\n执行完成!全链路输出如下:\n")
# 第1步: LLM 实体抽取
print("[第1步] LLM 原始抽取实体 (kg_entities):")
print(f" {result.get('kg_entities', [])}")
# 第2步: Milvus 对齐
print("\n[第2步] Milvus 对齐后实体 (kg_aligned_entities):")
print(f" {result.get('kg_aligned_entities', [])}")
# 第3步: Neo4j 种子节点
print("\n[第3步] Neo4j 命中的种子节点 (kg_seed_nodes):")
for seed in result.get("kg_seed_nodes", []):
print(f" - {seed}")
# 第4步: 三元组
triples = result.get("kg_triples", [])
print(f"\n[第4步] 扩展的一跳知识三元组 (共 {len(triples)} 条):")
for t in triples[:5]:
print(f" - {t}")
if len(triples) > 5:
print(f" - ... (省略其余 {len(triples) - 5} 条)")
# 第5步: 切片
chunks = result.get("kg_chunks", [])
print(f"\n[第5步] 最终召回的切片 (共 {len(chunks)} 条):")
for i, chunk in enumerate(chunks[:3], 1):
print(f" [{i}] ID: {chunk.get('chunk_id')}")
print(f" 商品: {chunk.get('item_name')}")
print(f" 内容: {chunk.get('content', '')[:80]}...")
except Exception as e:
print(f"\n执行失败: {e}")
import traceback
traceback.print_exc()5.3 预期输出
============================================================
知识图谱查询节点测试
============================================================
【输入状态】:
{
"original_query": "RS-12数字万用表怎么测电压?",
"rewritten_query": "RS-12数字万用表怎么测电压?",
"item_names": ["RS-12数字万用表"]
}
------------------------------------------------------------
2024-01-15 10:30:00 - query.query_kg - INFO - --- query_kg 开始 ---
2024-01-15 10:30:00 - query.query_kg - INFO - [step_1] LLM 抽取实体
2024-01-15 10:30:01 - query.query_kg - INFO - 抽取到 2 个实体: ['测电压', '电压']
2024-01-15 10:30:01 - query.query_kg - INFO - [step_2] Milvus 对齐实体
2024-01-15 10:30:02 - query.query_kg - INFO - 对齐后实体: ['电压测量', '直流电压']
2024-01-15 10:30:02 - query.query_kg - INFO - [step_3] Neo4j 种子节点 + 扩展一跳
2024-01-15 10:30:02 - query.query_kg - INFO - 种子节点: 2
2024-01-15 10:30:03 - query.query_kg - INFO - 三元组: 8
2024-01-15 10:30:03 - query.query_kg - INFO - [step_4] 整理输出
2024-01-15 10:30:03 - query.query_kg - INFO - 回填完成: 5 条切片
2024-01-15 10:30:03 - query.query_kg - INFO - --- query_kg 完成 ---
执行完成!全链路输出如下:
[第1步] LLM 原始抽取实体 (kg_entities):
['测电压', '电压']
[第2步] Milvus 对齐后实体 (kg_aligned_entities):
['电压测量', '直流电压']
[第3步] Neo4j 命中的种子节点 (kg_seed_nodes):
- {'name': '电压测量', 'item_name': 'RS-12数字万用表'}
- {'name': '直流电压', 'item_name': 'RS-12数字万用表'}
[第4步] 扩展的一跳知识三元组 (共 8 条):
- [RS-12数字万用表] 电压测量 -(REQUIRES)-> 表笔
- [RS-12数字万用表] 电压测量 -(HAS_STEP)-> 选择量程
- [RS-12数字万用表] 电压测量 -(HAS_STEP)-> 连接表笔
- [RS-12数字万用表] 直流电压 -(MEASURED_BY)-> V-档位
- [RS-12数字万用表] 直流电压 -(REQUIRES)-> 极性正确
- ... (省略其余 3 条)
[第5步] 最终召回的切片 (共 5 条):
[1] ID: 12345
商品: RS-12数字万用表
内容: 电压测量是万用表最常用的功能之一。测量直流电压时,将旋钮转到V-档位...
[2] ID: 12346
商品: RS-12数字万用表
内容: 测量交流电压时,将功能旋钮转到V~档位,选择适当量程...
[3] ID: 12347
商品: RS-12数字万用表
内容: 注意事项:测量高压时请确保量程足够,超量程测量可能损坏仪表...5.4 处理前后对比
| 阶段 | 数据 | 说明 |
|---|---|---|
| 输入 | "RS-12数字万用表怎么测电压?" | 用户原始问题 |
| 预处理后 | "怎么测电压?" | 移除商品名 |
| 实体抽取 | ["测电压", "电压"] | LLM 提取的实体 |
| 实体对齐 | ["电压测量", "直流电压"] | 对齐到图谱标准名 |
| 种子节点 | 2 个节点 | Neo4j 中找到的实体节点 |
| 三元组 | 8 条 | 一跳扩展的关系 |
| 切片 | 5 条 | 关联的文档切片(带完整文本) |
图谱查询的独特价值:
向量检索只能找到:
- 包含"电压"、"测量"关键词的文档
图谱查询额外提供:
- 电压测量需要的工具(表笔)
- 电压测量的步骤(选择量程、连接表笔)
- 直流电压对应的档位(V-档位)
- 测量的前提条件(极性正确)
→ 结构化知识使答案更完整、更准确6. 总结
6.1 节点功能概览
| 功能 | 说明 |
|---|---|
| 实体抽取 | 使用 LLM 从问题中提取关键实体 |
| 实体对齐 | 通过向量相似度将实体对齐到图谱 |
| 种子查找 | 在 Neo4j 中定位实体节点(精确 + 模糊) |
| 一跳扩展 | 获取种子节点的邻居关系 |
| 切片关联 | 通过 MENTIONED_IN 找到相关文档 |
| 文本回填 | 从 Milvus 获取切片完整内容 |
| 三元组转文本 | 将结构化关系转为自然语言描述 |
6.2 节点设计要点
1. 多阶段流水线架构
LLM抽取 → Milvus对齐 → Neo4j查询 → Milvus回填
↓ ↓ ↓ ↓
实体名 标准实体名 节点+关系 完整切片- 每个阶段独立,便于调试和优化
- 失败时优雅降级,返回空结果
2. 实体对齐的必要性
用户说: "电池更换"
图谱中: "电池安装"
直接查询: 无结果
向量对齐后: 匹配成功- 向量相似度容忍用词差异
- 提高图谱查询命中率
3. 精确匹配优先策略
# 先精确匹配
MATCH (n) WHERE n.name = $name ...
# 无结果则模糊匹配
MATCH (n) WHERE toLower(n.name) CONTAINS toLower($name) ...- 精确匹配性能更好、更准确
- 模糊匹配作为兜底
4. 带权重的切片排序
w_seed = 2.0 # 种子节点权重更高
w_neighbor = 1.0 # 邻居节点权重较低
# 切片得分 = 关联节点权重之和- 与用户问题直接相关的实体权重更高
- 多个实体关联的切片排名更靠前
5. 限流控制防止爆炸
per_entity = 3 # 每个实体最多 3 个种子
max_total = 30 # 种子总数上限 30
per_seed = 50 # 每个种子最多 50 条三元组
max_total = 200 # 三元组总数上限 200- 防止图遍历导致结果过多
- 保证查询性能
6. 丰富的调试信息
return {
"kg_chunks": kg_chunks, # 最终切片
"kg_triples": triples_docs, # 三元组文本
"kg_seed_nodes": seed_nodes, # 种子节点
"kg_entities": entities, # 原始实体
"kg_aligned_entities": aligned, # 对齐实体
"kg_alignments": alignments, # 对齐详情
}- 返回中间结果便于调试
- 支持可解释性分析
企业痛点映射
| 痛点 | 传统方案 | AI Agent KG 查询方案 | 效率提升 |
|---|---|---|---|
| 向量检索无法理解实体关系 | 搜不到"电池安装需要螺丝刀"这种关系 | 图遍历返回结构化的三元组 | 关系查询覆盖度从 ~20% 提升至 ~80% |
| 用户表述与图谱命名不一致 | 用户说"换电池",图谱里是"电池安装" | Milvus 向量对齐容错 | 实体命中率从 ~40% 提升至 ~85% |
| 图遍历结果过多不可控 | 一次查询返回数千条关系 | 4 级限流(per_entity/max_total/per_seed/max_total) | 结果可控在 200 条以内 |
| 种子 vs 邻居权重无区分 | 所有关联切片等权 | w_seed=2.0 / w_neighbor=1.0 带权排序 | 核心切片排名提升 ~200% |
Remote & Agent 应用场景价值
Remote 场景价值:同时协调 LLM API、Milvus 集群、Neo4j 服务器三个远程服务。团队通过
.env配置即可运行。execute_read事务机制确保读取一致性。Agent 落地场景:QueryKgNode 可拆分为 3 个顺序 Agent:(1) 实体抽取 Agent——调用 LLM 提取实体 (2) 实体对齐 Agent——向量搜索匹配标准名 (3) 图遍历 Agent——Neo4j 种子 + 一跳扩展 + 切片回填。每个 Agent 的输出是下一个 Agent 的输入,形成 Pipeline 模式。
Git Commit 对应
本节 KG 查询节点相关提交记录(参考值,以实际版本为准):
<待补充 — 建议搜索 "query_kg.py" 相关提交>cd shopkeeper_brain
git log --oneline --all -- knowledge/processor/query_process/nodes/query_kg.py