多 Agent 协作正在成为 AI 工程的主流范式:一个 Orchestrator 拆解任务,派发给多个 subagent 并行执行。但一旦"并行",就会撞上并发编程的祖传问题——lost update(更新丢失)资源争抢(race condition)。这篇文章讨论两个最典型的场景:多个 Agent 同时修改一个文件,以及多个 Agent 同时调用一个 Tool。

场景一:多 Agent 同时修改一个文件

问题本质

Agent A 和 Agent B 都读取了 config.yml(版本 v1),各自修改后先后写回。后写的会静默覆盖先写的修改——经典的 lost update:

1
2
A: read(v1) ──修改──> write(v2)   ✅ 
B: read(v1) ──修改──> write(v3) ⚠️ 覆盖了 A 的修改,且 A 毫不知情

更麻烦的是,LLM Agent 的写操作通常是全文件重写(Read → 改 → Write),而不是补丁式的,所以冲突面比人类编辑更大。

方案对比

方案 思路 优点 缺点 适用
悲观锁(文件锁) 写前先拿锁,拿不到就排队/失败 简单可靠,绝无冲突 并行度归零,锁泄漏风险 小项目、强一致要求
乐观并发(CAS) 读时记版本,写时校验版本,失败则重试/合并 不阻塞读,天然并行 冲突时需要重试逻辑 冲突率低、读多写少
任务分区(ownership) 每个文件只属于一个 Agent,规划期就隔离 从源头消灭冲突 需要精良的任务拆分 多任务并行重构
三方合并(3-way merge) 像 git merge 一样合并各自改动 冲突时保留双方语义 需要处理冲突块,LLM 可辅助 大规模并行改码

1. 悲观锁:最简单粗暴

写文件前抢锁,本质是把并发写降级为串行:

1
2
3
4
5
6
7
8
9
10
import fcntl

def safe_write(path: str, write_fn):
with open(path, "a+") as f:
fcntl.flock(f, fcntl.LOCK_EX) # 排他锁,拿不到就阻塞
content = f.seek(0) or f.read() # 拿到锁后再读最新内容
new_content = write_fn(content) # Agent 基于最新版本做修改
f.seek(0)
f.truncate()
f.write(new_content)

关键细节:必须在拿到锁之后再读文件。如果先读后抢锁,读到的仍是旧版本,锁形同虚设。

2. 乐观并发(CAS + 版本号)

每个文件带一个版本号(mtime 或 hash),写回时校验"我读的版本还是不是当前版本":

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
import hashlib, os, time

class OptimisticFileStore:
def read(self, path):
content = open(path, "rb").read()
version = hashlib.sha256(content).hexdigest()
return content.decode(), version

def write(self, path, new_content, expected_version, max_retry=3):
for i in range(max_retry):
current = open(path, "rb").read()
if hashlib.sha256(current).hexdigest() == expected_version:
with open(path, "w") as f:
f.write(new_content)
return True
# 版本变了:重读、重新应用 Agent 的编辑、再试
print(f"conflict detected, retry {i+1}/{max_retry}")
expected_version, _ = self._reapply(path, new_content)
raise WriteConflictError(path)

这正是数据库 MVCC、HTTP ETag/If-Match 的思路。冲突率低时性能极好;冲突率高时重试成本上升,那就该上合并了。

3. 任务分区:从源头隔离

在任务规划阶段就把"文件所有权"分配好:Agent 1 负责前端(src/web/),Agent 2 负责后端(src/api/)。Claude Code 的 git worktree 方案更彻底——每个 subagent 一个独立的 worktree,物理隔离,最后由主 Agent 或人来合并分支。这本质上是把并发控制问题转化为合并问题,把难题推迟到有 git 工具链支撑的地方解决。

4. 三方合并

保留共同祖先版本(base),对 A、B 的改动做 diff,非重叠区域自动合并,重叠区域冲突上报。对于语义冲突(比如 A 改了函数签名,B 在别处调用了旧签名),可以再让一个 LLM 做语义层面的冲突仲裁。

场景二:多 Agent 同时调用一个 Tool

先给 Tool 分类

处理并发问题的第一步不是加锁,而是分类——不同危险等级的 Tool 对待方式完全不同:

类型 例子 并发策略
只读(Read-only) 搜索、读文件、查文档 随便并行,多多益善
幂等写(Idempotent) set_config(key, value) 可并行,结果与调用次数无关
非幂等写(Non-idempotent) 发邮件、下单、扣款、创建资源 必须互斥或去重
顺序敏感 cd / git checkout 这类改变后续操作语义的 必须绑定同一会话上下文

方案一:按 Tool 声明并发度

处理并发问题的第一步是给每个 Tool 打标签——在工具注册时就声明它的并发安全等级,调度器据此路由:

等级 语义 示例
PARALLEL 只读,直接并行 search_webread_file
IDEMPOTENT 幂等写,可并行 set_config(key, value)
RESOURCE 资源级互斥,同 key 排队 write_file(按 path)、db_write(按 table)
GLOBAL 全局互斥 send_emaildeploy_service

RESOURCE 是细粒度的关键:write_file 不需要全局锁,只需要同一文件互斥——不同 Agent 写不同文件仍然并行。具体的声明和实现代码见方案二。

方案二:信号量 + 资源键(完整实现)

下面是一个可直接运行的完整实现,覆盖四个关键点:per-key 锁管理、资源键提取、锁超时防死锁、排队诊断日志

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
import asyncio
import time
from enum import Enum


class Safety(Enum):
PARALLEL = "parallel" # 只读,直接并行
IDEMPOTENT = "idempotent" # 幂等写,可并行
RESOURCE = "resource" # 资源级互斥(同 key 排队)
GLOBAL = "global" # 全局互斥


class ToolSpec:
def __init__(self, name, safety, resource_key=None, timeout=30.0):
self.name = name
self.safety = safety
self.resource_key = resource_key # 从 args 提取资源键的函数
self.timeout = timeout # 拿锁超时,防死锁


# ---------- 工具注册:声明并发安全等级 + 资源键提取规则 ----------
TOOLS = {
"search_web": ToolSpec("search_web", Safety.PARALLEL),
"read_file": ToolSpec("read_file", Safety.PARALLEL),
"write_file": ToolSpec("write_file", Safety.RESOURCE,
resource_key=lambda a: a["path"]),
"db_write": ToolSpec("db_write", Safety.RESOURCE,
resource_key=lambda a: f"table:{a['table']}"),
"send_email": ToolSpec("send_email", Safety.GLOBAL, timeout=10.0),
}


class ToolGateway:
"""Agent 与 Tool 之间的并发网关:
同一资源键串行,不同资源键并行;PARALLEL 直通。"""

def __init__(self):
self._locks: dict[str, asyncio.Lock] = {}
self._stats: dict[str, list] = {} # 每个锁的排队记录

def _lock(self, key: str) -> asyncio.Lock:
if key not in self._locks: # asyncio 单线程,无竞态
self._locks[key] = asyncio.Lock()
self._stats[key] = []
return self._locks[key]

async def call(self, agent_id: str, tool: str, args: dict):
spec = TOOLS[tool]
if spec.safety in (Safety.PARALLEL, Safety.IDEMPOTENT):
return await self._exec(tool, args)

key = (spec.resource_key(args) if spec.safety == Safety.RESOURCE
else f"{tool}::global")
lock = self._lock(key)
waiters = self._stats[key]

waiters.append(agent_id)
t0 = time.monotonic()
try:
# 带超时抢锁:避免某个 Tool 卡死导致全局排队
await asyncio.wait_for(lock.acquire(), timeout=spec.timeout)
except asyncio.TimeoutError:
waiters.remove(agent_id)
raise ToolBusyError(
f"{tool}[{key}] 排队超时({spec.timeout}s),"
f"当前持有者未释放。建议 Agent 改用重试策略。")
waited = time.monotonic() - t0
if waited > 0.05: # 只有真排队了才记日志,避免噪音
print(f"[lock] {tool}[{key}] {agent_id} 等待 {waited:.2f}s")
try:
return await self._exec(tool, args)
finally:
waiters.remove(agent_id)
lock.release()
if not waiters and not lock.locked():
del self._locks[key] # 锁清理,防内存泄漏
del self._stats[key]

async def _exec(self, tool, args):
await asyncio.sleep(0.2) # 模拟工具执行耗时
return {"ok": True, "tool": tool, "args": args}


class ToolBusyError(RuntimeError):
pass

用三个 Agent 模拟两种典型局面——写不同文件(应并行)写同一文件(应串行)

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
async def main():
gw = ToolGateway()

# 局面 1:Agent-1 和 Agent-2 写「不同」文件 → 并行,总耗时 ≈ 0.2s
# 局面 2:Agent-3 和 Agent-4 写「同一」文件 → 排队,总耗时 ≈ 0.4s
t0 = time.monotonic()
results = await asyncio.gather(
gw.call("A1", "write_file", {"path": "a.py", "content": "..."}),
gw.call("A2", "write_file", {"path": "b.py", "content": "..."}),
gw.call("A3", "write_file", {"path": "conf.yml", "content": "..."}),
gw.call("A4", "write_file", {"path": "conf.yml", "content": "..."}),
)
print(f"4 次调用总耗时 {time.monotonic() - t0:.2f}s")
# 输出(顺序可能不同):
# [lock] write_file[conf.yml] A4 等待 0.20s
# 4 次调用总耗时 0.40s
# —— a.py/b.py 并行 + conf.yml 排队执行,理论下界正好 0.4s

asyncio.run(main())

三个实现细节值得强调:

  1. asyncio 单线程模型下 _lock()dict 读改写不需要额外加锁——协程只在 await 处让出控制权,__init__if key not in self._locks 之间没有 await,天然原子。如果是多线程(threading)或多进程调度,这行要换成 defaultdict(threading.Lock) 或加全局保护。
  2. wait_for 超时是死锁保险丝:某个 Tool 执行中抛异常不释放、或外部资源 hang 住时,排队方会在超时后拿到明确的 ToolBusyError,Agent 可以退化为"稍后重试"而不是无限挂起。
  3. 锁清理(del self._locks[key])防止长期运行的进程内存泄漏:每个被写过的文件都会留下一把锁,跑几天的 Agent 系统里这会积累几十万个锁对象。

方案三:幂等键(Idempotency Key)

对非幂等的副作用型 Tool(发邮件、下单),光互斥还不够——网络重试也会导致重复执行。业界标准做法是幂等键:

1
2
3
4
5
6
7
8
9
import uuid

async def call_with_idempotency(tool, args, user_intent_hash):
key = f"{tool}:{user_intent_hash}" # 相同意图 → 相同 key
if result := await cache.get(key):
return result # 已执行过,直接返回上次结果
result = await execute(tool, args, headers={"Idempotency-Key": key})
await cache.set(key, result, ttl=3600)
return result

Stripe 的支付 API 就是这个模式:同一个 Idempotency-Key 的请求无论重发多少次,只扣一次款。Agent 重试风暴下这是唯一可靠的防线。

方案四:共享状态外置

还有一种釜底抽薪的思路:工具服务端自己保证并发安全,Agent 侧完全不管。比如所有 Tool 都是数据库操作,靠数据库的事务与行锁保证一致性;或者 Tool 是无状态的纯 HTTP 服务,靠服务端幂等。Agent 编排层只做"尽量并行 + 冲突重试",正确性交给有 ACID 保障的底层。

现实工程怎么做的

几个可参考的生产级实现:

  • Claude Code:subagent 默认物理隔离(独立 worktree/独立进程),主 Agent 充当合并点;写文件类操作在 UI 上有权限确认,人为充当"最终一致性的仲裁者"。
  • LangGraph:用图的状态结构显式建模共享状态,State 的每个 key 可以声明 Reducer(如 add 追加、overwrite 覆盖),并发更新按 Reducer 规则归并——本质是 CRDT 思想。
  • OpenAI Swarm / Handoff 模式:干脆限制同一时刻只有一个活跃 Agent,用 handoff 传递控制权,用"避免并发"来消灭并发问题。

总结:选型决策树

1
2
3
4
5
6
7
8
9
改动可分区吗(不同 Agent 不碰同一批文件/资源)?
├─ 可以 → 任务分区 + 最后合并(worktree/分支),最优解
└─ 不行 → 冲突频率高吗?
├─ 高 → 悲观锁/互斥队列,用并行度换正确性
└─ 低 → 乐观并发(版本号 + 冲突重试)
└─ 冲突了 → 三方合并 / LLM 语义仲裁
Tool 是副作用型吗?
├─ 只读/幂等 → 放心并行
└─ 非幂等 → 互斥 + 幂等键,缺一不可

一句话:并发问题没有银弹,但有成熟的分级处理框架——先分类(只读/幂等/互斥),再选粒度(全局锁/资源锁/无锁),最后给冲突兜底(重试/合并/人工确认)。这和数据库并发控制走的是同一条演化路径,只不过"事务"从 SQL 变成了自然语言驱动的 Tool 调用。