一批用户关系导进 Neo4j,脚本没报错,日志也全是成功,库里的节点却越跑越多。这种问题我一般不先怀疑 Neo4j,先看写入代码。果然,外层一个 for 循环,里面每次执行一条 CREATE。脚本重跑一次,数据就复制一遍。
图数据库写起来看着直观,真正容易出问题的地方,还是唯一约束、批量写入和事务边界。
Python 连接 Neo4j,用官方驱动就够了:
pip install neo4j
连接代码别散落在业务函数里,更不要每处理一条数据就创建一次 Driver。
import os
from neo4j import GraphDatabase
defbuild_driver():
uri = os.getenv("NEO4J_URI", "neo4j://127.0.0.1:7687")
user = os.getenv("NEO4J_USER", "neo4j")
password = os.environ["NEO4J_PASSWORD"]
driver = GraphDatabase.driver(uri, auth=(user, password))
driver.verify_connectivity()
return driver
Neo4j 官方 Python 驱动可以通过 Driver.execute_query() 执行 Cypher,并使用参数传值。密码别写死在代码里,这种东西迟早会被提交到仓库。
开始导数据前,我会先把唯一约束建好。
definit_schema(driver):
driver.execute_query(
"""
CREATE CONSTRAINT user_uid_unique IF NOT EXISTS
FOR (u:User)
REQUIRE u.uid IS UNIQUE
""",
database_="neo4j",
)
没有这个约束,下面这条语句虽然用了 MERGE,数据模型一旦写乱,后面排重复节点还是麻烦。
唯一约束不只是拦截脏数据,背后还会建立对应索引。已经存在重复数据时,约束创建可能直接失败,这时候别反复执行脚本,先把重复节点查出来处理。
真正写用户节点时,别这么干:
for user in users:
driver.execute_query(
"CREATE (:User {uid: $uid, name: $name})",
uid=user["uid"],
name=user["name"],
)
问题有两个:一条数据一次网络交互,而且脚本不能安全重跑。
我通常把数据切成小批次,再用 UNWIND 展开。
defimport_users(driver, users, chunk_size=800):
cypher = """
UNWIND $rows AS row
MERGE (u:User {uid: row.uid})
SET u.name = row.name,
u.department = row.department,
u.updated_at = datetime()
"""
for offset in range(0, len(users), chunk_size):
rows = users[offset: offset + chunk_size]
driver.execute_query(
cypher,
rows=rows,
database_="neo4j",
)
UNWIND 会把参数中的列表展开成多行,后面的 MERGE、SET 就能按行处理。比 Python 在外面一条条提交干净得多。
关系写入也差不多,不过这里有个地方我见过不少人写错:把会变化的属性塞进 MERGE。
defimport_follow_relations(driver, relations):
driver.execute_query(
"""
UNWIND $links AS link
MATCH (source:User {uid: link.source_uid})
MATCH (target:User {uid: link.target_uid})
MERGE (source)-[r:FOLLOWS]->(target)
ON CREATE SET r.created_at = datetime()
SET r.channel = link.channel,
r.updated_at = datetime()
""",
links=relations,
database_="neo4j",
)
这里 MERGE 只负责确定关系是否存在,channel 这种可能变化的字段放在 SET 里。
要是写成下面这样:
MERGE (source)-[:FOLLOWS {channel: link.channel}]->(target)
用户换一次关注渠道,就可能多出一条关系。语法没错,模型错了,而且这种问题通常要等查询结果重复时才会被发现。
查询图关系时,也别一兴奋就写无限深度路径。
deffind_following(driver, uid, max_rows=50):
records, _, _ = driver.execute_query(
"""
MATCH path =
(:User {uid: $uid})-[:FOLLOWS*1..2]->(target:User)
RETURN DISTINCT
target.uid AS uid,
target.name AS name,
length(path) AS distance
ORDER BY distance, uid
LIMIT $limit
""",
uid=uid,
limit=max_rows,
database_="neo4j",
)
return [record.data() for record in records]
*1..2 这个范围不能随手删。关系密集的图里,路径深度放开后,中间结果可能迅速膨胀。页面只展示几十条数据,就别让数据库先把整张关系网翻一遍。
还有一种操作,必须放在事务里:先删除旧关系,再重建新关系。
defreplace_document_tags(driver, document_id, tags):
defwrite_tags(tx):
tx.run(
"""
MATCH (:Document {id: $document_id})-[r:HAS_TAG]->()
DELETE r
""",
document_id=document_id,
).consume()
tx.run(
"""
MATCH (doc:Document {id: $document_id})
UNWIND $tags AS tag_name
MERGE (tag:Tag {name: tag_name})
MERGE (doc)-[:HAS_TAG]->(tag)
""",
document_id=document_id,
tags=tags,
).consume()
with driver.session(database="neo4j") as session:
session.execute_write(write_tags)
删除成功、重建失败,文档标签就空了。两个动作塞进同一个事务,要么一起提交,要么一起回滚。execute_query() 本身会自动使用事务;涉及多条查询和中间处理时,再使用 session.execute_write() 这种托管事务。
最后别忘了关连接:
defrun_import(users, relations):
driver = build_driver()
try:
init_schema(driver)
import_users(driver, users)
import_follow_relations(driver, relations)
finally:
driver.close()
Python 操作 Neo4j 的代码不算多,坑主要藏在 Cypher 里。
节点有没有稳定业务主键,MERGE 匹配了哪些字段,关系能不能重复,路径最大查几层,批量任务失败后能不能重跑——这些地方不先想清楚,驱动连接得再漂亮也没用。
尤其是那个套着 CREATE 的 for 循环,看见了就尽早删。