Agent 工程体系 · 第 12/98 篇。内容以 2026 年 9 月可验证的公开规范和稳定接口为基线;框架版本敏感能力会明确标注,不把实验行为写成通用保证。
Agent 并行工具调用:依赖图、只读并发、写入串行和合并
Agent 的一次行动,通常不是“模型调用一个工具,等待结果,再调用下一个工具”这么简单。模型可能在同一轮输出多个工具调用;执行器则必须判断这些调用之间是否存在数据依赖、资源冲突或业务顺序约束,然后决定哪些可以并行,哪些必须串行,哪些结果可以合并,哪些失败会使整个计划失效。
这一区分很重要:
- 模型决定要调用什么工具;
- 执行器决定以什么顺序、以什么并发度执行;
- 工具服务负责真正的外部副作用和一致性约束;
- 合并器负责把多个结果重新组织成下一轮模型可理解的事实。
不能因为模型在一次响应中返回了多个工具调用,就直接对它们执行 Promise.all 或 asyncio.gather。并行调用只是一个候选执行集合,不等于这些操作在语义上彼此独立。
OpenAI 的 Function Calling 文档把工具调用描述为模型向应用发出的工具请求,并允许模型在一个回合中返回多个函数调用;在支持的模型上,可以使用 parallel_tool_calls: false 强制每轮最多执行一个工具调用。这个参数只能控制模型输出形式,不能替代执行器对依赖和冲突的判断。(developers.openai.com)
MCP 则定义了模型应用与外部工具、数据源之间的标准化连接方式。MCP 规定工具通过 tools/list 发现、通过 tools/call 调用,并用 JSON Schema 描述输入和可选的输出结构;它并不替应用决定多个工具调用之间的调度策略。(modelcontextprotocol.io)
一、先区分四个容易混淆的概念
1. 工具调用不是工具执行
一个工具调用可以抽象为:
其中:
- :本次调用的唯一标识;
- :工具名称;
- :模型生成的参数;
- :会话、用户、授权、租户、追踪信息等执行上下文。
工具执行则是:
执行结果还需要带上状态:
因此,模型给出:
[
{
"id": "call_weather_hz",
"name": "get_weather",
"arguments": { "city": "杭州" }
},
{
"id": "call_calendar",
"name": "list_calendar",
"arguments": { "date": "2026-09-01" }
}
]
只说明模型请求了两个工具,不说明:
- 两个工具是否能同时访问;
- 第二个工具是否依赖第一个工具的结果;
- 两个工具是否会修改相同资源;
- 一个失败后是否应该取消另一个;
- 两个结果是否可以直接拼接给模型。
执行器需要在“调用”与“执行”之间增加一个计划分析阶段。
2. 并发不等于并行
并发描述的是多个任务存在重叠执行的可能;并行通常表示它们在同一时间真正占用不同的执行资源。
在 Agent 系统中,工程上通常不需要区分到 CPU 指令级别。只需要保证:
- 允许并发的任务可以同时处于
running; - 受限任务不能同时执行;
- 所有依赖满足后,任务才可以进入
ready; - 最终结果顺序不能依赖网络返回顺序。
因此,本文使用“并行工具调用”表示执行器允许多个互不冲突的调用重叠执行,但具体是线程、协程、进程还是远程服务并行,由运行时决定。
3. 只读不是“函数名里有 read”
只读工具的关键不是名称,而是可观察副作用。
一个工具只有在以下条件都成立时,才能被视为只读:
- 不修改数据库、缓存、文件、队列、外部 API 状态;
- 不触发发送消息、计费、下单、审批等隐式动作;
- 不依赖会被其他并发调用修改的共享状态;
- 失败、重试、超时不会改变外部系统状态;
- 读取结果可以接受其一致性级别,例如允许读取副本或稍旧快照。
例如:
get_user_profile(user_id) 通常是只读
search_orders(user_id) 通常是只读
send_email(to, body) 写入
create_payment_intent(amount) 写入,即使接口名称不是 save
reserve_inventory(sku, count) 写入
get_next_sequence() 可能是写入,若会消耗序列号
get_next_sequence() 是一个常见反例。它看起来像读取,但如果每次调用都会递增数据库序列,那么并发执行会消耗多个号码,重试还可能产生额外缺口。
4. 写入串行不是所有写入都必须全局排队
“写入串行”通常不是指整个系统只能同时执行一个写工具,而是指:
对同一个冲突域内的写操作,必须按照定义好的顺序提交。
可以定义多个冲突域:
tenant:acme
account:user_123
order:order_456
inventory:sku_ABC
global:payment_provider
于是:
- 修改不同订单的写入可以并发;
- 修改同一订单的写入需要串行;
- 修改同一库存 SKU 的扣减通常需要串行或由数据库原子操作协调;
- 涉及全局支付额度的操作可能需要更宽的锁域。
如果执行器只设计一个全局写锁,正确性简单,但吞吐量会很差;如果完全不设锁,吞吐量高,却可能产生竞态。冲突域是两者之间的结构化折中。
二、把工具调用建模为依赖图
1. 图的基本定义
将一次 Agent 行动中的工具调用建模为有向图:
其中:
- :工具调用节点集合;
- :依赖边集合;
- :表示调用 必须等待调用 的某个结果或提交状态。
例如,用户要求:
查询订单详情和物流状态,如果订单存在且未签收,再发送一条提醒。
可以拆成:
A: 查询订单详情
B: 查询物流状态
C: 判断订单是否满足提醒条件
D: 发送提醒
依赖关系是:
A ─┐
├──> C ───> D
B ─┘
对应边集合:
因为 A 和 B 都只是查询,并且互不依赖,所以它们可以并行。C 必须等待 A、B。D 必须等待 C 的判断结果。
2. 依赖边与资源边
依赖图中的边不只有一种来源。
数据依赖
调用 B 需要调用 A 的输出:
A: find_customer(email) -> customer_id
B: list_orders(customer_id)
边为:
写后读依赖
调用 A 修改数据,调用 B 读取修改后的数据:
A: update_order_status(order_id, "shipped")
B: get_order(order_id)
如果 B 要观察 A 的新状态,就必须有:
否则 B 可能读到旧值。
写写冲突
两个写调用修改同一资源:
A: set_order_status(order_id, "paid")
B: set_order_status(order_id, "cancelled")
它们可能没有数据依赖,但存在写写冲突。执行器必须:
- 依据业务顺序建立边;
- 或拒绝同时执行;
- 或交给数据库事务、版本号、条件更新解决。
如果确定业务顺序是先支付、后取消,则加入:
如果顺序不确定,就不能随意选择一个顺序并声称结果正确。此时应该返回冲突,让规划器或用户决定。
读写冲突
一个调用读取资源,另一个调用修改相同资源:
A: get_account_balance(account_id)
B: withdraw(account_id, 100)
若 A 的结果用于决定 B 是否可以执行,则是数据依赖:
若 A 只是生成日志或展示信息,则可以在某些一致性要求下并行,但读取结果不应被解释为扣款前后的确定快照。
3. 为什么依赖图必须是 DAG
调度器通常需要一个有向无环图,即 DAG(Directed Acyclic Graph,有向无环图)。
如果出现:
A -> B -> C -> A
则不存在满足所有依赖的第一个节点,拓扑排序失败。循环依赖常见于:
- 工具 A 等待工具 B 返回资源 ID;
- 工具 B 又要求工具 A 先创建确认状态;
- 模型把两个互相需要的动作同时放进计划;
- 执行器错误地把“结果合并”反向当成了工具依赖。
调度器必须在执行前检测环,而不是让任务全部进入等待状态后才发现死锁。
三、从依赖图推导可执行批次
1. 入度决定任务是否就绪
对每个节点 ,定义入度:
入度为 0 的节点没有未完成前置依赖,可以进入 ready 状态。
执行过程如下:
- 计算所有节点的入度;
- 将入度为 0 的节点放入就绪队列;
- 取出满足资源和授权条件的节点执行;
- 节点成功完成后,删除其出边;
- 后继节点入度降为 0 时,加入就绪队列;
- 重复直到全部完成或无法继续。
伪代码:
def topological_batches(nodes, edges):
outgoing = {node: [] for node in nodes}
indegree = {node: 0 for node in nodes}
for before, after in edges:
outgoing[before].append(after)
indegree[after] += 1
ready = sorted(node for node in nodes if indegree[node] == 0)
batches = []
visited = 0
while ready:
batch = ready
batches.append(batch)
next_ready = []
for node in batch:
visited += 1
for successor in outgoing[node]:
indegree[successor] -= 1
if indegree[successor] == 0:
next_ready.append(successor)
ready = sorted(next_ready)
if visited != len(nodes):
raise ValueError("dependency graph contains a cycle")
return batches
输入:
nodes = {"A", "B", "C", "D"}
edges = {
("A", "C"),
("B", "C"),
("C", "D"),
}
输出:
[
["A", "B"],
["C"],
["D"],
]
第一批次中的 A、B 可以并发;C 只能等待两者都完成;D 只能等待 C。
2. 拓扑层不是最终并发批次
上面的分层只体现依赖关系,还没有体现:
- 只读/写入属性;
- 冲突域;
- 并发配额;
- 工具级限流;
- 用户授权;
- 超时和取消;
- 优先级与公平性。
例如:
A: read(account:1)
B: read(account:1)
C: write(account:1)
D: write(account:2)
即使四个节点都没有显式数据依赖,也不能简单地把它们全部放入一个批次。更合理的执行策略可能是:
批次 1:A、B、D
批次 2:C
这里 D 与 C 的资源不同,可以并发;A、B 是只读,可以并发;C 是写入 account:1,需要避开对该资源的其他操作。
3. 调度条件的形式化
一个节点 可以在时刻 启动,当且仅当:
各项含义是:
- :所有前置节点已经完成并满足该节点要求;
- :当前用户、Agent、租户和工具权限允许执行;
- :会话、租户、工具和全局资源配额允许执行;
- :没有与运行中任务发生不允许的资源冲突;
- :启动后仍有足够时间完成或达到可接受的截止时间。
只要其中一项不满足,节点就应保持 blocked 或 queued,而不是强行执行。
四、只读并发:安全条件与边界
1. 只读调用可并行的充分条件
设两个调用 都是只读。若满足:
并且它们不依赖彼此的结果:
同时读取的一致性要求允许它们观察不同时间点的数据,则可以并行:
这个条件是工程上的安全近似,而不是数据库理论中的唯一条件。
例如:
get_user_profile(user_id=123)
list_recent_orders(user_id=123)
二者都读取数据,且一个不需要另一个的输出,可以并发。
但如果业务要求“用户资料和订单必须来自同一个数据库快照”,那么仅仅都是只读还不够。它们需要:
- 共享同一个事务快照;
- 或使用同一个版本号;
- 或由后端提供聚合查询;
- 或接受一致性声明中的“可能跨时间点”。
2. 只读并发的时间收益
设三个独立只读工具耗时分别为:
串行执行的理想耗时是:
并行执行的理想耗时接近:
实际耗时还要加上调度、连接、限流、序列化和结果合并开销:
因此,只有在工具耗时足够大、并发资源可用、服务端不会因突发并发而显著排队时,并行才有实际收益。
3. 只读工具也可能造成副作用
以下实现不能被简单归类为只读:
def get_report(report_id):
report = db.query(report_id)
db.execute(
"INSERT INTO report_access_log(report_id, accessed_at) VALUES (?, now())",
report_id,
)
return report
它虽然返回查询结果,但会写访问日志。如果日志写入:
- 需要顺序;
- 可能触发审计告警;
- 与删除操作存在竞态;
- 使用共享计数器;
那么它就不是纯只读工具。
另一个常见问题是缓存:
get_user_profile()
如果缓存未命中时会回填共享缓存,严格来说它包含写操作。实践中可以把这类工具标记为:
read_with_cache_write
并允许它与其他读取并发,但不能假设它与所有写入都无冲突。
4. 只读并发与数据快照
考虑初始库存为 10:
R: get_inventory() -> 10
W: reserve_inventory(7)
如果 R 与 W 并发,R 可能返回 10,也可能返回 3,取决于数据库读取时机。两个结果都可能是数据库层面合法的,但 Agent 不能把 R 的结果解释成“预留操作前的库存”或“预留操作后的库存”,除非系统明确规定了快照关系。
需要稳定观察结果时,可以把版本号纳入工具接口:
{
"name": "get_inventory",
"arguments": {
"sku": "A-100",
"as_of_version": 42
}
}
或者让写操作返回:
{
"sku": "A-100",
"reserved": 7,
"remaining": 3,
"version": 43
}
这样后续读取可以明确请求版本 43,而不是依赖不可靠的时间顺序。
五、写入串行:顺序、冲突域和提交点
1. 写入为什么需要顺序
设两个操作都读取旧值 ,再写回计算结果:
A: x = x + 10
B: x = x + 20
初始:
若 A、B 并发且都执行“读—计算—写”:
A 读到 100
B 读到 100
A 写入 110
B 写入 120
最终结果是 120,而正确的累加结果应为:
这就是丢失更新。
如果串行执行:
A 读 100,写 110
B 读 110,写 130
则得到正确结果 130。
但如果数据库使用原子更新:
UPDATE account
SET balance = balance + :delta
WHERE account_id = :account_id;
那么两个操作可以由数据库并发协调,执行器不一定需要把它们完全串行化。此时正确性来自数据库的原子语义,而不是 Agent 调度器的顺序。
因此,“写入串行”应理解为:
在没有更强一致性机制证明并发安全之前,对同一冲突域的写操作按顺序提交。
2. 冲突域键
为每个工具调用声明资源访问集合:
ToolCall(
id="reserve_1",
mode="write",
resources={("inventory", "sku:A-100")},
)
ToolCall(
id="update_order",
mode="write",
resources={("order", "order:9001")},
)
两个调用是否冲突,可以定义为:
或:
对于两个写调用,还可以采用更严格的判断:
工程上通常不把完整的读写集合交给模型自行填写,而是由工具注册表提供。例如:
{
"name": "reserve_inventory",
"execution": {
"mode": "write",
"resource_keys": ["inventory:{sku}"],
"commutativity": "non_commutative",
"idempotency": "required"
}
}
这里的 resource_keys 是执行器元数据,不应仅依赖工具描述文本。MCP 的工具描述包括名称、描述和输入 Schema,也允许工具提供行为注解;但规范明确要求客户端把工具注解视为不可信,除非工具来自可信服务器。(modelcontextprotocol.io)
3. 写入顺序来自哪里
写入顺序可以有四种来源:
规划顺序
模型或上层规划器明确生成:
先创建订单,再扣库存,再发送确认邮件
执行器将其转成:
资源顺序
多个任务操作同一个资源时,执行器依据提交序列号排序:
order:9001:
seq=17 -> update_status("paid")
seq=18 -> update_status("shipped")
业务状态机
工具只能接受合法状态迁移:
pending -> paid -> shipped -> delivered
即使执行器把 shipped 放在 paid 前面,服务端也应拒绝非法迁移,而不能依赖调用方永远正确排序。
条件写入
使用版本号避免覆盖并发更新:
UPDATE orders
SET status = :new_status,
version = version + 1
WHERE order_id = :order_id
AND version = :expected_version;
受影响行数为 0 时,说明版本冲突,工具应返回可恢复错误,而不是伪装成成功。
4. 串行执行不等于串行生成
模型可以一次生成:
A: update_profile
B: update_preferences
C: send_confirmation
执行器可以先并行完成 A、B,再串行执行 C:
A ─┐
├──> C
B ─┘
所以:
- 生成阶段可以一次产生多个调用;
- 调度阶段可以并发执行独立节点;
- 提交阶段必须遵守资源和业务顺序;
- 结果阶段必须恢复稳定顺序。
六、结果合并:不是简单拼接字符串
1. 合并的输入和目标
合并器接收多个工具结果:
输出一个供 Agent 状态机消费的结构:
其中:
facts:可供模型使用的事实;errors:失败、超时、取消和冲突;effects:哪些写入已提交;provenance:每条事实来自哪个调用;next_actions:是否需要重试、补偿、澄清或继续规划。
一个不可靠的合并方式是:
context += tool_result.text
它会丢失:
- 调用 ID;
- 成功或失败状态;
- 结果顺序;
- 结构化字段;
- 是否发生副作用;
- 是否只是部分完成。
2. 使用稳定的结果信封
建议把每个结果包装成统一信封:
{
"call_id": "reserve_1",
"tool": "reserve_inventory",
"status": "succeeded",
"committed": true,
"resource_keys": ["inventory:sku:A-100"],
"output": {
"sku": "A-100",
"reserved": 2,
"remaining": 8
},
"error": null,
"started_at": "2026-09-01T10:00:01.120Z",
"finished_at": "2026-09-01T10:00:01.410Z"
}
失败结果:
{
"call_id": "send_email_1",
"tool": "send_email",
"status": "failed",
"committed": "unknown",
"resource_keys": ["mail:user:123"],
"output": null,
"error": {
"kind": "timeout",
"retryable": true,
"message": "provider response was not received before deadline"
}
}
committed: "unknown" 很重要。网络超时只说明执行器没有收到确认,不说明外部写入一定没有发生。对于发送邮件、扣款、下单等操作,不能看到超时就盲目重试。
MCP 将错误分为协议错误和工具执行错误:未知工具、格式错误等属于 JSON-RPC 协议错误;业务失败、输入校验失败和外部 API 失败则应作为工具结果中的 isError: true 返回,使模型能够理解和修正。(modelcontextprotocol.io)
3. 结果顺序必须稳定
假设并行任务完成顺序是:
B 完成
A 完成
C 完成
不能直接按照完成顺序把结果送回模型,因为模型上下文会出现非确定性:
[B, A, C]
下一次可能变成:
[A, C, B]
应按照以下优先级之一排序:
- 计划中的节点序号;
- 拓扑序;
- 工具调用在模型响应中的索引;
- 稳定的
call_id排序。
例如:
def stable_merge(results, plan_order):
position = {call_id: i for i, call_id in enumerate(plan_order)}
return sorted(results, key=lambda item: position[item["call_id"]])
这不会改变工具实际完成顺序,只会使模型看到的结果顺序稳定。
4. 结构化合并与冲突
假设两个只读工具返回:
A: { "user_id": "u1", "name": "张三", "plan": "basic" }
B: { "user_id": "u1", "name": "张三", "plan": "pro" }
不能使用“后写覆盖前写”的字典合并:
merged.update(result)
因为这会静默丢失冲突。应保留来源并报告不一致:
{
"user_id": "u1",
"name": {
"value": "张三",
"sources": ["A", "B"]
},
"plan": {
"conflict": true,
"values": [
{ "value": "basic", "source": "A" },
{ "value": "pro", "source": "B" }
]
}
}
合并策略必须按字段类型定义:
- 集合:可做并集,但要去重;
- 列表:通常按业务排序键合并;
- 计数:只有确认各结果是互斥分片时才能求和;
- 快照:必须保留版本和读取时间;
- 状态:不能任意覆盖,应按状态机或优先级合并;
- 写入结果:按提交序列合并,并保留失败和未知状态。
七、一个完整算例:查询并发、写入串行、结果合并
需求:
查询订单详情、查询库存和查询用户偏好;如果订单属于用户本人、库存足够且用户允许通知,则预留库存并发送确认消息。
拆分工具:
A = get_order(order_id)
B = get_inventory(sku)
C = get_user_preferences(user_id)
D = reserve_inventory(sku, quantity)
E = send_confirmation(user_id, order_id)
依赖关系:
A ─┐
B ─┼──> D ───> E
C ─┘
A ───────────> E
C ───────────> E
更具体地说:
- D 需要 A 提供 SKU 和数量,B 提供当前库存;
- D 需要 C 判断是否允许执行相关业务;
- E 需要 D 返回预留成功;
- E 还需要 A、C 的订单和通知信息。
图可以表示为:
flowchart LR
A[get_order] --> D[reserve_inventory]
B[get_inventory] --> D
C[get_user_preferences] --> D
D --> E[send_confirmation]
A --> E
C --> E
第一步:初始化节点
[
{ "id": "A", "mode": "read", "resources": ["order:9001"] },
{ "id": "B", "mode": "read", "resources": ["inventory:A-100"] },
{ "id": "C", "mode": "read", "resources": ["user:u1"] },
{ "id": "D", "mode": "write", "resources": ["inventory:A-100"] },
{ "id": "E", "mode": "write", "resources": ["mail:user:u1"] }
]
A、B、C 的入度均为 0,因此进入第一批次:
ready = [A, B, C]
它们是只读,且资源互不产生不允许的冲突,可以并行。
第二步:并行执行只读节点
假设返回:
A -> order_id=9001, user_id=u1, sku=A-100, quantity=2, owner=u1
B -> sku=A-100, available=5
C -> user_id=u1, notification_allowed=true
合并器不直接生成自然语言,而是生成内部事实:
{
"order": {
"id": "9001",
"owner": "u1",
"sku": "A-100",
"quantity": 2
},
"inventory": {
"sku": "A-100",
"available": 5
},
"preferences": {
"notification_allowed": true
},
"sources": {
"order": "A",
"inventory": "B",
"preferences": "C"
}
}
第三步:执行写入 D
校验:
D 获得执行资格。由于 D 修改 inventory:A-100,同一冲突域的其他写任务必须等待。
工具最好使用幂等键:
{
"sku": "A-100",
"quantity": 2,
"reservation_id": "agent-run-7:reserve-D"
}
若执行器因超时重试,服务端应根据 reservation_id 返回原来的预留结果,而不是再次扣减库存。
假设 D 返回:
{
"reservation_id": "res-7788",
"sku": "A-100",
"reserved": 2,
"remaining": 3,
"committed": true
}
第四步:执行通知 E
E 只能在 D 确认提交后执行。因为“库存预留成功”是发送确认消息的前置条件:
如果 D 失败,E 不应发送“预留成功”的确认消息。
如果 E 超时:
D 已提交
E 状态未知
此时不能回滚库存并简单重试,除非系统确实支持可靠补偿。更安全的处理是:
- 使用相同幂等键重试发送;
- 查询消息服务的发送状态;
- 如果无法确认,写入待处理任务;
- 向模型报告“库存已预留,通知状态未知”。
最终合并结果:
{
"status": "partially_succeeded",
"committed_effects": [
{
"type": "inventory_reservation",
"reservation_id": "res-7788",
"sku": "A-100",
"quantity": 2
}
],
"uncertain_effects": [
{
"type": "confirmation_message",
"idempotency_key": "agent-run-7:notify-E"
}
],
"next_action": "check_or_retry_notification"
}
这比返回一句“操作失败”准确得多,因为操作已经产生了部分不可逆副作用。
八、执行器的生命周期与状态机
一个工具调用至少需要区分以下状态:
planned
-> blocked
-> ready
-> queued
-> running
-> succeeded
-> failed
-> timed_out
-> cancelled
-> unknown
状态转换不能任意发生。例如:
planned -> running:通常不允许,必须先经过授权、依赖和资源检查;running -> cancelled:只有底层工具支持取消,才能认为执行已停止;running -> timed_out:只表示本地等待超时,不代表远端动作未发生;timed_out -> retrying:必须检查幂等性和副作用状态;succeeded -> compensated:表示成功之后执行了补偿,不应覆盖原始成功事实。
推荐把本地超时和远端提交状态分开:
{
"local_status": "timed_out",
"remote_commit_status": "unknown",
"retry_policy": "query_status_before_retry"
}
MCP 规范提供取消、进度和错误报告等通用能力;长时间运行操作还可以使用可选的 Tasks 扩展,但这些协议能力并不会自动解决 Agent 层的依赖图、写入顺序或业务补偿。MCP 的 Tasks 扩展属于需要双方显式协商的可选扩展。(modelcontextprotocol.io)
九、一个可运行的简化调度器
下面的示例使用 Python 标准库,展示:
- DAG 拓扑调度;
- 只读任务并发;
- 同一资源的写入互斥;
- 结果按计划顺序合并;
- 任意节点失败后阻止依赖它的后继节点。
它没有实现生产级授权、重试、分布式锁和远端幂等,只用于说明核心机制。
from __future__ import annotations
import asyncio
from dataclasses import dataclass, field
from typing import Any, Callable
@dataclass(frozen=True)
class Call:
call_id: str
mode: str # "read" or "write"
resources: frozenset[str]
run: Callable[[], Any]
@dataclass
class Result:
call_id: str
status: str
output: Any = None
error: str | None = None
class Executor:
def __init__(self, calls: dict[str, Call], edges: set[tuple[str, str]]):
self.calls = calls
self.edges = edges
self.outgoing = {call_id: set() for call_id in calls}
self.remaining = {call_id: set() for call_id in calls}
for before, after in edges:
self.outgoing[before].add(after)
self.remaining[after].add(before)
self.results: dict[str, Result] = {}
self.active_resources: set[str] = set()
self.resource_lock = asyncio.Lock()
async def can_start(self, call: Call) -> bool:
async with self.resource_lock:
if call.mode == "read":
# 简化规则:只读任务不与写资源重叠。
# 生产环境可使用读写锁和更细粒度的冲突矩阵。
return not (call.resources & self.active_resources)
return not (call.resources & self.active_resources)
async def reserve_resources(self, call: Call) -> None:
async with self.resource_lock:
self.active_resources.update(call.resources)
async def release_resources(self, call: Call) -> None:
async with self.resource_lock:
self.active_resources.difference_update(call.resources)
async def run_one(self, call: Call) -> Result:
await self.reserve_resources(call)
try:
value = await call.run()
return Result(call.call_id, "succeeded", output=value)
except Exception as exc:
return Result(call.call_id, "failed", error=str(exc))
finally:
await self.release_resources(call)
async def run(self) -> list[Result]:
pending = set(self.calls)
while pending:
ready = [
call_id for call_id in sorted(pending)
if all(
dep in self.results and self.results[dep].status == "succeeded"
for dep in self.remaining[call_id]
)
]
# 如果某个依赖失败,后继节点不会被执行。
blocked = [
call_id for call_id in sorted(pending)
if any(
dep in self.results and self.results[dep].status != "succeeded"
for dep in self.remaining[call_id]
)
]
for call_id in blocked:
self.results[call_id] = Result(
call_id,
"cancelled",
error="blocked by failed dependency",
)
pending.remove(call_id)
runnable = []
for call_id in ready:
call = self.calls[call_id]
if await self.can_start(call):
runnable.append(call)
if not runnable:
if pending:
raise RuntimeError(
"no runnable task: cycle or resource scheduling deadlock"
)
break
# 同一轮启动所有当前无冲突的任务。
batch_results = await asyncio.gather(
*(self.run_one(call) for call in runnable)
)
for result in batch_results:
self.results[result.call_id] = result
pending.remove(result.call_id)
# 按调用定义的稳定顺序输出,而不是按完成顺序输出。
return [self.results[call_id] for call_id in self.calls]
示例工具:
async def get_order():
await asyncio.sleep(0.20)
return {"order_id": "9001", "sku": "A-100", "quantity": 2}
async def get_inventory():
await asyncio.sleep(0.10)
return {"sku": "A-100", "available": 5}
async def reserve_inventory():
await asyncio.sleep(0.15)
return {"reservation_id": "res-7788", "reserved": 2}
calls = {
"A": Call(
call_id="A",
mode="read",
resources=frozenset({"order:9001"}),
run=get_order,
),
"B": Call(
call_id="B",
mode="read",
resources=frozenset({"inventory:A-100"}),
run=get_inventory,
),
"D": Call(
call_id="D",
mode="write",
resources=frozenset({"inventory:A-100"}),
run=reserve_inventory,
),
}
edges = {
("A", "D"),
("B", "D"),
}
results = asyncio.run(Executor(calls, edges).run())
for result in results:
print(result)
预期输出顺序是:
Result(call_id='A', status='succeeded', output=...)
Result(call_id='B', status='succeeded', output=...)
Result(call_id='D', status='succeeded', output=...)
虽然 B 可能比 A 更早完成,但结果仍按 A、B、D 的计划顺序返回。D 不能在 A、B 完成前启动;A、B 由于是不同资源上的只读操作,可以重叠执行。
这个示例有一个重要限制:active_resources 使用简单集合,因此读操作之间也会互相阻塞同一资源。生产系统通常需要显式的读写锁:
- 多个只读调用可以共享读锁;
- 写入调用需要独占写锁;
- 若写入需要读取并基于读取结果提交,则应升级为事务或条件写入;
- 锁的粒度应与冲突域一致,而不是默认全局锁。
十、与 Function Calling 和 MCP 的衔接
1. Function Calling 层
模型输出的工具调用通常包含:
tool_call_id
function.name
function.arguments
流式响应时,工具调用参数可能分散在多个增量事件中,执行器必须先按调用索引聚合参数,再进行 JSON 解码和 Schema 校验,不能在收到第一个片段时执行工具。OpenAI 文档也展示了需要按工具调用索引累积 arguments 的处理方式。(developers.openai.com)
推荐流水线:
模型响应
-> 聚合流式 tool call
-> 严格 JSON 解码
-> Schema 校验
-> 工具注册表查找
-> 授权检查
-> 依赖提取
-> 冲突分析
-> 调度执行
-> 结果信封
-> 按调用 ID 回填工具结果
-> 下一轮模型请求
启用严格 Schema 有助于把参数形状错误挡在执行前。OpenAI 文档说明,strict: true 要求对象 Schema 使用 additionalProperties: false,并要求 properties 中的字段标记为 required;可选值可以通过允许 null 表示。(developers.openai.com)
但严格 Schema 只保证参数结构更可靠,不保证:
- 参数代表的资源真实存在;
- 用户有权访问该资源;
- 工具可以安全重试;
- 两个调用没有业务冲突;
- 工具结果与其他读取处于同一快照。
2. MCP 层
MCP 的典型生命周期是:
初始化连接
-> 协商能力
-> tools/list
-> 缓存工具定义
-> 生成或暴露给模型的工具集合
-> tools/call
-> 返回结构化或非结构化结果
MCP 工具通过 tools/list 暴露 inputSchema,并可以提供 outputSchema。如果服务端声明了输出 Schema,服务端必须返回符合该 Schema 的结构化结果,客户端应进行验证。(modelcontextprotocol.io)
MCP 工具列表可能随时间变化;服务端可以声明 listChanged,并发送 notifications/tools/list_changed 通知。执行器因此不能无限期缓存工具权限和工具元数据,至少应在工具列表变化、授权变化或连接初始化时重新确认。(modelcontextprotocol.io)
MCP 对状态工具的建议也与并发调度直接相关:协议层没有隐式的跨调用会话状态,需要延续状态时,应通过显式句柄传递,例如 basket_id、transaction_id 或 reservation_id。服务端必须在每次调用时重新校验句柄授权,不能因为句柄曾经由某个调用创建,就默认后续调用永远有权使用。(modelcontextprotocol.io)
十一、失败路径:部分成功比整体失败更危险
并行执行带来一个串行执行没有的情况:
A 成功
B 成功
C 失败
D 仍在执行
此时“取消整个批次”未必能撤销 A、B,也未必能停止 D。系统必须区分:
- 未开始:可以取消;
- 正在运行:只能发出取消请求,不能假设取消成功;
- 已成功且无副作用:可以忽略或重新计算;
- 已提交副作用:需要幂等查询、回滚或补偿;
- 状态未知:必须先确认远端状态。
例如:
reserve_inventory 成功
send_email 超时
错误的处理是:
认为整个流程失败
重新执行 reserve_inventory
重新执行 send_email
这可能导致库存被预留两次,邮件发送两次。
正确的处理是把流程拆成两个独立事实:
库存预留:已提交
确认消息:状态未知
然后对 E 使用幂等键或状态查询:
notify:user:u1:order:9001:reservation:res-7788
如果工具不支持幂等,也无法查询状态,那么执行器不应自动重试不可逆写入,只能将状态交给人工或专门的恢复流程。
失败传播规则
对依赖图中的节点 ,可以定义:
如果某个前置节点失败,则后继节点默认进入 blocked,但并非所有失败都必须终止整张图:
- 关键数据查询失败:阻止依赖它的写入;
- 非关键推荐查询失败:可以继续主流程;
- 审计记录失败:可能需要阻止高风险操作;
- 通知失败:主业务可能成功,但整体状态是部分成功。
因此,边还可以附带失败策略:
{
"from": "C",
"to": "D",
"failure_policy": "block"
}
{
"from": "D",
"to": "E",
"failure_policy": "block"
}
{
"from": "A",
"to": "optional_analytics",
"failure_policy": "continue"
}
十二、常见错误与诊断方法
错误一:看到多个 tool call 就全部并发
表现:
- 订单状态偶尔倒退;
- 库存出现负数或重复预留;
- 同一个用户收到重复通知;
- 测试稳定通过,生产偶发失败。
诊断:
记录每个调用的:
run_id
call_id
tool_name
resource_keys
mode
dependency_ids
start_time
finish_time
commit_status
idempotency_key
如果两个同时运行的调用拥有交集资源键,且至少一个是写入,就应检查冲突规则是否失效。
错误二:把工具描述当成安全声明
表现:
- 工具描述写着“只读”,实际调用触发缓存写入或审计写入;
- 外部 MCP Server 修改了工具注解后,执行器错误地放宽并发;
- 远程工具被替换后,原有调度策略仍继续使用。
诊断:
工具的并发属性应来自可信注册表或服务端策略,而不是仅来自模型可见的自然语言描述。MCP 规范也明确指出工具注解不应被无条件信任。(modelcontextprotocol.io)
错误三:按完成顺序合并结果
表现:
- 同样的请求每次模型回答略有不同;
- 模型错误地把一个工具的结果当成另一个工具的结果;
- 回放日志无法重现线上问题。
诊断:
检查结果合并是否保留:
call_id
plan_index
dependency_level
tool_name
source
合并应使用计划序号,而不是网络返回时间排序。
错误四:超时后直接重试写入
表现:
- 重复扣款;
- 重复创建订单;
- 重复发送消息;
- 本地显示失败,远端实际已经成功。
诊断:
检查超时信封中是否区分:
local timeout
remote commit unknown
没有幂等键、状态查询或补偿机制的写工具,不应自动重试。
错误五:只在执行器加锁,不在工具服务端校验
表现:
- 单实例执行器正确,多实例部署后出现重复写入;
- 两个 Agent 进程各自持有本地锁,仍然同时修改同一资源;
- 任务迁移到另一节点后破坏顺序。
原因:
本地锁只能约束本地进程。跨进程、跨机器、跨服务的正确性必须由以下至少一种机制保证:
- 数据库事务和行锁;
- 分布式锁;
- 条件更新和版本号;
- 幂等键;
- 服务端状态机;
- 单写者队列。
执行器的依赖图是调度层保证,不能替代外部系统的最终一致性边界。
十三、生产取舍:并发度、顺序和可恢复性
并发度越高,不一定越快。系统需要同时考虑:
当大量只读任务同时访问同一个下游服务时,瓶颈可能从 Agent 端转移到数据库连接池或 API 限流器。此时应在调度器中增加:
- 每会话并发上限;
- 每租户并发上限;
- 每工具并发上限;
- 每资源冲突域队列;
- 高风险写入的审批门;
- 截止时间传播;
- 失败后的取消和恢复策略。
读写策略可以按风险分级:
| 操作类型 | 默认策略 | 原因 |
|---|---|---|
| 独立纯读取 | 并发 | 无共享写副作用 |
| 同一快照的多个读取 | 共享事务或聚合查询 | 需要一致观察点 |
| 不同资源的写入 | 可并发 | 冲突域不相交 |
| 同一资源的非交换写入 | 串行 | 顺序影响结果 |
| 同一资源的可交换原子更新 | 可并发但由服务端保证 | 依赖原子语义 |
| 不可逆外部动作 | 串行、幂等、可查询 | 超时后状态可能未知 |
| 需要用户确认的高风险动作 | 先暂停 | 并发不能绕过授权 |
这里的“可交换”是指:
例如两个对计数器执行原子加法的操作,在满足溢出和业务限制不影响结果的前提下可能可交换;但“设置状态为已支付”和“设置状态为已取消”通常不可交换。
十四、最终设计原则
Agent 并行工具调用的核心不是把等待时间压缩到最短,而是让执行顺序与业务语义一致。
一个可靠的执行器应遵循以下逻辑:
- 把每个模型工具调用解析为带 ID 的节点;
- 从显式参数、规划关系和工具元数据中建立依赖图;
- 检测循环依赖;
- 区分纯读取、带隐式副作用的读取和写入;
- 为调用计算资源键和冲突域;
- 只并发执行依赖已满足且资源不冲突的节点;
- 对同一冲突域的非交换写入按确定顺序提交;
- 为外部写入保留幂等键和远端提交状态;
- 按调用 ID 和计划顺序稳定合并结果;
- 对部分成功、超时和状态未知进行显式建模;
- 将协议错误、参数错误、业务错误和远端未知状态分开;
- 把最终一致性责任落实到工具服务端,而不是只依赖 Agent 进程内的锁。
模型可以提出一组行动,但只有依赖图、冲突分析、授权检查、串行提交和结构化合并共同成立时,这组行动才真正构成一个可安全执行的 Agent 计划。
系列导航与关联阅读
- 系列入口:Agent 工程完整路线:从运行循环、记忆与协议到安全、评测和生产交付
- 上一篇:Agent 工具结果契约:结构、大小、引用、不可信内容和回写
- 下一篇:Agent 工具注册表:能力发现、租户过滤、版本和动态装配
- 延伸:Agent 工具执行器:严格解码、授权、超时、幂等和错误信封
- 延伸:Agent 并发与调度:会话隔离、任务池、优先级、公平和资源配额
官方资料
本文依据 Agent、模型、协议与框架官方资料重新梳理;正文、示例与生产清单由 WR BLOG 编写。

评论
0 条讨论