Agent 工程体系 · 第 68/98 篇。内容以 2026 年 9 月可验证的公开规范和稳定接口为基线;框架版本敏感能力会明确标注,不把实验行为写成通用保证。
Agent 并发与调度:会话隔离、任务池、优先级、公平和资源配额
Agent 并发不是简单地把多个 async 函数同时执行。一个生产级 Agent 系统通常同时面对四种不同的并发:
- 会话并发:多个用户、多个会话同时运行。
- 任务并发:一个会话内部同时推进多个子任务。
- 工具并发:一次模型决策产生多个工具调用。
- 输出并发:生成、事件推送、消息发送和客户端连接状态彼此独立。
如果只在最外层增加线程数或协程数,系统可能获得更高吞吐,却同时失去会话隔离、写入一致性、优先级控制和成本边界。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,则当前结果不能直接覆盖状态,必须:
- 丢弃;
- 重新基于最新状态生成;
- 作为独立消息发送,不写入主历史;
- 或进入人工/程序合并流程。
形式化地,状态更新可以写成:
其中:
- 是会话状态;
- 是任务读取状态时的版本;
- 是任务产生的状态变更;
conflict表示任务基于旧状态,不能直接提交。
版本检查只能防止静默覆盖,不能自动解决语义冲突。 两个任务都给同一账户退款时,即使版本检查正确,也仍需要业务幂等键和事务约束。
三、工具调用并发:依赖图决定执行顺序
1. 从工具列表建立依赖图
模型返回多个工具调用时,不应默认全部并发,也不应默认全部串行。应先构造有向图:
其中:
- 是工具调用;
- 是依赖边;
- 表示工具调用 必须等待 完成。
例如:
fetch_user ───────┐
├── build_report ── send_report
fetch_orders ─────┘
可以同时执行 fetch_user 和 fetch_orders,但 build_report 必须等待二者完成。
若图无环,则可按拓扑层级执行:
第 0 层:fetch_user, fetch_orders
第 1 层:build_report
第 2 层:send_report
若图中出现环:
A -> B -> C -> A
则说明规划结果无法直接执行,应拒绝、改写或中止,而不是让 Worker 永久等待。
2. 只读并发
只读工具只观察外部状态,不产生可见副作用。例如:
- 查询用户资料;
- 查询订单;
- 读取知识库;
- 获取监控指标;
- 读取文件内容。
如果两个工具满足以下条件,就可以并发:
还应补充一个经常被忽略的条件:二者不能依赖同一个会改变结果的外部快照。如果一个查询要求“在同一数据库快照上读取”,就需要事务快照或显式版本号,而不是普通并发请求。
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"})
但只有满足以下条件时才可合并:
- 操作满足幂等或可交换;
- 中间状态不会触发外部事件;
- 最终状态足以表达全部意图;
- 合并过程不会丢失条件检查。
反例:
扣款 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 ───────┘
执行过程:
| 时刻 | 状态 | 可执行任务 |
|---|---|---|
| T1、T2 就绪 | T1、T2 并发 | |
| T1 完成,T2 未完成 | 等待 T2 | |
| T1、T2 都完成 | 执行 T3 | |
| T3 成功 | 执行 T4 | |
| 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”负责:
- 读取会话上下文;
- 调用模型;
- 解析工具调用或子任务;
- 生成依赖关系;
- 把可执行节点放入执行队列。
“执行 Worker”负责:
- 获取一个已获准执行的任务;
- 检查取消和过期状态;
- 获取资源锁或配额;
- 调用工具;
- 记录结果;
- 释放资源并通知依赖节点。
“发送器”负责:
- 从结果事件流读取事件;
- 判断事件是否属于当前会话版本;
- 检查客户端是否仍连接;
- 按序发送;
- 在发送失败时重连、重放或丢弃。
将发送器独立出来很重要。模型生成完成不代表客户端一定还能接收;如果执行 Worker 直接写 WebSocket,网络阻塞会反向占用执行资源。
五、背压和过期丢弃:队列不能无限增长
1. 背压是什么
当任务进入速度大于系统处理速度时:
队列长度会不断增长。若没有背压,系统会以以下方式失败:
- 内存被队列占满;
- 所有请求延迟变大;
- 旧任务占用 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. 为什么必须丢弃过期任务
不是所有排队任务都值得执行。例如:
- 用户已经发送了新的问题;
- 前端已关闭页面;
- 语音转写结果已经被更新;
- 搜索建议已经超过有效时间;
- 旧的流式生成事件已经被新版本替代。
可以给任务定义截止时间:
但“过期检查”必须发生在多个阶段:
- 入队前;
- 从队列取出后;
- 获取资源配额前;
- 调用外部工具前;
- 结果提交前;
- 发送前。
只在入队时检查是不够的,因为任务可能在队列中等待很久。
一个安全的提交条件可以写成:
即使模型调用已经完成,如果会话已取消或事件版本过旧,也不应发送。
六、优先级:决定“先做谁”,不等于“永远插队”
1. 优先级队列的基本定义
优先级队列按照任务优先级选择下一个任务。若任务 的优先级高于任务 ,则通常希望:
其中 表示在资源可用时,优先调度 。
实际排序通常还需要加入创建时间和序列号:
排序键 = (-priority, created_at, sequence)
负号表示优先级越大越靠前;sequence 用于保证相同优先级下的稳定顺序。
2. 纯优先级的反例:饥饿
假设有两个队列:
低优先级:L1, L2, L3
高优先级:H1, H2, H3, ...
只要高优先级任务持续到达,调度器就永远执行 H:
H1 -> H2 -> H3 -> H4 -> ...
L1 从未执行
这称为饥饿。因此,“高优先级优先”不应等于“低优先级没有服务保证”。
3. 老化机制
一种简单方法是让等待时间提升有效优先级:
其中:
base_priority是任务初始优先级;wait(t)是等待秒数;- 是老化周期;
- 每等待 秒,有效优先级增加 1。
例如:
| 任务 | 初始优先级 | 等待时间 | 秒 | 有效优先级 |
|---|---|---|---|---|
| H1 | 10 | 0 秒 | 0 | 10 |
| M1 | 5 | 20 秒 | 2 | 7 |
| L1 | 1 | 100 秒 | 10 | 11 |
此时 L1 可能超过 H1,但这也可能不符合安全要求。因此老化通常设置上限:
对于支付、删除、权限变更等敏感操作,优先级只能影响排队顺序,不能绕过授权和人工审批。
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. 资源公平的形式化
令任务 的估计资源成本为 ,主体 的累计服务量为:
若主体权重为 ,调度器可以优先选择:
选择 score_k 最低的主体,再从其队列中取任务。
直觉是:权重越大,允许获得的服务量越多;但不能只统计任务个数,而要统计模型 token、工具执行时间、GPU 时间或其他受限资源。
4. 公平不等于所有任务相同待遇
在线交互请求和离线批处理不应共享完全相同的队列。常见划分是:
交互队列:短任务、低延迟、有限重试
后台队列:长任务、可暂停、低优先级
控制队列:取消、暂停、人工审批
控制任务通常需要“优先于普通任务”,否则用户点击取消后,取消消息也要排队等待长任务结束,系统就失去了可控性。
八、资源配额:限制并发数量还不够
1. 并发上限和资源配额的区别
max_workers = 20 只限制同时运行的任务数,不限制:
- 每个任务消耗多少 token;
- 一次运行最多调用多少工具;
- 一个会话可以生成多少子任务;
- 一个租户可以使用多少模型请求;
- 某个工具的并发连接数;
- 单次任务的最大执行时长;
- 失败重试产生的额外成本。
因此需要至少三类配额:
- 数量配额:同时运行多少任务;
- 速率配额:单位时间允许多少请求或 token;
- 总量配额:一次 Run 或一个 Session 最多使用多少资源。
2. 分层配额
可以建立以下层级:
全局
├── 租户
│ ├── 用户
│ │ └── 会话
│ └── 工具类别
└── 模型供应商
任务只有在所有相关配额都允许时才能启动:
如果任务同时调用模型和数据库工具,它可能需要同时获得:
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)
令牌桶适合限制时间窗口内的速率。设:
- 桶容量为 ;
- 令牌产生速率为 ;
- 每次请求消耗 个令牌。
当当前令牌数小于 时,请求必须等待或被拒绝。令牌桶允许短时突发,但长期平均速率不会超过 。
只使用信号量的反例:
每个请求都很慢,但同时只有 8 个请求
系统可能不超出并发数,却持续超出供应商的每分钟 token 限制。只使用令牌桶的反例:
请求速率合法,但每个请求都持有一个数据库连接
最终仍可能耗尽连接池。因此实际系统通常组合多种限制。
4. 配额扣减的时机
配额至少有三种状态:
reserved -> consumed
\-> released
reserved:任务获准启动,资源已预留;consumed:实际使用量已经确认;released:任务取消、过期或失败,未使用部分归还。
若等任务完成后才扣减配额,多个任务可能同时看到“余额充足”,造成超卖。若一开始按最大值扣减,又可能因为任务提前结束而浪费容量。
常见折中是:
- 启动前预留估计值;
- 执行中按实际使用追加;
- 完成后释放未使用部分;
- 对无法准确估计的资源使用保守上限。
九、一个可运行的调度器骨架
下面的示例演示四个机制:
- 同一会话串行;
- 不同会话并发;
- 优先级排序;
- 任务过期丢弃。
它不是模型 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
旧任务即使稍后产生输出,也必须在发送器处被丢弃:
这比仅取消模型请求更可靠,因为外部请求不一定能立刻停止,网络缓冲区中也可能已经存在旧事件。
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,任务数量可能呈树状增长:
其中:
- 是平均分支数;
depth是最大委派深度;- 是潜在任务数。
即使实际任务没有完全展开,也必须设置:
- 最大委派深度;
- 单个 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,等待租约恢复,再检查:
外部资源只改变一次
任务最终状态可判定
重试次数有上限
十五、生产取舍:先确定语义,再选择实现
一个可操作的决策顺序是:
- 先定义状态边界:哪些状态属于会话,哪些属于 Run,哪些属于外部资源。
- 再标注工具属性:只读、写入、可交换、幂等、可取消还是不可取消。
- 再建立依赖图:明确并行层和串行层。
- 再决定队列维度:全局、租户、会话、资源是否需要独立队列。
- 再配置公平策略:优先级、轮转、老化和权重如何组合。
- 最后设置配额:并发、速率、token、工具次数和总时长分别限制。
如果任务主要是固定步骤,代码工作流往往比动态多 Agent 更容易验证。Anthropic 明确区分了 workflow 与 agent,并指出 Agent 的自主性会带来更高成本和错误累积风险;OpenAI 也将“由应用自行控制循环”与“由 Agents SDK 管理 Agent 循环”区分开来。(anthropic.com)
最终,Agent 并发的正确抽象不是“启动更多模型请求”,而是一个带有状态、依赖、资源和时效性的调度系统:
会话隔离
保证谁能看到和修改什么
依赖调度
保证哪些操作可以同时做
任务池与背压
保证系统不会因输入过快而失控
优先级与公平
保证重要任务及时执行,同时避免主体饥饿
资源配额
保证单个会话、租户和全局成本可控
版本、取消和过期
保证已经失去意义的结果不会继续产生影响
这些机制共同决定了 Agent 系统在高并发、长任务、工具失败和客户端断连时,是否仍然保持可解释、可恢复和可验证。
系列导航与关联阅读
- 系列入口:Agent 工程完整路线:从运行循环、记忆与协议到安全、评测和生产交付
- 上一篇:事件驱动 Agent:Topic、消费者、关联 ID、顺序和最终一致性
- 下一篇:Agent 持久化执行:事件日志、Checkpoint、租约、恢复和确定性
- 延伸:Agent 并行工具调用:依赖图、只读并发、写入串行和合并
- 延伸:Agent 队列与背压:会话任务、生成 Worker、发送器和过期丢弃
官方资料
本文依据 Agent、模型、协议与框架官方资料重新梳理;正文、示例与生产清单由 WR BLOG 编写。

评论
0 条讨论