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

Agent 并发与调度:会话隔离、任务池、优先级、公平和资源配额

Agent 并发不是简单地把多个 async 函数同时执行。一个生产级 Agent 系统通常同时面对四种不同的并发:

  1. 会话并发:多个用户、多个会话同时运行。
  2. 任务并发:一个会话内部同时推进多个子任务。
  3. 工具并发:一次模型决策产生多个工具调用。
  4. 输出并发:生成、事件推送、消息发送和客户端连接状态彼此独立。

如果只在最外层增加线程数或协程数,系统可能获得更高吞吐,却同时失去会话隔离、写入一致性、优先级控制和成本边界。Agent 运行时需要解决的核心问题是:

哪些工作可以同时执行,哪些工作必须有序执行;谁可以获得资源,获得多少资源;当系统过载、任务过期或执行失败时,哪些结果仍然有效。

OpenAI Agents SDK 将 Agent 描述为能够规划、调用工具、协作并维护足够状态以完成多步工作的应用;其 SDK 运行器负责 Agent 循环、工具调用和 Agent 切换,而 Responses API 则允许应用自行控制循环和分支。这个区别很重要:调度策略属于应用运行时,不应被误认为模型自动提供的并发保证。(developers.openai.com)

Anthropic 将固定代码路径编排称为 workflow,将由模型动态决定下一步和工具使用的系统称为 agent,并把并行化、路由、orchestrator-workers 等视为可组合的系统模式。其建议是从简单方案开始,只有在复杂度确实改善结果时才引入更多 Agent 和调度层。(anthropic.com)


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

1. 并发、并行、异步不是同一个概念

并发表示多个任务在时间上交错推进;并行表示多个任务在同一时刻使用不同计算资源执行;异步是程序组织等待和回调的方式。

例如,一个进程只有一个事件循环,也可以并发执行多个网络请求:

请求 A:发送 ───── 等待 ───── 收到结果
请求 B:       发送 ───── 等待 ───── 收到结果

这里有并发,但不一定有 CPU 并行。

对 Agent 而言,网络等待通常不是主要难点。真正需要调度的是:

  • 模型调用是否可以同时发起;
  • 工具是否访问同一资源;
  • 工具调用是否具有副作用;
  • 子任务是否属于同一个会话;
  • 输出是否必须按版本顺序发送;
  • 任务是否已经失去业务意义。

因此,asyncio.gather() 只是一个执行原语,不是完整的并发策略。

2. 会话、运行、任务和工具调用的层级

可以把一个 Agent 系统抽象成四层:

用户会话 Session
└── 一次或多次 Agent Run
    └── 一个或多个 Task
        └── 一个或多个 Tool Call

含义分别是:

  • Session:用户或业务实体的长期交互边界,包含会话身份、权限、历史、取消状态和资源预算。
  • Run:一次从输入到结果或暂停状态的 Agent 执行。
  • Task:可调度的工作单元,例如“查询订单”“分析日志”“生成回复”。
  • Tool Call:模型或程序请求调用某个工具的具体操作。

一个会话可以有多个 Run,但通常不应让同一个会话的两个 Run 无条件修改同一份会话状态。否则,后完成的旧请求可能覆盖先完成的新请求。


二、会话隔离:并发安全首先是状态安全

1. 什么是会话隔离

会话隔离要求一个会话的状态、权限、取消信号、任务预算和输出通道不能被另一个会话错误使用。

至少应隔离以下对象:

对象 典型内容 不能共享的原因
身份上下文 用户 ID、租户 ID、权限 防止越权访问
对话状态 消息历史、工具结果、摘要 防止上下文串线
运行状态 run_id、取消标记、暂停点 防止错误恢复
资源预算 token、工具次数、执行时长 防止单会话耗尽全局资源
输出序列 事件序号、消息版本 防止旧结果覆盖新结果
临时文件和沙箱 工作目录、凭证、缓存 防止数据泄露和文件冲突

隔离并不意味着每个会话都必须拥有独立进程。实际系统通常使用共享 Worker,但要求每个任务携带不可变的 SessionContext,并在所有外部访问处验证上下文。

2. 不要把“会话 ID”当成全部隔离

仅在日志中增加 session_id 不构成隔离。下面的实现仍然是不安全的:

history = {}  # session_id -> messages

async def run_agent(session_id, user_text):
    messages = history.setdefault(session_id, [])
    messages.append({"role": "user", "content": user_text})

    result = await call_model(messages)

    messages.append({"role": "assistant", "content": result})
    return result

问题在于同一会话的两个请求可能交错:

请求 A:读取 [旧历史]
请求 B:读取 [旧历史]
请求 A:追加问题 A
请求 B:追加问题 B
请求 B:写入答案 B
请求 A:写入答案 A

最终历史可能变成:

问题 A
问题 B
答案 B
答案 A

如果答案 A 依赖的是问题 A 之后的上下文,顺序就已经失真。

更严重的是,messages 是可变对象,模型调用期间其他协程可能修改它;即使底层字典是线程安全的,也不能保证“读取历史—生成结果—提交结果”这一整个事务的原子性。

3. 两种会话并发策略

策略 A:同一会话串行化

对同一个 session_id 使用互斥锁或单会话队列:

Session A: Run 1 ───── Run 2 ───── Run 3
Session B: Run 1 ── Run 2

不同会话之间仍然可以并发。

优点:

  • 语义简单;
  • 对话顺序清晰;
  • 不需要合并冲突。

缺点:

  • 一个长任务会阻塞同一会话的后续请求;
  • 用户发送“停止”“修改问题”等控制消息时响应可能变慢。

策略 B:会话内允许并发,但使用版本提交

给会话状态设置单调递增的版本号:

读取版本 v=10
生成答案
提交时要求 current_version == 10

若提交时版本已经变成 11,则当前结果不能直接覆盖状态,必须:

  1. 丢弃;
  2. 重新基于最新状态生成;
  3. 作为独立消息发送,不写入主历史;
  4. 或进入人工/程序合并流程。

形式化地,状态更新可以写成:

commit(s,v,Δ)={sapply(s,Δ), versionv+1current_version=vconflictcurrent_versionvcommit(s, v, \Delta) = \begin{cases} s \leftarrow apply(s, \Delta),\ version \leftarrow v+1 & current\_version=v \\ conflict & current\_version\neq v \end{cases}

其中:

  • ss 是会话状态;
  • vv 是任务读取状态时的版本;
  • Δ\Delta 是任务产生的状态变更;
  • conflict 表示任务基于旧状态,不能直接提交。

版本检查只能防止静默覆盖,不能自动解决语义冲突。 两个任务都给同一账户退款时,即使版本检查正确,也仍需要业务幂等键和事务约束。


三、工具调用并发:依赖图决定执行顺序

1. 从工具列表建立依赖图

模型返回多个工具调用时,不应默认全部并发,也不应默认全部串行。应先构造有向图:

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

其中:

  • VV 是工具调用;
  • EE 是依赖边;
  • aba\rightarrow b 表示工具调用 bb 必须等待 aa 完成。

例如:

fetch_user ───────┐
                  ├── build_report ── send_report
fetch_orders ─────┘

可以同时执行 fetch_userfetch_orders,但 build_report 必须等待二者完成。

若图无环,则可按拓扑层级执行:

第 0 层:fetch_user, fetch_orders
第 1 层:build_report
第 2 层:send_report

若图中出现环:

A -> B -> C -> A

则说明规划结果无法直接执行,应拒绝、改写或中止,而不是让 Worker 永久等待。

2. 只读并发

只读工具只观察外部状态,不产生可见副作用。例如:

  • 查询用户资料;
  • 查询订单;
  • 读取知识库;
  • 获取监控指标;
  • 读取文件内容。

如果两个工具满足以下条件,就可以并发:

parallel(a,b)=readOnly(a)readOnly(b)noDataDependency(a,b)parallel(a,b) = readOnly(a)\land readOnly(b)\land noDataDependency(a,b)

还应补充一个经常被忽略的条件:二者不能依赖同一个会改变结果的外部快照。如果一个查询要求“在同一数据库快照上读取”,就需要事务快照或显式版本号,而不是普通并发请求。

3. 写入串行化

写入工具会改变外部状态,例如:

  • 创建订单;
  • 修改工单;
  • 扣款;
  • 发送邮件;
  • 删除文件;
  • 发布消息。

两个写入操作之间只要存在以下任一关系,就不能无条件并发:

  • 写同一资源;
  • 一个操作依赖另一个的结果;
  • 两者的业务顺序有意义;
  • 外部 API 不支持安全重试;
  • 操作不可逆。

写入串行不一定意味着整个系统只有一个 Worker,而是要求同一资源的写入经过同一个顺序化边界

全局 Worker 池
├── resource:user:42  -> 串行队列
├── resource:order:9  -> 串行队列
└── resource:file:x   -> 串行队列

不同资源仍可并发。

4. 写入合并

某些写操作可以先在内存或临时状态中合并,再一次提交。例如:

set_label("urgent")
set_label("customer-visible")
set_label("needs-review")

可以转换为:

replace_labels({"urgent", "customer-visible", "needs-review"})

但只有满足以下条件时才可合并:

  1. 操作满足幂等或可交换;
  2. 中间状态不会触发外部事件;
  3. 最终状态足以表达全部意图;
  4. 合并过程不会丢失条件检查。

反例:

扣款 100 元
扣款 50 元

不能简单合并成“扣款 150 元”,因为两次操作可能对应不同订单、不同幂等键或不同授权检查。

5. 完整执行示例

假设模型规划出以下调用:

T1 = get_customer(customer_id=42)          # 只读
T2 = list_open_tickets(customer_id=42)     # 只读
T3 = update_ticket(ticket_id=8, priority=1) # 写入
T4 = send_email(to="user@example.com")     # 写入

依赖关系:

T1 ───────┐
          ├── T3 ─── T4
T2 ───────┘

执行过程:

时刻 状态 可执行任务
t0t_0 T1、T2 就绪 T1、T2 并发
t1t_1 T1 完成,T2 未完成 等待 T2
t2t_2 T1、T2 都完成 执行 T3
t3t_3 T3 成功 执行 T4
t4t_4 T4 成功 Run 完成

如果 T3 失败,T4 不能继续,因为发送邮件会声明一个并未成功完成的更新。若 T4 失败,则工单更新可能已经成功,系统应记录为“业务成功、通知失败”,而不是把整个 Run 简化成一个失败状态。


四、任务池:把 Agent 推理和实际执行解耦

1. 任务池的职责

任务池是保存待执行任务的调度结构。它至少需要处理:

  • 入队;
  • 出队;
  • 优先级排序;
  • 并发限制;
  • 取消;
  • 超时;
  • 重试;
  • 过期丢弃;
  • Worker 崩溃后的恢复;
  • 结果回传;
  • 队列长度和等待时间观测。

任务池中的任务不应只是一个函数指针。生产任务通常需要携带完整元数据:

from dataclasses import dataclass, field
from typing import Any

@dataclass
class Task:
    task_id: str
    session_id: str
    run_id: str
    kind: str
    payload: dict[str, Any]
    priority: int
    created_at: float
    deadline: float | None = None
    cost_tokens: int = 0
    resource_keys: tuple[str, ...] = ()
    idempotency_key: str | None = None
    cancelled: bool = False
    attempt: int = 0

这些字段分别解决不同问题:

  • session_id:会话隔离;
  • run_id:同一会话内区分一次执行;
  • priority:调度顺序;
  • deadline:任务是否仍有意义;
  • cost_tokens:资源预算;
  • resource_keys:写入冲突控制;
  • idempotency_key:重试不重复产生副作用;
  • attempt:重试次数和退避策略。

2. 生成 Worker、执行 Worker 和发送器

Agent 系统常见的组件关系如下:

flowchart LR
    U[用户请求] --> S[会话入口]
    S --> Q[任务队列]
    Q --> G[生成 Worker]
    G --> P[规划结果 / 子任务]
    P --> D[依赖图调度器]
    D --> E[执行 Worker 池]
    E --> R[结果事件流]
    R --> A[会话状态提交器]
    R --> O[发送器]
    O --> C[客户端]
    Q --> X[过期与取消清理]
    A --> DB[(状态存储)]

这里的“生成 Worker”负责:

  1. 读取会话上下文;
  2. 调用模型;
  3. 解析工具调用或子任务;
  4. 生成依赖关系;
  5. 把可执行节点放入执行队列。

“执行 Worker”负责:

  1. 获取一个已获准执行的任务;
  2. 检查取消和过期状态;
  3. 获取资源锁或配额;
  4. 调用工具;
  5. 记录结果;
  6. 释放资源并通知依赖节点。

“发送器”负责:

  1. 从结果事件流读取事件;
  2. 判断事件是否属于当前会话版本;
  3. 检查客户端是否仍连接;
  4. 按序发送;
  5. 在发送失败时重连、重放或丢弃。

将发送器独立出来很重要。模型生成完成不代表客户端一定还能接收;如果执行 Worker 直接写 WebSocket,网络阻塞会反向占用执行资源。


五、背压和过期丢弃:队列不能无限增长

1. 背压是什么

当任务进入速度大于系统处理速度时:

λin>λout\lambda_{in} > \lambda_{out}

队列长度会不断增长。若没有背压,系统会以以下方式失败:

  • 内存被队列占满;
  • 所有请求延迟变大;
  • 旧任务占用 Worker;
  • 新任务即使更重要也无法及时执行;
  • 重试任务进一步放大负载。

背压是让上游感知系统容量不足的机制。常见做法有:

  • 限制队列长度;
  • 入队超时;
  • 按会话限制待处理任务数;
  • 按租户限制并发;
  • 拒绝低优先级任务;
  • 返回“排队中”状态;
  • 把任务转为异步后台作业;
  • 对重复任务做去重或合并。

2. 有界队列的基本实现

import asyncio

class BoundedTaskPool:
    def __init__(self, maxsize: int, workers: int):
        self.queue = asyncio.Queue(maxsize=maxsize)
        self.workers = [
            asyncio.create_task(self._worker(i))
            for i in range(workers)
        ]

    async def submit(self, task, timeout: float = 0.2) -> bool:
        try:
            await asyncio.wait_for(self.queue.put(task), timeout)
            return True
        except asyncio.TimeoutError:
            return False

    async def _worker(self, worker_id: int):
        while True:
            task = await self.queue.get()
            try:
                if task.is_expired():
                    continue
                await execute(task)
            finally:
                self.queue.task_done()

输入是一个带容量上限的任务队列;预期行为是:队列未满时入队,队列持续满载超过 timeout 时拒绝提交。

这段代码仍不是完整生产实现,因为它没有处理:

  • 优先级;
  • 取消;
  • Worker 崩溃;
  • 重试;
  • 资源锁;
  • 跨进程持久化;
  • 任务结果确认。

它适合说明一个关键因果关系:有界队列把无限等待转换成明确的拒绝或降级行为

3. 为什么必须丢弃过期任务

不是所有排队任务都值得执行。例如:

  • 用户已经发送了新的问题;
  • 前端已关闭页面;
  • 语音转写结果已经被更新;
  • 搜索建议已经超过有效时间;
  • 旧的流式生成事件已经被新版本替代。

可以给任务定义截止时间:

expired(t,now)    deadline(t)nowdeadline(t)expired(t, now) \iff deadline(t)\neq\varnothing \land now \ge deadline(t)

但“过期检查”必须发生在多个阶段:

  1. 入队前;
  2. 从队列取出后;
  3. 获取资源配额前;
  4. 调用外部工具前;
  5. 结果提交前;
  6. 发送前。

只在入队时检查是不够的,因为任务可能在队列中等待很久。

一个安全的提交条件可以写成:

deliver(e)=¬expired(e)session_active(e.session)version_valid(e)not_cancelled(e.run)deliver(e) = \neg expired(e) \land session\_active(e.session) \land version\_valid(e) \land not\_cancelled(e.run)

即使模型调用已经完成,如果会话已取消或事件版本过旧,也不应发送。


六、优先级:决定“先做谁”,不等于“永远插队”

1. 优先级队列的基本定义

优先级队列按照任务优先级选择下一个任务。若任务 aa 的优先级高于任务 bb,则通常希望:

priority(a)>priority(b)abpriority(a) > priority(b) \Rightarrow a \prec b

其中 aba \prec b 表示在资源可用时,优先调度 aa

实际排序通常还需要加入创建时间和序列号:

排序键 = (-priority, created_at, sequence)

负号表示优先级越大越靠前;sequence 用于保证相同优先级下的稳定顺序。

2. 纯优先级的反例:饥饿

假设有两个队列:

低优先级:L1, L2, L3
高优先级:H1, H2, H3, ...

只要高优先级任务持续到达,调度器就永远执行 H:

H1 -> H2 -> H3 -> H4 -> ...
L1 从未执行

这称为饥饿。因此,“高优先级优先”不应等于“低优先级没有服务保证”。

3. 老化机制

一种简单方法是让等待时间提升有效优先级:

effective_priority(t)=base_priority(t)+wait(t)Aeffective\_priority(t) = base\_priority(t) + \left\lfloor \frac{wait(t)}{A} \right\rfloor

其中:

  • base_priority 是任务初始优先级;
  • wait(t) 是等待秒数;
  • AA 是老化周期;
  • 每等待 AA 秒,有效优先级增加 1。

例如:

任务 初始优先级 等待时间 A=10A=10 有效优先级
H1 10 0 秒 0 10
M1 5 20 秒 2 7
L1 1 100 秒 10 11

此时 L1 可能超过 H1,但这也可能不符合安全要求。因此老化通常设置上限:

effective_priority=min(prioritymax,base+wait/A)effective\_priority = \min(priority_{max}, base + \lfloor wait/A \rfloor)

对于支付、删除、权限变更等敏感操作,优先级只能影响排队顺序,不能绕过授权和人工审批。

4. 优先级反映业务价值,而不是模型自评

不应直接信任模型输出的 "priority": 100。优先级应由服务端根据可验证属性计算,例如:

最终优先级 =
  用户套餐权重
+ 业务类型权重
+ 是否在线请求
+ 是否接近截止时间
- 预估成本惩罚

模型可以提出“这很紧急”,但服务端必须重新分类。否则提示词注入可能把普通任务伪装成紧急任务,挤占其他用户资源。


七、公平:全局公平、租户公平和会话公平

1. 为什么有优先级还需要公平

优先级回答“哪个任务更重要”,公平回答“某个用户或会话是否长期得不到服务”。

如果一个租户可以同时提交数千个低优先级任务,它仍可能占满 Worker、连接池或模型速率限制。公平调度需要在任务优先级之外增加服务主体维度

常见主体包括:

  • 用户;
  • 会话;
  • 租户;
  • 产品功能;
  • 模型供应商;
  • 工具资源。

2. 加权公平轮转

将任务按会话分组:

Session A: A1, A2, A3
Session B: B1, B2
Session C: C1

轮转时:

A1 -> B1 -> C1 -> A2 -> B2 -> A3

如果不同主体有不同权重,可以使用加权轮转:

租户 Pro:Pro, Pro, Pro, ...
租户 Free:Free, ...

但加权轮转只能提供近似公平。若任务成本差异很大:

A1 成本 1 秒
B1 成本 60 秒

按任务数量公平不代表按资源消耗公平。更合理的做法是按消耗量或虚拟完成时间调度。

3. 资源公平的形式化

令任务 ii 的估计资源成本为 cic_i,主体 kk 的累计服务量为:

servicek=ikciservice_k = \sum_{i \in k} c_i

若主体权重为 wkw_k,调度器可以优先选择:

scorek=servicekwkscore_k = \frac{service_k}{w_k}

选择 score_k 最低的主体,再从其队列中取任务。

直觉是:权重越大,允许获得的服务量越多;但不能只统计任务个数,而要统计模型 token、工具执行时间、GPU 时间或其他受限资源。

4. 公平不等于所有任务相同待遇

在线交互请求和离线批处理不应共享完全相同的队列。常见划分是:

交互队列:短任务、低延迟、有限重试
后台队列:长任务、可暂停、低优先级
控制队列:取消、暂停、人工审批

控制任务通常需要“优先于普通任务”,否则用户点击取消后,取消消息也要排队等待长任务结束,系统就失去了可控性。


八、资源配额:限制并发数量还不够

1. 并发上限和资源配额的区别

max_workers = 20 只限制同时运行的任务数,不限制:

  • 每个任务消耗多少 token;
  • 一次运行最多调用多少工具;
  • 一个会话可以生成多少子任务;
  • 一个租户可以使用多少模型请求;
  • 某个工具的并发连接数;
  • 单次任务的最大执行时长;
  • 失败重试产生的额外成本。

因此需要至少三类配额:

  1. 数量配额:同时运行多少任务;
  2. 速率配额:单位时间允许多少请求或 token;
  3. 总量配额:一次 Run 或一个 Session 最多使用多少资源。

2. 分层配额

可以建立以下层级:

全局
├── 租户
│   ├── 用户
│   │   └── 会话
│   └── 工具类别
└── 模型供应商

任务只有在所有相关配额都允许时才能启动:

admit(t)=global_ok(t)tenant_ok(t)session_ok(t)tool_ok(t)deadline_ok(t)admit(t) = global\_ok(t) \land tenant\_ok(t) \land session\_ok(t) \land tool\_ok(t) \land deadline\_ok(t)

如果任务同时调用模型和数据库工具,它可能需要同时获得:

model_tokens: 2,000
db_connection: 1
external_api_request: 1

这属于多资源准入,而不是单一信号量。

3. 令牌桶用于速率,信号量用于并发

信号量适合限制同时运行数量:

model_slots = asyncio.Semaphore(8)

async def call_model_limited(request):
    async with model_slots:
        return await call_model(request)

令牌桶适合限制时间窗口内的速率。设:

  • 桶容量为 BB
  • 令牌产生速率为 rr
  • 每次请求消耗 qq 个令牌。

当当前令牌数小于 qq 时,请求必须等待或被拒绝。令牌桶允许短时突发,但长期平均速率不会超过 rr

只使用信号量的反例:

每个请求都很慢,但同时只有 8 个请求

系统可能不超出并发数,却持续超出供应商的每分钟 token 限制。只使用令牌桶的反例:

请求速率合法,但每个请求都持有一个数据库连接

最终仍可能耗尽连接池。因此实际系统通常组合多种限制。

4. 配额扣减的时机

配额至少有三种状态:

reserved -> consumed
        \-> released
  • reserved:任务获准启动,资源已预留;
  • consumed:实际使用量已经确认;
  • released:任务取消、过期或失败,未使用部分归还。

若等任务完成后才扣减配额,多个任务可能同时看到“余额充足”,造成超卖。若一开始按最大值扣减,又可能因为任务提前结束而浪费容量。

常见折中是:

  1. 启动前预留估计值;
  2. 执行中按实际使用追加;
  3. 完成后释放未使用部分;
  4. 对无法准确估计的资源使用保守上限。

九、一个可运行的调度器骨架

下面的示例演示四个机制:

  • 同一会话串行;
  • 不同会话并发;
  • 优先级排序;
  • 任务过期丢弃。

它不是模型 SDK 的代码,而是可以包在任意 Agent 运行器外层的运行时骨架。

import asyncio
import heapq
import itertools
import time
from dataclasses import dataclass, field


@dataclass(order=True)
class QueueItem:
    # heapq 默认最小值优先,所以 priority 使用负数
    sort_key: tuple = field(init=False, repr=False)
    priority: int
    sequence: int
    task: "AgentTask" = field(compare=False)

    def __post_init__(self):
        self.sort_key = (-self.priority, self.sequence)


@dataclass
class AgentTask:
    task_id: str
    session_id: str
    priority: int
    deadline: float | None
    duration: float
    text: str

    def expired(self) -> bool:
        return self.deadline is not None and time.monotonic() >= self.deadline


class Scheduler:
    def __init__(self, worker_count: int = 2):
        self._queue: list[QueueItem] = []
        self._condition = asyncio.Condition()
        self._sequence = itertools.count()
        self._session_locks: dict[str, asyncio.Lock] = {}
        self._workers = [
            asyncio.create_task(self._worker(i))
            for i in range(worker_count)
        ]

    def _session_lock(self, session_id: str) -> asyncio.Lock:
        return self._session_locks.setdefault(session_id, asyncio.Lock())

    async def submit(self, task: AgentTask):
        if task.expired():
            print(f"drop {task.task_id}: expired before enqueue")
            return

        async with self._condition:
            item = QueueItem(
                priority=task.priority,
                sequence=next(self._sequence),
                task=task,
            )
            heapq.heappush(self._queue, item)
            self._condition.notify()

    async def _next_task(self) -> AgentTask:
        async with self._condition:
            while not self._queue:
                await self._condition.wait()
            return heapq.heappop(self._queue).task

    async def _worker(self, worker_id: int):
        while True:
            task = await self._next_task()

            # 取出后再次检查,因为排队等待期间可能已经过期
            if task.expired():
                print(f"drop {task.task_id}: expired in queue")
                continue

            # 同一个会话串行,不同会话仍可由其他 worker 并发处理
            async with self._session_lock(task.session_id):
                if task.expired():
                    print(f"drop {task.task_id}: expired before run")
                    continue

                print(
                    f"worker={worker_id} start "
                    f"task={task.task_id} session={task.session_id}"
                )
                await asyncio.sleep(task.duration)
                print(f"worker={worker_id} done task={task.task_id}")


async def main():
    scheduler = Scheduler(worker_count=2)
    now = time.monotonic()

    await scheduler.submit(AgentTask(
        task_id="A1", session_id="S1",
        priority=1, deadline=now + 10,
        duration=0.4, text="first request",
    ))
    await scheduler.submit(AgentTask(
        task_id="A2", session_id="S1",
        priority=10, deadline=now + 10,
        duration=0.1, text="second request",
    ))
    await scheduler.submit(AgentTask(
        task_id="B1", session_id="S2",
        priority=5, deadline=now + 10,
        duration=0.2, text="other session",
    ))
    await scheduler.submit(AgentTask(
        task_id="OLD", session_id="S3",
        priority=100, deadline=now - 1,
        duration=0.1, text="already stale",
    ))

    await asyncio.sleep(1)


asyncio.run(main())

可能输出:

drop OLD: expired before enqueue
worker=0 start task=B1 session=S2
worker=1 start task=A2 session=S1
worker=1 done task=A2
worker=1 start task=A1 session=S1
worker=0 done task=B1
worker=1 done task=A1

这里有一个值得注意的边界:代码按全局优先级取出任务,因此 A2 可能先于 A1 出队;但 S1 的会话锁保证它们不会同时执行。若业务要求“同一会话严格按提交顺序”,就不能只加锁,还必须在会话队列中保持 FIFO。

换句话说:

  • 解决“不能同时执行”;
  • 队列顺序解决“谁先执行”;
  • 版本检查解决“旧结果能否提交”;
  • 幂等键解决“重试是否重复产生副作用”。

十、生成、发送和过期事件的生命周期

一个交互式 Agent 通常不是“收到请求后阻塞到返回字符串”,而是事件流:

accepted
  -> queued
  -> planning
  -> tool_started
  -> tool_completed
  -> generating
  -> output_delta
  -> completed

事件必须带有足够的身份和顺序信息:

{
  "session_id": "s-42",
  "run_id": "r-9",
  "generation": 3,
  "sequence": 18,
  "type": "output_delta",
  "payload": "..."
}

其中:

  • generation 表示同一会话的回复版本;
  • sequence 保证同一版本内部有序;
  • run_id 区分不同执行;
  • session_id 用于路由和权限检查。

1. 新输入如何使旧生成失效

假设用户快速输入两次:

generation=1:帮我查订单
generation=2:不用查了,改为总结退款政策

如果 generation=1 的模型仍在生成,系统应把它标记为取消或过期:

旧任务:running -> cancelled
新任务:queued -> running

旧任务即使稍后产生输出,也必须在发送器处被丢弃:

deliver(event)    event.generation=session.current_generationdeliver(event) \iff event.generation = session.current\_generation

这比仅取消模型请求更可靠,因为外部请求不一定能立刻停止,网络缓冲区中也可能已经存在旧事件。

2. 发送器故障不应回滚业务写入

以下顺序可能发生:

更新工单成功
发送通知失败

这不是一个原子操作。系统应拆成两个状态:

业务状态:committed
通知状态:pending / failed

随后由发送器重试通知,或者把通知写入可靠事件队列。若通知操作可重复,必须使用幂等键:

idempotency_key = f"{run_id}:ticket-updated:{ticket_id}"

否则发送器重试可能导致重复邮件、重复消息或重复 webhook。


十一、故障路径:失败不是一种状态

Agent 调度中至少要区分以下失败:

失败类型 示例 是否应重试
排队失败 队列已满 通常由上游降级或稍后重试
过期 用户已发送新请求 不应重试
取消 用户主动停止 不应重试
限流 供应商返回速率超限 退避后重试
网络瞬态错误 连接断开 仅幂等操作可重试
工具业务错误 订单不存在 通常不重试
权限错误 无权访问资源 不应盲目重试
版本冲突 会话状态已更新 重新读取或丢弃
Worker 崩溃 进程被杀 需要租约恢复

1. 租约防止 Worker 崩溃后任务永久丢失

持久化任务通常使用租约:

queued
  -> leased(worker=W1, expires_at=t)
  -> completed

如果 W1 崩溃,租约过期后任务回到 queued

leased -> retryable

但恢复任务必须结合幂等键。否则 W1 可能已经完成外部写入,只是在提交结果前崩溃;新 Worker 重试会重复写入。

2. 重试预算必须独立计费

一次失败任务可能产生:

首次模型调用
+ 工具调用
+ 第二次模型调用
+ 重试模型调用

因此重试不是免费的控制流。应为每个 Run 设置:

max_attempts
max_retry_tokens
max_elapsed_time
max_tool_calls

并且把重试次数纳入资源预算,否则“自动恢复”可能成为成本放大器。


十二、多 Agent 调度:并不是 Worker 越多越好

Anthropic 的 orchestrator-workers 模式由一个中心模型动态拆分任务、交给 Worker,再综合结果;它适合子任务无法预先确定的复杂问题。与固定的并行化相比,它的差异在于:子任务不是事先写死,而是由协调者根据输入动态决定。(anthropic.com)

但多 Agent 会引入新的调度成本:

主 Agent
├── 规划调用
├── Worker A
├── Worker B
├── Worker C
└── 汇总调用

如果每个 Worker 又动态生成更多 Worker,任务数量可能呈树状增长:

NdepthbdepthN_{depth} \leq b^{depth}

其中:

  • bb 是平均分支数;
  • depth 是最大委派深度;
  • NN 是潜在任务数。

即使实际任务没有完全展开,也必须设置:

  • 最大委派深度;
  • 单个 Run 最大子任务数;
  • 单个 Session 最大活动任务数;
  • 全局任务生成速率;
  • 子任务总 token 预算。

1. Worker 之间的结果合并

只读分析结果通常可以并行收集,但合并器必须面对:

  • 结果互相矛盾;
  • 某个 Worker 超时;
  • 某个 Worker 使用过时上下文;
  • 不同 Worker 对同一资源提出写入建议;
  • 部分结果已提交、部分结果未提交。

安全的模式是:

Worker 只产生 proposal
        ↓
验证器检查
        ↓
单独的提交器执行写入

这样可以把“分析并行”和“副作用提交”分开:

并行:查资料、生成计划、验证候选方案
串行:扣款、发布、删除、修改共享状态

OpenAI Agents SDK 支持 Agent 之间通过 agents-as-tools 或 handoffs 组织协作,但这些机制解决的是 Agent 责任转移和调用关系;会话级公平、资源配额和写入锁仍然需要应用运行时定义。(developers.openai.com)


十三、常见错误及其失败表现

错误一:对所有任务使用一个全局锁

表现是数据安全,但吞吐极低:

查询用户 A
等待
查询用户 B
等待
查询订单 C
等待

正确边界通常是资源级锁或会话级锁,而不是全局锁。

错误二:所有工具调用都并发

表现是:

  • 写入覆盖;
  • 重复扣款;
  • 发送通知早于业务提交;
  • 第二个工具读取到不一致状态;
  • 外部 API 触发竞态错误。

应先分析依赖图和副作用,再决定并行层级。

错误三:只限制 Worker 数量

表现是 Worker 数量稳定,但:

  • token 使用失控;
  • 单个任务产生大量子任务;
  • 某个租户长期占满资源;
  • 重试使实际请求量超过预期。

必须同时控制并发、速率和总量。

错误四:只在前端取消任务

表现是用户界面显示“已停止”,后端仍继续调用模型和工具。取消应贯穿:

入队 -> 出队 -> 资源等待 -> 工具调用 -> 状态提交 -> 事件发送

工具本身若不能取消,则至少在副作用前再次检查取消状态,并使用幂等键。

错误五:把旧输出发送给客户端

表现是用户看到回复倒退、内容交错或“已经取消的回答又出现”。发送器必须验证会话版本、Run 状态和事件序号。

错误六:把任务优先级交给模型完全决定

表现是提示词注入导致普通任务插队,或某类任务被模型持续标为紧急。模型可以提供语义信号,但最终优先级必须由服务端策略裁决。


十四、如何验证调度器是否真的可靠

单元测试不足以覆盖并发错误。至少应验证以下性质。

1. 会话互不串线

并发提交多个会话,检查:

结果中的 session_id
工具调用携带的 tenant_id
状态提交的版本
发送事件的目标连接

任何一个字段错配都应使测试失败。

2. 同一资源写入保持顺序

构造:

W1(resource=R, value=1)
W2(resource=R, value=2)

检查最终结果是否符合定义的顺序,而不是仅检查两个任务都返回成功。

3. 不相关资源可以并发

构造两个长时间任务:

W1(resource=A)
W2(resource=B)

检查总耗时接近 max(duration(W1), duration(W2)),而不是两者之和。这里测试的是并行边界是否过窄。

4. 过期任务不会产生副作用

让任务在队列中等待到截止时间,验证:

工具未被调用
状态未被提交
事件未被发送
配额已释放

5. 低优先级任务不会永久饥饿

持续注入高优先级任务,检查低优先级任务的最大等待时间是否有上界,或者老化机制是否生效。

6. Worker 崩溃可恢复且不重复写入

在外部写入完成、任务确认之前杀死 Worker,等待租约恢复,再检查:

外部资源只改变一次
任务最终状态可判定
重试次数有上限

十五、生产取舍:先确定语义,再选择实现

一个可操作的决策顺序是:

  1. 先定义状态边界:哪些状态属于会话,哪些属于 Run,哪些属于外部资源。
  2. 再标注工具属性:只读、写入、可交换、幂等、可取消还是不可取消。
  3. 再建立依赖图:明确并行层和串行层。
  4. 再决定队列维度:全局、租户、会话、资源是否需要独立队列。
  5. 再配置公平策略:优先级、轮转、老化和权重如何组合。
  6. 最后设置配额:并发、速率、token、工具次数和总时长分别限制。

如果任务主要是固定步骤,代码工作流往往比动态多 Agent 更容易验证。Anthropic 明确区分了 workflow 与 agent,并指出 Agent 的自主性会带来更高成本和错误累积风险;OpenAI 也将“由应用自行控制循环”与“由 Agents SDK 管理 Agent 循环”区分开来。(anthropic.com)

最终,Agent 并发的正确抽象不是“启动更多模型请求”,而是一个带有状态、依赖、资源和时效性的调度系统:

会话隔离
  保证谁能看到和修改什么

依赖调度
  保证哪些操作可以同时做

任务池与背压
  保证系统不会因输入过快而失控

优先级与公平
  保证重要任务及时执行,同时避免主体饥饿

资源配额
  保证单个会话、租户和全局成本可控

版本、取消和过期
  保证已经失去意义的结果不会继续产生影响

这些机制共同决定了 Agent 系统在高并发、长任务、工具失败和客户端断连时,是否仍然保持可解释、可恢复和可验证。


系列导航与关联阅读

官方资料

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