Neo4j 实战

一、为什么用 Neo4j

给湖北理工计算机学院官网助手做 Graph RAG 时,第一步是把文档抽成三元组,第二步就得找个地方存这些三元组。

最开始我想用 NetworkX——Python 内存图,简单。跑了几次多跳查询之后放弃了。不是 NetworkX 不好,是我的场景不对。学生问"要学人工智能方向,得先学哪些课",这个查询在图上是沿着 HAS_PREREQUISITE 边跳 3-4 步——NetworkX 在内存里构建邻接表再遍历,100 个节点还行,1000 个节点就慢了。而且数据在内存里,服务一重启就没了。

换成关系数据库(SQLite)呢?存三元组倒是简单——建个三元组表,存 (主体, 关系, 客体)。但查多跳的时候 SQL 得写 WITH RECURSIVE,三步就绕晕了。

Neo4j 专门为这种场景设计的——存的是节点和边,查的是路径。多跳查询就是一行 Cypher。

二、安装

Docker Compose 一行搞定:

# docker-compose.yml
services:
  neo4j:
    image: neo4j:5.21
    ports:
      - "7474:7474"   # HTTP
      - "7687:7687"   # Bolt(程序连这个)
    environment:
      NEO4J_AUTH: neo4j/your_password_here
      NEO4J_PLUGINS: '["apoc"]'   # APOC 扩展,备份/批量操作用
    volumes:
      - neo4j_data:/data
      - neo4j_logs:/logs

volumes:
  neo4j_data:
  neo4j_logs:
docker compose up -d
# 等 10 秒
curl http://localhost:7474   # 能看到 Neo4j Browser

浏览器打开 http://localhost:7474,用 neo4j/your_password_here 登录,顶上那个输入框就是 Cypher 命令行——后面所有查询都可以先在这里试,再写到代码里。

Python 连接:

pip install neo4j
from neo4j import GraphDatabase

driver = GraphDatabase.driver(
    "bolt://localhost:7687",
    auth=("neo4j", "your_password_here")
)

# 测试连接
driver.verify_connectivity()
print("连上了")

三、Schema 设计——实体和关系怎么建模

湖北理工官网的知识图谱场景:从教务处文档里抽取出教师、课程、院系、实验室等实体,以及它们之间的关系。

核心原则:所有实体一个标签,用属性区分类型。

// 实体节点
CREATE (:Entity {
    id: "e001",
    name: "王建国",
    category: "Professor",       // 类型:Professor / Course / Department / Lab
    aliases: ["王建国教授", "王老师"],
    source_chunks: ["chunk_042", "chunk_108"],
    properties: "{title: '教授', research_area: 'NLP'}",
    created_at: 1710000000
})

为什么要全部放一个 Entity 标签而不是拆成 ProfessorCourse 多个标签?因为 LLM 抽出来的实体类型是动态的——今天可能是 Scholarship,明天可能是 Competition。拆标签的话每次都要改 schema,代码也要跟着改。一个标签 + category 字段,LLM 随便输出什么类型都能存。

关系边:

// 王建国 -[TEACHES]-> 数据结构与算法
MATCH (t:Entity {name: "王建国"}), (c:Entity {name: "数据结构与算法"})
CREATE (t)-[:TEACHES {
    since: "2020",
    chunk_ref: "chunk_042"    // 这条关系是从哪个 chunk 抽出来的
}]->(c)

每条边必须带 chunk_ref——知道它来自哪个 chunk。后面增量更新时会靠这个字段定位和删除旧关系。

索引——最容易被跳过的步骤:

// 按名称查实体是最频繁的操作,必须建索引
CREATE INDEX entity_name FOR (n:Entity) ON (n.name);
CREATE INDEX entity_category FOR (n:Entity) ON (n.category);

// 社区检测后每个实体有个 community_id,也要检索
CREATE INDEX entity_community FOR (n:Entity) ON (n.community_id);

没索引的话,每次 MATCH (n:Entity {name: "xxx"}) 是全表扫描。300 个节点时感知不到,3000 个节点时每次查询多 50ms,积累起来就明显了。

四、写入——批量导入 + 去重

从文档里抽取出几百条三元组后,需要写入 Neo4j。核心操作是 MERGE——节点或边已存在就跳过,不存在才创建。

from neo4j import GraphDatabase

def batch_insert_triples(driver, triples: list[dict], batch_size: int = 100):
    """批量写入三元组,用 MERGE 自动去重"""
    
    with driver.session() as session:
        for i in range(0, len(triples), batch_size):
            batch = triples[i:i + batch_size]
            
            for t in batch:
                session.run("""
                    // 主体节点:name 匹配就复用,否则创建
                    MERGE (src:Entity {name: $src_name})
                    ON CREATE SET 
                        src.id = $src_id,
                        src.category = $src_type,
                        src.source_chunks = [$chunk_id],
                        src.created_at = timestamp()
                    ON MATCH SET 
                        src.aliases = CASE 
                            WHEN $src_alias IS NOT NULL 
                            THEN apoc.coll.union(coalesce(src.aliases, []), [$src_alias])
                            ELSE src.aliases 
                        END,
                        src.source_chunks = apoc.coll.union(
                            coalesce(src.source_chunks, []), [$chunk_id]
                        )
                    
                    // 客体节点
                    MERGE (tgt:Entity {name: $tgt_name})
                    ON CREATE SET 
                        tgt.id = $tgt_id,
                        tgt.category = $tgt_type,
                        tgt.source_chunks = [$chunk_id],
                        tgt.created_at = timestamp()
                    
                    // 关系:src 和 tgt 的组合已存在就跳过
                    MERGE (src)-[r:RELATES {type: $rel_type}]->(tgt)
                    ON CREATE SET r.chunk_ref = $chunk_id
                """, {
                    "src_name": t["source_name"],
                    "src_id": t["source_id"],
                    "src_type": t["source_type"],
                    "src_alias": t.get("source_alias"),
                    "tgt_name": t["target_name"],
                    "tgt_id": t["target_id"],
                    "tgt_type": t["target_type"],
                    "rel_type": t["relation"],
                    "chunk_id": t["chunk_id"],
                })

ON CREATE SET 只在节点首次创建时执行,ON MATCH SET 在节点已存在时执行——给已有实体追加别名和来源 chunk。

写入性能: 单条 MERGE 大概 3-5ms。1000 条三元组串行写入要 4-5 秒。如果上万条,开事务批量提交:

def bulk_insert_with_transaction(driver, triples: list[dict]):
    """用事务批量写,比单条 MERGE 快 3-5 倍"""
    with driver.session() as session:
        tx = session.begin_transaction()
        for i, t in enumerate(triples):
            tx.run("""
                MERGE (src:Entity {name: $src_name})
                ON CREATE SET src.id = $src_id, src.category = $src_type
                MERGE (tgt:Entity {name: $tgt_name})
                ON CREATE SET tgt.id = $tgt_id, tgt.category = $tgt_type
                MERGE (src)-[r:RELATES {type: $rel_type}]->(tgt)
            """, {
                "src_name": t["source_name"], "src_id": t["source_id"],
                "src_type": t["source_type"],
                "tgt_name": t["target_name"], "tgt_id": t["target_id"],
                "tgt_type": t["target_type"], "rel_type": t["relation"],
            })
            
            if (i + 1) % 500 == 0:
                tx.commit()   # 每 500 条提交一次
                tx = session.begin_transaction()
        
        tx.commit()  # 提交剩余的

五、查询——从单跳到多跳

5.1 点查询

-- 查一个实体
MATCH (e:Entity {name: "王建国"})
RETURN e.name, e.category, e.aliases, e.properties

5.2 单跳——沿着一条边找邻居

-- "王建国教哪些课?"
MATCH (t:Entity {name: "王建国"})-[r:RELATES {type: "teaches"}]->(c:Entity)
RETURN c.name, c.aliases
# Python 代码里执行同样的查询
def query_teacher_courses(driver, teacher_name: str) -> list[dict]:
    with driver.session() as session:
        result = session.run("""
            MATCH (t:Entity {name: $name})-[r:RELATES {type: "teaches"}]->(c:Entity)
            RETURN c.name AS course, c.aliases AS aliases
        """, {"name": teacher_name})
        return [record.data() for record in result]

5.3 多跳——沿着边走几步

这是 Neo4j 跟关系数据库拉开差距的地方。查"学人工智能之前得先学哪些课":

-- 沿着 HAS_PREREQUISITE 边跳 1 到 4 步
MATCH path = (start:Entity {name: "人工智能"})
             -[:RELATES {type: "has_prerequisite"}*1..4]->(prereq:Entity)
RETURN [node in nodes(path) | node.name] AS 课程路径,
       length(path) AS 跳数
ORDER BY 跳数

*1..4 是可变长度路径——最少 1 跳,最多 4 跳。Neo4j 内部用 BFS 遍历,每一步只检查邻接边,不像关系数据库要递归 JOIN。

def query_prerequisite_chain(driver, course_name: str, max_hops: int = 4):
    """查一门课的先修链"""
    with driver.session() as session:
        result = session.run("""
            MATCH path = (start:Entity {name: $name})
                         -[:RELATES {type: "has_prerequisite"}*1..$max_hops]->(prereq:Entity)
            RETURN [node in nodes(path) | node.name] AS chain,
                   length(path) AS hops
            ORDER BY hops
        """, {"name": course_name, "max_hops": max_hops})
        return [record.data() for record in result]

5.4 反向查——谁指向我

-- "数据结构与算法有哪些先修课?" 反过来查"哪些课把它当先修课"
MATCH (course:Entity {name: "数据结构与算法"})
      <-[:RELATES {type: "has_prerequisite"}]-(dependent:Entity)
RETURN dependent.name AS 后续课程

5.5 两跳内找全部关联

-- "跟王建国有关的都有什么?"(两跳内所有实体和关系)
MATCH (t:Entity {name: "王建国"})-[r1:RELATES*1..2]-(related:Entity)
RETURN DISTINCT related.name, related.category, related.aliases

六、跟 RAG 系统对接

Agent 接到用户问题后,怎么决定是查向量库还是查 Neo4j?

def route_query(user_question: str) -> str:
    """判断走向量检索还是图查询"""
    graph_keywords = [
        "谁教", "老师", "教授", "先修课", "前置课程",
        "属于哪个", "有哪些课", "课程体系", "关系",
    ]
    
    for kw in graph_keywords:
        if kw in user_question:
            return "graph"     # 走 Neo4j
    
    return "vector"            # 走 Milvus

图查询走 Neo4j 拿到实体关系后,再带上原始 chunk 一起给 LLM:

def answer_with_graph(driver, user_question: str) -> str:
    """结合图查询 + 向量检索回答"""
    
    # 1. 从问题中提取实体名(简单做法:送 LLM 抽取)
    entities = extract_entity_names(user_question)
    
    # 2. 在 Neo4j 中查关联
    graph_context = ""
    for name in entities:
        result = driver.session().run("""
            MATCH (e:Entity {name: $name})-[r:RELATES*1..2]-(related:Entity)
            RETURN DISTINCT e.name AS 起点, type(r) AS 关系, 
                   related.name AS 终点, related.category AS 类型
        """, {"name": name})
        
        for record in result:
            graph_context += f"{record['起点']} --[{record['关系']}]--> {record['终点']} ({record['类型']})\n"
    
    # 3. 把图结构 + 原始 chunk 一起给 LLM
    prompt = f"""根据以下知识图谱关系回答问题:

{graph_context}

用户问题:{user_question}

如果图中有答案,直接回答。如果图中信息不够,说明需要补充什么信息。"""
    
    return llm.invoke(prompt)

七、增量更新——文档变了图怎么跟着变

新文档上传后,不是全量重建。以 chunk 为粒度做增量:

def incremental_update(driver, new_chunks: list[dict]):
    """chunk 级增量更新图"""
    
    with driver.session() as session:
        for chunk in new_chunks:
            chunk_id = chunk["id"]
            chunk_md5 = chunk["md5"]
            
            # 1. 检查这个 chunk 变没变
            existing = session.run("""
                MATCH ()-[r {chunk_ref: $chunk_id}]->()
                RETURN count(r) AS edge_count, r.chunk_md5 AS old_md5 LIMIT 1
            """, {"chunk_id": chunk_id}).single()
            
            if existing and existing.get("old_md5") == chunk_md5:
                continue  # 内容没变,跳过
            
            # 2. 内容变了——删旧边
            if existing and existing["edge_count"] > 0:
                session.run("""
                    MATCH ()-[r {chunk_ref: $chunk_id}]->()
                    DELETE r
                """, {"chunk_id": chunk_id})
            
            # 3. 重新抽取 + 写入新三元组
            triples = extract_triples_from_chunk(chunk["text"])
            for t in triples:
                session.run("""
                    MERGE (src:Entity {name: $src_name})
                    ON CREATE SET src.id = $src_id, src.category = $src_type
                    MERGE (tgt:Entity {name: $tgt_name})
                    ON CREATE SET tgt.id = $tgt_id, tgt.category = $tgt_type
                    MERGE (src)-[r:RELATES {type: $rel_type}]->(tgt)
                    SET r.chunk_ref = $chunk_id, r.chunk_md5 = $chunk_md5
                """, {
                    "src_name": t["source_name"], "src_id": t["source_id"],
                    "src_type": t["source_type"],
                    "tgt_name": t["target_name"], "tgt_id": t["target_id"],
                    "tgt_type": t["target_type"], "rel_type": t["relation"],
                    "chunk_id": chunk_id, "chunk_md5": chunk_md5,
                })

两个要点:

  • 只删边不删节点——实体节点即使暂时没有边连着也留着。"王建国以前教过什么课?"这种历史查询还需要旧实体
  • 边上加 chunk_md5——判断内容是否变化的依据。跟 同一份文件多个版本的管理策略 的思路一致

八、踩过的坑

1. MERGE 不加索引慢十倍。 最开始没给 name 建索引,2000 条三元组写入从 3 秒变成 30 秒。MERGE 内部先做 MATCH 查找是否存在——没索引就是全表扫描。

2. 关系类型别动态生成。 最早我把关系类型(teacheshas_prerequisite)直接当关系标签存:(t)-[:TEACHES]->(c)。看起来干净,但 LLM 抽出来的关系类型是动态的——今天多一个 cooperates_with,Cypher 查询就得跟着改。后来统一用一个 RELATES 标签 + type 属性:(t)-[:RELATES {type: "teaches"}]->(c)。查的时候 WHERE r.type = "xxx" 就行,不用每次新增关系类型就改查询。

3. 事务开太大 OOM。 一次 commit 塞了 5000 条,Neo4j 事务日志撑到 2GB,直接报错。改成每 500 条 commit 一次就好了。

4. APOC 插件没装导致 apoc.coll.union 调用失败。 docker-compose.ymlNEO4J_PLUGINS: '["apoc"]' 别忘了。

5. 多跳查询不加深度限制。 [*] 不设上限的话在图密集区域会遍历所有节点,一次查询十几秒。永远加 *1..N 限制跳数。

九、总结

Neo4j 在 Graph RAG 里的定位很清楚:它不存原始文档,只存实体和关系。 回答问题时,图给出"谁跟谁什么关系",向量库给出"原文是怎么说的",两者拼起来才是完整答案。

核心操作就四板斧:

  1. MERGE 写入 + 去重
  2. 索引 给 name/category/community_id 建索引
  3. 可变长度路径 *1..N 查多跳
  4. chunk_ref 挂边上做增量更新的锚点

2000 行代码不到,但给 RAG 加了一层"理解关系"的能力——从"查一段话"升级到了"查一张图"。