引言
「两个事务同时改同一条边,会发生什么?」在关系型数据库里,这个问题有标准答案:MVCC + 行锁 + 隔离级别,一套体系覆盖了绝大多数场景。到了图数据库,答案变得微妙:锁的粒度落在节点与关系上,而不是行;事务里一次遍历会碰到成百上千个实体,写写冲突的概率被放大;更麻烦的是,图上的业务语义(「更新这个用户的所有关注关系」)天然是范围写,范围写正是死锁的温床。
与此同时,图库的使用场景又特别吃批量更新:每天从上游同步几百万条实体与关系、把图算法的结果写回、给全量节点补一个属性。这些批量任务与在线读写并发执行时,如果事务边界、批大小、重试策略没设计好,会同时出现三类症状——吞吐上不去、偶发死锁、以及最难查的「丢更新」(lost update,两次并发写只有一次生效)。
本文按「并发模型 → 冲突类型 → 控制手段 → 批量吞吐 → 集群一致性 → 幂等实现」的顺序把这个问题拆开讲,重点是给出可落地的参数与代码,而不是理论推导。事务与索引的基础见 事务与索引调优 ;批量写入的具体手法可参考 图导入与 ETL ;通用的事务理论横向对照 数据库事务与隔离级别 与 并发安全容器 。
1. 图库的并发模型
1.1 锁粒度:节点与关系
Neo4j 的写锁粒度是单个节点与单条关系。当你执行 SET n.x = 1 时,只有 n 被锁住;执行 CREATE (a)-[:R]->(b) 时,会锁住 a、b 以及新关系。这与关系型库的「行锁」形似,但图的特殊性在于:
- 一次 Cypher 会触碰大量实体。
MATCH (a)-[:F*1..3]->(b) SET b.v = 1可能锁住几万个节点。 - 关系也是独立加锁对象。删关系要锁关系本身,而不只是两端节点。
- 索引更新也要加锁。写属性时,属性所在的索引条目会被更新,高并发热点属性(例如全局计数器节点)会造成索引争用。
1.2 MVCC 与读不阻塞写
Neo4j 采用 MVCC(多版本并发控制):读操作不加共享锁,不阻塞写。读事务看到的是一致性快照。这带来两个直接后果:
- 长读事务不会拖慢写,但会阻止旧版本被回收,长读事务会让存储膨胀(旧版本必须保留到没有读事务引用为止)。
- 写写冲突不会被读放大,但写写之间仍然串行化。
| 操作组合 | 是否阻塞 | 说明 |
|---|---|---|
| 读 + 读 | 否 | MVCC 快照 |
| 读 + 写 | 否 | 读不阻塞写 |
| 写 + 写(不同实体) | 否 | 锁粒度细 |
| 写 + 写(同一实体) | 是 | 后到者等待或失败 |
| 写 + 写(同一索引条目) | 是 | 索引级争用 |
1.3 写写冲突的两种结局
同一实体上的并发写有两种处理方式,取决于数据库配置与版本:
策略一:等待(阻塞)
- 后到的事务在锁上排队,直到前者提交
- 风险:等待超时(默认 lock acquisition timeout 常为 10s)
- 症状:DeadlockDetectedException / TransientError
策略二:失败(乐观)
- 后到者立即抛 TransientError(Neo.TransientError.Transaction.DeadlockDetected)
- 由应用层重试
- 风险:重试风暴(高并发下反复撞车)
关键认知:Neo4j 把「死锁检测」也归到 TransientError 这一类,意味着应用层必须区分「可重试错误」与「不可重试错误」。
2. 隔离级别与幻读
Neo4j 实际提供的隔离级别是 Read Committed(读已提交)。它保证:
- 不会读到未提交的数据(无脏读)。
- 同一事务内重复读可能不一致(因为其他事务可能已提交),也就是说
READ COMMITTED下允许不可重复读。
图场景下的一个典型陷阱是**「检查-写入」竞态**:
// 反模式:检查与写入之间没有原子性保证
MATCH (u:User {id: $id})
WHERE u.balance >= $amount
SET u.balance = u.balance - $amount
RETURN u.balance;
两个并发事务可能都通过了 WHERE u.balance >= $amount 检查(因为它们各自读到的是快照),然后都执行扣减,导致余额变负。这类「先读后写」必须在同一事务内完成,且依赖写锁来串行化——Neo4j 的写锁在 SET 时才获取,所以上面的查询确实会被串行化到 SET 那一步;但如果检查逻辑放在应用层(先查余额、再发 SET),就没有任何保护。
正确做法是把检查与写入放在同一条 Cypher 里,让数据库在同一个事务内完成:
MATCH (u:User {id: $id})
WHERE u.balance >= $amount // 条件与写入在同一语句、同一事务
SET u.balance = u.balance - $amount
RETURN u.balance;
如果业务上需要「读出来给用户看、过一会儿再写」,那就必须显式加锁或做版本校验(见第 4 节)。
3. 死锁:图上的高发场景
3.1 死锁是怎么形成的
死锁需要两个条件同时满足:互相持有对方需要的锁 + 加锁顺序不一致。图场景下最容易触发的是「批量更新时遍历顺序不同」:
事务 A:更新 User-1 的所有关系,再更新 User-2 的
事务 B:更新 User-2 的所有关系,再更新 User-1 的
→ A 持有 1 的锁等 2,B 持有 2 的锁等 1 → 死锁
3.2 规避手段
手段一:统一加锁顺序。 所有批量任务按同一个键(例如 id 升序)处理:
// 按 id 升序加锁,避免与并发任务交叉
MATCH (u:User)
WHERE u.id IN $ids
WITH u ORDER BY u.id // 关键:统一顺序
SET u.syncedAt = datetime();
手段二:用 MERGE 或显式锁提前占位。 在事务开始时先锁住关键节点:
// 显式加锁:apoc.lock.nodes 立即获取节点写锁
MATCH (u:User {id: $id})
CALL apoc.lock.nodes([u])
SET u.updatedAt = datetime();
手段三:缩短事务。 事务越短,持锁时间越短,撞车窗口越小。批量任务必须分批(见第 5 节)。
手段四:捕获并重试。 死锁无法完全消除,必须靠重试兜底:
from neo4j.exceptions import TransientError
import time, random
def run_with_retry(session, cypher, params, max_attempts=5):
for attempt in range(max_attempts):
try:
return session.execute_write(lambda tx: tx.run(cypher, params).data())
except TransientError as e:
if attempt == max_attempts - 1:
raise
# 指数退避 + 抖动,避免重试风暴
sleep = min(0.05 * (2 ** attempt), 1.0) * (0.5 + random.random())
time.sleep(sleep)
raise RuntimeError("unreachable")
抖动(jitter)不是可选项:固定退避会让所有失败的事务在同一时刻重试,形成新的撞车高峰。
| 手段 | 适用 | 代价 |
|---|---|---|
| 统一加锁顺序 | 批量任务 | 需业务可排序 |
| 显式加锁 | 热点节点 | 提前持锁,降低并发度 |
| 缩短事务 | 所有场景 | 需改批量逻辑 |
| 捕获重试 | 兜底 | 增加延迟与复杂度 |
4. 乐观重试 vs 显式锁
4.1 乐观并发:版本号校验
给节点加一个版本号,写入时校验版本未变:
MATCH (u:User {id: $id})
WHERE u._version = $expectedVersion // 乐观校验
SET u.balance = $newBalance,
u._version = u._version + 1
RETURN u._version AS newVersion;
如果 RETURN 没有结果,说明版本已被别人改过,应用层需要重新读取并重试。优点是不长时间持锁,适合冲突率低的场景;缺点是冲突率一高就变成重试风暴。
4.2 悲观并发:显式锁
// 用 apoc.lock.nodes 显式获取写锁,后续操作在该事务内独占
MATCH (u:User {id: $id})
CALL apoc.lock.nodes([u])
WITH u
SET u.balance = u.balance - $amount;
显式锁适合冲突率高的热点(例如全局计数器、库存扣减),代价是降低并发度。选择依据:
| 场景 | 冲突率 | 推荐 |
|---|---|---|
| 用户资料更新 | 极低 | 乐观版本号 |
| 库存扣减 | 高 | 悲观显式锁 |
| 计数器累加 | 极高 | 分片计数器(见 4.3) |
| 批量同步 | 中 | 统一加锁顺序 + 重试 |
4.3 热点计数器的分片
单一节点的计数器会成为全局瓶颈。解法是分片:创建 N 个分片节点,写入时随机选一个,读取时求和:
// 写入:随机选分片($shard 由应用层随机生成 0..N-1)
MATCH (c:Counter {name: $name, shard: $shard})
SET c.value = c.value + $delta;
// 读取:求和
MATCH (c:Counter {name: $name})
RETURN sum(c.value) AS total;
分片数 N 的经验取值是「并发写线程数 × 2」,太少起不到分散作用,太多会让求和变慢。这是一个典型的用读放大换写吞吐的权衡。
5. 批量更新的吞吐优化
5.1 批大小与事务边界
批量写入的吞吐与批大小不是单调关系,存在一个峰值:
批太小(如 10):事务开销(日志、提交、锁)占主导 → 吞吐低
批适中(1000~10000):开销摊薄,吞吐接近峰值
批太大(如 100000):事务日志膨胀、锁持有时间长、GC 压力大 → 吞吐反降
推荐从 1000 行/批起步,逐步加倍测吞吐,直到吞吐不再提升或延迟开始恶化。
def batch_upsert(driver, rows, batch_size=1000):
CYPHER = """
UNWIND $rows AS row
MERGE (u:User {id: row.id})
SET u.name = row.name,
u.syncedAt = datetime()
"""
with driver.session() as session:
for i in range(0, len(rows), batch_size):
chunk = rows[i:i + batch_size]
session.execute_write(lambda tx: tx.run(CYPHER, rows=chunk))
5.2 用 UNWIND 而非逐条
UNWIND 把一批数据作为参数一次性送入,同一条查询模板可以复用执行计划:
// 正确:一次网络往返、一个执行计划、一个事务
UNWIND $rows AS row
MERGE (u:User {id: row.id})
SET u.name = row.name;
// 反模式:N 次网络往返、N 个计划(字面量不同)、N 个事务
CREATE (:User {id: 'u1', name: 'a'});
CREATE (:User {id: 'u2', name: 'b'});
吞吐差异通常在一到两个数量级。
5.3 建立关系时避免重复扫描
批量建关系时,两端节点的查找是最贵的部分。先建节点、再建关系,并确保两端都有唯一约束:
// 阶段一:批量建节点
UNWIND $nodes AS n
MERGE (u:User {id: n.id})
SET u.name = n.name;
// 阶段二:批量建关系(两端走唯一索引定位)
UNWIND $edges AS e
MATCH (a:User {id: e.from})
MATCH (b:User {id: e.to})
MERGE (a)-[:FOLLOWS]->(b);
注意 MERGE 关系时如果关系上没有唯一约束,Neo4j 会用「先 MATCH 再 CREATE」的方式模拟,可能产生重复关系。需要幂等就用 MERGE,但要接受它比 CREATE 慢(多一次查找)。
5.4 关闭不必要的开销
批量导入期可以临时调整:
| 配置 | 作用 | 风险 |
|---|---|---|
用 neo4j-admin database import | 离线导入,绕过事务 | 需停机 |
| 关闭约束与索引后导入再重建 | 写入更快 | 失去唯一性保护 |
| 增大事务日志大小 | 支持更大批次 | 磁盘占用 |
用 CALL {} IN TRANSACTIONS | 自动分批 | 需注意不返回聚合 |
CALL { ... } IN TRANSACTIONS OF 1000 ROWS 让数据库自动分批提交,适合「一次性处理大量数据但不想手写分批循环」:
MATCH (u:User)
WHERE u.syncedAt IS NULL
CALL (u) {
SET u.syncedAt = datetime()
} IN TRANSACTIONS OF 1000 ROWS;
6. 集群下的并发与一致性
6.1 写入路径
Neo4j 集群是 leader 写、follower 读(单写多读)模型:
写请求 → 路由到 leader → 提交 → 通过 Raft 复制到 follower → follower 应用后可见
读请求 → 可路由到任意 follower → 可能读到「落后」的数据
这意味着写完立刻读,可能读到旧值。要避免这一点,必须用书签(bookmark)。
6.2 书签一致性
with driver.session() as session:
# 写事务:拿到书签
bookmark = session.execute_write(lambda tx: tx.run(CREATE_CYPHER, p).single())
# 关键:把书签传给下一个读会话,保证读到「至少包含该书签」的状态
with driver.session(bookmarks=[bookmark]) as session:
result = session.execute_read(lambda tx: tx.run(READ_CYPHER, p).data())
书签是因果一致性的实现:它不要求强一致,只要求「不会读到比我自己写之前更旧的状态」。这是分布式图库里最容易被忽略的一环——不加书签的「写后读」在单机测试时永远正确,一上集群就间歇性出错。
6.3 集群下的批量写入
批量任务在集群里要特别注意:
- 只往 leader 写,不要指望 follower 能接受写请求(会被路由或拒绝)。
- 大批量任务会占用 leader 的复制带宽,影响在线写延迟。建议在低峰执行,或限制批速率。
- 重试要重新路由。leader 切换(选举)期间的写请求会失败,重试逻辑要允许重新发现路由表。
7. 幂等批量更新的完整实现
生产级的批量 upsert 需要三个要素:幂等键 + 批量 + 重试。下面是一个完整骨架:
import time, random
from neo4j import GraphDatabase
from neo4j.exceptions import TransientError
UPSERT = """
UNWIND $rows AS row
MERGE (u:User {id: row.id})
ON CREATE SET u.createdAt = datetime()
SET u.name = row.name,
u.email = row.email,
u._syncBatch = $batchId,
u.syncedAt = datetime()
"""
class BatchUpserter:
def __init__(self, uri, auth, batch_size=1000, max_attempts=5):
self.driver = GraphDatabase.driver(uri, auth=auth)
self.batch_size = batch_size
self.max_attempts = max_attempts
def upsert(self, rows, batch_id):
with self.driver.session() as session:
for i in range(0, len(rows), self.batch_size):
chunk = rows[i:i + self.batch_size]
self._write_with_retry(session, chunk, batch_id)
def _write_with_retry(self, session, chunk, batch_id):
for attempt in range(self.max_attempts):
try:
session.execute_write(
lambda tx: tx.run(UPSERT, rows=chunk, batchId=batch_id).consume()
)
return
except TransientError:
if attempt == self.max_attempts - 1:
raise
time.sleep(min(0.05 * 2 ** attempt, 1.0) * (0.5 + random.random()))
三个设计要点:
MERGE依赖唯一约束。没有CREATE CONSTRAINT ... REQUIRE u.id IS UNIQUE,MERGE会退化成全标签扫描,且并发下会产生重复节点。ON CREATE SET与SET分离。创建时间只在首次写入时设置,后续更新不覆盖。_syncBatch标记批次。任务失败后可以按批次定位「哪些行已写、哪些没写」,实现断点续跑。
断点续跑的实现思路:
// 查某个批次已写入的行数
MATCH (u:User {_syncBatch: $batchId})
RETURN count(u) AS done;
// 续跑时跳过已处理的 id 区间(前提:id 有序且批次内按序处理)
8. 监控与排错
8.1 必须盯的指标
1. 事务失败率(按错误类型分组)
- Neo.TransientError.* → 可重试,看重试次数分布
- Neo.ClientError.* → 不可重试,看是否语法/约束问题
2. 锁等待时间
- 持续升高说明存在热点或长事务
3. 活跃事务数与平均事务时长
- 平均时长上升 = 事务边界过大
4. 批量任务吞吐(行/秒)与批延迟 p99
- 吞吐停滞 + 延迟升高 = 批大小过大或锁争用
5. 集群复制延迟
- 影响 follower 读的新鲜度
8.2 定位正在阻塞的事务
// 列出当前所有查询及其状态
CALL dbms.listQueries() YIELD queryId, query, elapsedTime, status, transactionId
WHERE status <> 'completed'
RETURN queryId, elapsedTime, status, substring(query, 0, 80) AS q
ORDER BY elapsedTime DESC;
// 必要时中止阻塞源(谨慎:会回滚其事务)
CALL dbms.killQuery('query-123');
8.3 常见故障对照
| 症状 | 根因 | 处理 |
|---|---|---|
偶发 DeadlockDetected | 加锁顺序不一致 | 统一顺序 + 重试 |
| 批量任务吞吐上不去 | 批太小或未用 UNWIND | 调大批 + 参数化 |
| 写延迟尖刺 | 长事务持锁 | 缩小事务边界 |
| 集群写后读旧值 | 未用书签 | 传 bookmark |
| 出现重复节点 | MERGE 缺唯一约束 | 建约束 + 去重 |
| 计数器写入瓶颈 | 单节点热点 | 分片计数器 |
| 存储持续增长 | 长读事务阻止版本回收 | 缩短读事务 |
小结
图库并发控制的核心认知有三条:锁粒度是节点与关系,遍历会放大冲突面;读不阻塞写(MVCC),但写写之间串行;死锁无法消除,只能靠统一加锁顺序 + 退避重试兜底。批量更新的落地纪律是:一律 UNWIND 参数化、从 1000 行/批起步测吞吐、MERGE 前先建唯一约束、用批次标记实现断点续跑。集群环境额外加一条——写后读必须传书签,否则单机测试永远测不出问题。把这些做对,批量同步任务就能在不停机的前提下与在线流量共存。
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。