Agent 工程体系 · 第 12/98 篇。内容以 2026 年 9 月可验证的公开规范和稳定接口为基线;框架版本敏感能力会明确标注,不把实验行为写成通用保证。

Agent 并行工具调用:依赖图、只读并发、写入串行和合并

Agent 的一次行动,通常不是“模型调用一个工具,等待结果,再调用下一个工具”这么简单。模型可能在同一轮输出多个工具调用;执行器则必须判断这些调用之间是否存在数据依赖、资源冲突或业务顺序约束,然后决定哪些可以并行,哪些必须串行,哪些结果可以合并,哪些失败会使整个计划失效。

这一区分很重要:

  • 模型决定要调用什么工具
  • 执行器决定以什么顺序、以什么并发度执行
  • 工具服务负责真正的外部副作用和一致性约束
  • 合并器负责把多个结果重新组织成下一轮模型可理解的事实

不能因为模型在一次响应中返回了多个工具调用,就直接对它们执行 Promise.allasyncio.gather。并行调用只是一个候选执行集合,不等于这些操作在语义上彼此独立。

OpenAI 的 Function Calling 文档把工具调用描述为模型向应用发出的工具请求,并允许模型在一个回合中返回多个函数调用;在支持的模型上,可以使用 parallel_tool_calls: false 强制每轮最多执行一个工具调用。这个参数只能控制模型输出形式,不能替代执行器对依赖和冲突的判断。(developers.openai.com)

MCP 则定义了模型应用与外部工具、数据源之间的标准化连接方式。MCP 规定工具通过 tools/list 发现、通过 tools/call 调用,并用 JSON Schema 描述输入和可选的输出结构;它并不替应用决定多个工具调用之间的调度策略。(modelcontextprotocol.io)


一、先区分四个容易混淆的概念

1. 工具调用不是工具执行

一个工具调用可以抽象为:

ci=(idi,namei,argsi,ctxi)c_i = (id_i, name_i, args_i, ctx_i)

其中:

  • idiid_i:本次调用的唯一标识;
  • nameiname_i:工具名称;
  • argsiargs_i:模型生成的参数;
  • ctxictx_i:会话、用户、授权、租户、追踪信息等执行上下文。

工具执行则是:

ei=execute(ci)e_i = \operatorname{execute}(c_i)

执行结果还需要带上状态:

ri=(idi,statusi,outputi,errori,metadatai)r_i = (id_i, status_i, output_i, error_i, metadata_i)

因此,模型给出:

[
  {
    "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”

只读工具的关键不是名称,而是可观察副作用

一个工具只有在以下条件都成立时,才能被视为只读:

  1. 不修改数据库、缓存、文件、队列、外部 API 状态;
  2. 不触发发送消息、计费、下单、审批等隐式动作;
  3. 不依赖会被其他并发调用修改的共享状态;
  4. 失败、重试、超时不会改变外部系统状态;
  5. 读取结果可以接受其一致性级别,例如允许读取副本或稍旧快照。

例如:

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 行动中的工具调用建模为有向图:

G=(V,E)G=(V,E)

其中:

  • VV:工具调用节点集合;
  • EE:依赖边集合;
  • uvu \rightarrow v:表示调用 vv 必须等待调用 uu 的某个结果或提交状态。

例如,用户要求:

查询订单详情和物流状态,如果订单存在且未签收,再发送一条提醒。

可以拆成:

A: 查询订单详情
B: 查询物流状态
C: 判断订单是否满足提醒条件
D: 发送提醒

依赖关系是:

A ─┐
   ├──> C ───> D
B ─┘

对应边集合:

E={AC, BC, CD}E=\{A\rightarrow C,\ B\rightarrow C,\ C\rightarrow D\}

因为 A 和 B 都只是查询,并且互不依赖,所以它们可以并行。C 必须等待 A、B。D 必须等待 C 的判断结果。

2. 依赖边与资源边

依赖图中的边不只有一种来源。

数据依赖

调用 B 需要调用 A 的输出:

A: find_customer(email) -> customer_id
B: list_orders(customer_id)

边为:

ABA \rightarrow B

写后读依赖

调用 A 修改数据,调用 B 读取修改后的数据:

A: update_order_status(order_id, "shipped")
B: get_order(order_id)

如果 B 要观察 A 的新状态,就必须有:

ABA \rightarrow B

否则 B 可能读到旧值。

写写冲突

两个写调用修改同一资源:

A: set_order_status(order_id, "paid")
B: set_order_status(order_id, "cancelled")

它们可能没有数据依赖,但存在写写冲突。执行器必须:

  • 依据业务顺序建立边;
  • 或拒绝同时执行;
  • 或交给数据库事务、版本号、条件更新解决。

如果确定业务顺序是先支付、后取消,则加入:

ABA \rightarrow B

如果顺序不确定,就不能随意选择一个顺序并声称结果正确。此时应该返回冲突,让规划器或用户决定。

读写冲突

一个调用读取资源,另一个调用修改相同资源:

A: get_account_balance(account_id)
B: withdraw(account_id, 100)

若 A 的结果用于决定 B 是否可以执行,则是数据依赖:

ABA \rightarrow B

若 A 只是生成日志或展示信息,则可以在某些一致性要求下并行,但读取结果不应被解释为扣款前后的确定快照。

3. 为什么依赖图必须是 DAG

调度器通常需要一个有向无环图,即 DAG(Directed Acyclic Graph,有向无环图)。

如果出现:

A -> B -> C -> A

则不存在满足所有依赖的第一个节点,拓扑排序失败。循环依赖常见于:

  • 工具 A 等待工具 B 返回资源 ID;
  • 工具 B 又要求工具 A 先创建确认状态;
  • 模型把两个互相需要的动作同时放进计划;
  • 执行器错误地把“结果合并”反向当成了工具依赖。

调度器必须在执行前检测环,而不是让任务全部进入等待状态后才发现死锁。


三、从依赖图推导可执行批次

1. 入度决定任务是否就绪

对每个节点 vv,定义入度:

indegree(v)={uuv}indegree(v)=|\{u \mid u\rightarrow v\}|

入度为 0 的节点没有未完成前置依赖,可以进入 ready 状态。

执行过程如下:

  1. 计算所有节点的入度;
  2. 将入度为 0 的节点放入就绪队列;
  3. 取出满足资源和授权条件的节点执行;
  4. 节点成功完成后,删除其出边;
  5. 后继节点入度降为 0 时,加入就绪队列;
  6. 重复直到全部完成或无法继续。

伪代码:

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. 调度条件的形式化

一个节点 vv 可以在时刻 tt 启动,当且仅当:

Ready(v,t)=DepsDone(v,t)Authz(v,t)QuotaAvailable(v,t)ConflictFree(v,t)DeadlineValid(v,t)Ready(v,t) = DepsDone(v,t) \land Authz(v,t) \land QuotaAvailable(v,t) \land ConflictFree(v,t) \land DeadlineValid(v,t)

各项含义是:

  • DepsDoneDepsDone:所有前置节点已经完成并满足该节点要求;
  • AuthzAuthz:当前用户、Agent、租户和工具权限允许执行;
  • QuotaAvailableQuotaAvailable:会话、租户、工具和全局资源配额允许执行;
  • ConflictFreeConflictFree:没有与运行中任务发生不允许的资源冲突;
  • DeadlineValidDeadlineValid:启动后仍有足够时间完成或达到可接受的截止时间。

只要其中一项不满足,节点就应保持 blockedqueued,而不是强行执行。


四、只读并发:安全条件与边界

1. 只读调用可并行的充分条件

设两个调用 r1,r2r_1,r_2 都是只读。若满足:

Writes(r1)=Writes(r2)=Writes(r_1)=Writes(r_2)=\varnothing

并且它们不依赖彼此的结果:

r1↛r2,r2↛r1r_1 \not\rightarrow r_2,\quad r_2 \not\rightarrow r_1

同时读取的一致性要求允许它们观察不同时间点的数据,则可以并行:

Parallelizable(r1,r2)=trueParallelizable(r_1,r_2)=true

这个条件是工程上的安全近似,而不是数据库理论中的唯一条件。

例如:

get_user_profile(user_id=123)
list_recent_orders(user_id=123)

二者都读取数据,且一个不需要另一个的输出,可以并发。

但如果业务要求“用户资料和订单必须来自同一个数据库快照”,那么仅仅都是只读还不够。它们需要:

  • 共享同一个事务快照;
  • 或使用同一个版本号;
  • 或由后端提供聚合查询;
  • 或接受一致性声明中的“可能跨时间点”。

2. 只读并发的时间收益

设三个独立只读工具耗时分别为:

T1=300ms,T2=500ms,T3=800msT_1=300ms,\quad T_2=500ms,\quad T_3=800ms

串行执行的理想耗时是:

Tserial=T1+T2+T3=1600msT_{serial}=T_1+T_2+T_3=1600ms

并行执行的理想耗时接近:

Tparallel=max(T1,T2,T3)=800msT_{parallel}=\max(T_1,T_2,T_3)=800ms

实际耗时还要加上调度、连接、限流、序列化和结果合并开销:

Tactual=max(Ti)+Tschedule+Tmerge+TqueueT_{actual} = \max(T_i)+T_{schedule}+T_{merge}+T_{queue}

因此,只有在工具耗时足够大、并发资源可用、服务端不会因突发并发而显著排队时,并行才有实际收益。

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. 写入为什么需要顺序

设两个操作都读取旧值 x0x_0,再写回计算结果:

A: x = x + 10
B: x = x + 20

初始:

x0=100x_0=100

若 A、B 并发且都执行“读—计算—写”:

A 读到 100
B 读到 100
A 写入 110
B 写入 120

最终结果是 120,而正确的累加结果应为:

100+10+20=130100+10+20=130

这就是丢失更新。

如果串行执行:

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")},
)

两个调用是否冲突,可以定义为:

Conflict(a,b)=Writes(a)(Reads(b)Writes(b))Conflict(a,b) = Writes(a)\cap (Reads(b)\cup Writes(b))\neq\varnothing

或:

Writes(b)(Reads(a)Writes(a))Writes(b)\cap (Reads(a)\cup Writes(a))\neq\varnothing

对于两个写调用,还可以采用更严格的判断:

Conflictww(a,b)=Writes(a)Writes(b)Conflict_{ww}(a,b)=Writes(a)\cap Writes(b)\neq\varnothing

工程上通常不把完整的读写集合交给模型自行填写,而是由工具注册表提供。例如:

{
  "name": "reserve_inventory",
  "execution": {
    "mode": "write",
    "resource_keys": ["inventory:{sku}"],
    "commutativity": "non_commutative",
    "idempotency": "required"
  }
}

这里的 resource_keys 是执行器元数据,不应仅依赖工具描述文本。MCP 的工具描述包括名称、描述和输入 Schema,也允许工具提供行为注解;但规范明确要求客户端把工具注解视为不可信,除非工具来自可信服务器。(modelcontextprotocol.io)

3. 写入顺序来自哪里

写入顺序可以有四种来源:

规划顺序

模型或上层规划器明确生成:

先创建订单,再扣库存,再发送确认邮件

执行器将其转成:

CreateOrderReserveInventorySendEmailCreateOrder \rightarrow ReserveInventory \rightarrow SendEmail

资源顺序

多个任务操作同一个资源时,执行器依据提交序列号排序:

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. 合并的输入和目标

合并器接收多个工具结果:

R={r1,r2,,rn}R=\{r_1,r_2,\dots,r_n\}

输出一个供 Agent 状态机消费的结构:

M(R)=(facts, errors, effects, provenance, next_actions)M(R)= ( facts,\ errors,\ effects,\ provenance,\ next\_actions )

其中:

  • 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]

应按照以下优先级之一排序:

  1. 计划中的节点序号;
  2. 拓扑序;
  3. 工具调用在模型响应中的索引;
  4. 稳定的 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

校验:

owner=u1owner=u1

available=5quantity=2available=5 \ge quantity=2

notification_allowed=truenotification\_allowed=true

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 确认提交后执行。因为“库存预留成功”是发送确认消息的前置条件:

DED \rightarrow E

如果 D 失败,E 不应发送“预留成功”的确认消息。

如果 E 超时:

D 已提交
E 状态未知

此时不能回滚库存并简单重试,除非系统确实支持可靠补偿。更安全的处理是:

  1. 使用相同幂等键重试发送;
  2. 查询消息服务的发送状态;
  3. 如果无法确认,写入待处理任务;
  4. 向模型报告“库存已预留,通知状态未知”。

最终合并结果:

{
  "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_idtransaction_idreservation_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

如果工具不支持幂等,也无法查询状态,那么执行器不应自动重试不可逆写入,只能将状态交给人工或专门的恢复流程。

失败传播规则

对依赖图中的节点 vv,可以定义:

Runnable(v)=uPred(v)status(u)=SucceededRunnable(v)= \bigwedge_{u\in Pred(v)} status(u)=Succeeded

如果某个前置节点失败,则后继节点默认进入 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 进程各自持有本地锁,仍然同时修改同一资源;
  • 任务迁移到另一节点后破坏顺序。

原因:

本地锁只能约束本地进程。跨进程、跨机器、跨服务的正确性必须由以下至少一种机制保证:

  • 数据库事务和行锁;
  • 分布式锁;
  • 条件更新和版本号;
  • 幂等键;
  • 服务端状态机;
  • 单写者队列。

执行器的依赖图是调度层保证,不能替代外部系统的最终一致性边界。


十三、生产取舍:并发度、顺序和可恢复性

并发度越高,不一定越快。系统需要同时考虑:

Throughputmin(agent_quota,tool_quota,connection_pool,database_capacity,conflict_free_capacity)Throughput \approx \min( agent\_quota, tool\_quota, connection\_pool, database\_capacity, conflict\_free\_capacity )

当大量只读任务同时访问同一个下游服务时,瓶颈可能从 Agent 端转移到数据库连接池或 API 限流器。此时应在调度器中增加:

  • 每会话并发上限;
  • 每租户并发上限;
  • 每工具并发上限;
  • 每资源冲突域队列;
  • 高风险写入的审批门;
  • 截止时间传播;
  • 失败后的取消和恢复策略。

读写策略可以按风险分级:

操作类型 默认策略 原因
独立纯读取 并发 无共享写副作用
同一快照的多个读取 共享事务或聚合查询 需要一致观察点
不同资源的写入 可并发 冲突域不相交
同一资源的非交换写入 串行 顺序影响结果
同一资源的可交换原子更新 可并发但由服务端保证 依赖原子语义
不可逆外部动作 串行、幂等、可查询 超时后状态可能未知
需要用户确认的高风险动作 先暂停 并发不能绕过授权

这里的“可交换”是指:

Apply(a,Apply(b,S))=Apply(b,Apply(a,S))Apply(a, Apply(b,S)) = Apply(b, Apply(a,S))

例如两个对计数器执行原子加法的操作,在满足溢出和业务限制不影响结果的前提下可能可交换;但“设置状态为已支付”和“设置状态为已取消”通常不可交换。


十四、最终设计原则

Agent 并行工具调用的核心不是把等待时间压缩到最短,而是让执行顺序与业务语义一致。

一个可靠的执行器应遵循以下逻辑:

  1. 把每个模型工具调用解析为带 ID 的节点;
  2. 从显式参数、规划关系和工具元数据中建立依赖图;
  3. 检测循环依赖;
  4. 区分纯读取、带隐式副作用的读取和写入;
  5. 为调用计算资源键和冲突域;
  6. 只并发执行依赖已满足且资源不冲突的节点;
  7. 对同一冲突域的非交换写入按确定顺序提交;
  8. 为外部写入保留幂等键和远端提交状态;
  9. 按调用 ID 和计划顺序稳定合并结果;
  10. 对部分成功、超时和状态未知进行显式建模;
  11. 将协议错误、参数错误、业务错误和远端未知状态分开;
  12. 把最终一致性责任落实到工具服务端,而不是只依赖 Agent 进程内的锁。

模型可以提出一组行动,但只有依赖图、冲突分析、授权检查、串行提交和结构化合并共同成立时,这组行动才真正构成一个可安全执行的 Agent 计划。


系列导航与关联阅读

官方资料

本文依据 Agent、模型、协议与框架官方资料重新梳理;正文、示例与生产清单由 WR BLOG 编写。