Skip to content

知识库查询 —— 知识图谱查询节点

本文档详细介绍知识图谱查询节点(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" vs git log --committer="me"——直接匹配的(作者)权重更高,间接关联的(提交者)权重减半。

1. 任务目标

1.1 本章目标

通过本章学习,你将掌握:

  1. 理解知识图谱在 RAG 中的作用:掌握结构化知识与向量检索的结合方式
  2. 学会 LLM 实体抽取:从自然语言问题中提取关键实体
  3. 掌握实体对齐技术:将抽取的实体与图谱中的标准实体名对齐
  4. 理解图遍历策略:种子节点查找与一跳扩展的实现方式
  5. 学会多数据源协作:LLM + Milvus + Neo4j 的联合查询模式
  6. 实现完整的图谱查询节点:通过 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,从问题中移除商品名避免干扰实体抽取。

python
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 关系与原文切片关联:

cypher
(:Entity {name: "电池安装"}) -[:MENTIONED_IN]-> (:Chunk {id: 12345})

通过这个关系,可以找到与实体相关的原文切片,获取详细的上下文信息。

2.6 带权重的节点排序

🔥 【P0 必须要学】 Step2 LLM 实体抽取——使用 ENTITY_EXTRACT_SYSTEM_PROMPT 约束 LLM 输出 JSON 数组。关键设计enable_thinking=False(禁止推理链,直接输出 JSON),set 去重。

python
# 种子节点权重更高
w_seed = 2.0      # 直接匹配的实体

# 邻居节点权重较低
w_neighbor = 1.0  # 一跳扩展的实体

# 切片得分 = 关联节点权重之和
# 多个高权重实体关联的切片排名更靠前

3. 知识图谱查询业务处理流程(总)

3.1 整体流程图

3.2 数据流转详解


4. 知识图谱查询业务处理流程(分)

4.1 目标

实现一个能够:

  1. 从用户问题中智能抽取关键实体
  2. 将抽取实体与知识图谱中的标准实体对齐
  3. 在图数据库中找到相关节点和关系
  4. 将图谱信息转化为可用于答案生成的文本
  5. 回填原文切片,提供详细上下文

4.2 需求分析

4.2.1 功能需求

  1. 实体抽取:使用 LLM 从自然语言问题中提取实体
  2. 实体对齐:通过向量相似度将实体对齐到知识图谱
  3. 种子查找:在 Neo4j 中定位对应节点(精确 + 模糊匹配)
  4. 关系扩展:获取种子节点的一跳邻居和关系
  5. 切片关联:通过 MENTIONED_IN 关系找到相关文档切片
  6. 文本回填:从 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)不影响其他实体的对齐。

python
# 实体对齐
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: 预处理输入

目的: 清理输入数据,为后续处理做准备。

实现逻辑:

  1. 获取查询文本(优先 rewritten_query)
  2. 清理商品名称列表(处理 None/str/list)
  3. 从问题中移除商品名,避免干扰实体抽取

代码片段:

🟡 【P1 看注释就行】 Step4 Neo4j 种子查找——先精确匹配 n.name = $name,失败则模糊匹配 CONTAINS toLower($name)per_entity=3 + max_total=30 防止结果爆炸。

python
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() 确保大小写不敏感。

python
@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 实体抽取

目的: 从用户问题中提取可能存在于知识图谱中的实体名称。

实现逻辑:

  1. 使用专门的实体抽取提示词
  2. 调用 LLM(JSON 模式)
  3. 解析 JSON 响应,提取 entities 字段
  4. 去重和清理

代码片段:

🟡 【P1 看注释就行】 Step5 一跳扩展——双向遍历(出边 + 入边),排除 MENTIONED_IN 关系。per_seed=50 + max_total=200 限流。

python
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 抽取的实体与知识图谱中的标准实体名对齐。

实现逻辑:

  1. 为每个实体生成混合向量
  2. 在 entity_collection 中执行混合搜索
  3. 根据相似度分数选择最佳匹配
  4. 低于阈值的视为未命中
  5. 记录对齐详情供调试

代码片段:

🟡 【P1 看注释就行】 Cypher 扩展——CALL { MATCH ... UNION MATCH ... } 的双向查询模式。注意排除 MENTIONED_IN 的 WHERE 条件。

python
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 排序。这个权重策略决定了哪些切片排在前面。

python
{
    "aligned_entities": ["电池安装"],
    "alignments": [
        {"original": "电池", "aligned": "电池安装", "score": 0.92},
        {"original": "更换", "aligned": None, "reason": "no_hit"}
    ]
}

Step 4: Neo4j 种子节点查找

目的: 在 Neo4j 图数据库中找到对齐实体对应的节点。

实现逻辑:

  1. 使用事务读取(execute_read)
  2. 先尝试精确匹配(name = $name)
  3. 精确匹配失败则模糊匹配(CONTAINS)
  4. 限制每个实体的种子数和总数

代码片段:

🟡 【P1 看注释就行】 Step7 Milvus 文本回填——去重 chunk_id → fetch_chunks_by_ids 批量查询 → 按原顺序回填。注意 row_map 的映射构建。

python
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 去重确保无重复。

python
@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: 一跳扩展

目的: 获取种子节点的一跳邻居关系,形成三元组。

实现逻辑:

  1. 双向遍历(出边 + 入边)
  2. 排除 MENTIONED_IN 关系(这是切片引用关系)
  3. 限制每个种子的三元组数和总数
  4. 去重处理

代码片段:

🔥 【P0 必须要学】 QueryKgNode 完整架构——8 步流程:预处理 → LLM抽取 → Milvus对齐 → Neo4j种子 → 一跳扩展 → 切片关联(权重排序) → 文本回填 → 三元组转文本。特别注意 process 返回 6 个调试字段(kg_entities / kg_aligned_entities / kg_seed_nodes / kg_triples / kg_chunks / kg_alignments),便于链路易链追溯。

python
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 步执行链输出,看预期输出中的三元组和切片回填即可。

python
@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 关系找到与实体相关的文档切片。

实现逻辑:

  1. 构建带权重的节点列表(种子权重 > 邻居权重)
  2. 查询 MENTIONED_IN 关系
  3. 按权重和计数排序
  4. 返回切片引用列表

代码片段:

python
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 []

权重计算详解:

python
@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 获取完整的切片内容。

实现逻辑:

  1. 提取去重的 chunk_id 列表
  2. 调用 fetch_chunks_by_ids 批量查询
  3. 按原顺序回填内容
  4. 合并 item_name 信息

代码片段:

python
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 result

Step 8: 三元组转文本

目的: 将结构化的三元组转换为自然语言描述,用于答案生成。

实现逻辑:

  1. 遍历三元组列表
  2. 格式化为 "[商品名] 头实体 -(关系)-> 尾实体" 格式
  3. 去重处理

代码片段:

python
@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,此处展示核心结构:

python
# 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 运行知识图谱查询节点测试

bash
# 进入项目目录
cd knowledge

# 激活虚拟环境
source .venv/bin/activate  # Linux/Mac
# 或
.venv\Scripts\activate     # Windows

# 运行测试
python -m knowledge.processor.query_process.nodes.query_kg

5.2 测试代码

python
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. 精确匹配优先策略

python
# 先精确匹配
MATCH (n) WHERE n.name = $name ...

# 无结果则模糊匹配
MATCH (n) WHERE toLower(n.name) CONTAINS toLower($name) ...
  • 精确匹配性能更好、更准确
  • 模糊匹配作为兜底

4. 带权重的切片排序

python
w_seed = 2.0      # 种子节点权重更高
w_neighbor = 1.0  # 邻居节点权重较低

# 切片得分 = 关联节点权重之和
  • 与用户问题直接相关的实体权重更高
  • 多个实体关联的切片排名更靠前

5. 限流控制防止爆炸

python
per_entity = 3      # 每个实体最多 3 个种子
max_total = 30      # 种子总数上限 30
per_seed = 50       # 每个种子最多 50 条三元组
max_total = 200     # 三元组总数上限 200
  • 防止图遍历导致结果过多
  • 保证查询性能

6. 丰富的调试信息

python
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" 相关提交>
bash
cd shopkeeper_brain
git log --oneline --all -- knowledge/processor/query_process/nodes/query_kg.py

OPC 超级个体实战指南