Agent 工程体系 · 第 67/98 篇。内容以 2026 年 9 月可验证的公开规范和稳定接口为基线;框架版本敏感能力会明确标注,不把实验行为写成通用保证。
事件驱动 Agent:Topic、消费者、关联 ID、顺序和最终一致性
事件驱动 Agent 不是“让多个 Agent 互相发消息”这么简单。真正需要解决的是:一个事件应该被谁消费,如何判断它属于哪个任务,重复投递是否会造成重复副作用,多个事件的先后关系由谁保证,以及不同 Agent 看到的状态暂时不一致时,系统如何最终收敛。
OpenAI 将 Agent 描述为能够规划、调用工具、协作并保留足够状态以完成多步骤工作的应用;Anthropic 则区分了固定代码路径编排的 workflow 与由模型动态决定过程和工具调用的 agent。无论采用哪种形式,一旦任务跨越多个进程、Worker 或 Agent,可靠性问题就不再由一次模型调用解决,而要由事件、状态和执行协议共同解决。(developers.openai.com)
1. 先建立一个准确模型:事件驱动 Agent 到底驱动什么
1.1 Agent 的一次运行不是一个函数调用
设一次用户任务为:
它可能经历以下步骤:
- 接收用户请求;
- 判断任务类型;
- 查询订单、知识库或内部系统;
- 让一个或多个专门 Agent 分析;
- 生成回复或执行外部动作;
- 发送结果;
- 将执行过程和最终状态写入持久化存储。
在同步程序中,这些步骤可能都发生在一次 HTTP 请求里:
HTTP 请求
└── 调用模型
└── 调用工具
└── 更新数据库
└── 返回响应
但在生产环境中,模型调用可能持续几十秒甚至更久,外部工具可能超时,服务实例可能重启,用户也可能在任务完成前再次发送消息。因此更合理的结构是:
用户请求
└── 创建任务和事件
└── Topic
├── 路由消费者
├── 查询消费者
├── 生成消费者
└── 发送消费者
这里的事件不是“模型思考内容”,而是系统中已经发生、需要其他组件观察或处理的事实。例如:
{
"event_type": "support.requested",
"event_id": "evt_01J...",
"task_id": "task_01J...",
"conversation_id": "conv_01J...",
"correlation_id": "corr_01J...",
"causation_id": "evt_01H...",
"aggregate_type": "support_case",
"aggregate_id": "case_123",
"aggregate_version": 7,
"occurred_at": "2026-09-01T10:00:00Z",
"payload": {
"text": "我的订单为什么还没有退款?"
}
}
事件应描述“已经发生的事实”,而不是“请某个 Agent 做某事”的模糊指令。后者也可以通过事件传递,但事件类型应该表达清楚的业务意图,例如 refund.review.requested,而不是 agent.do_something。
1.2 事件驱动不等于模型自主决策
事件总线只负责传递事实或任务信号,不负责替代 Agent 的推理。一个典型消费者可以执行如下闭环:
读取事件
↓
加载任务状态和当前版本
↓
调用模型或工具
↓
验证结果
↓
提交状态变更
↓
发布后续事件
Anthropic 对 Agent 的描述也强调了这一点:Agent 通常是在工具调用和环境反馈之间循环,通过每一步获得“ground truth”,并在完成、遇到阻塞或达到最大迭代次数时停止。(anthropic.com)
因此,不应把模型输出直接当作最终事实:
模型说“退款已经完成”
不能自动等价于:
支付系统确认退款成功
后者必须由支付系统返回、被系统持久化,并通过事件通知其他组件。
2. Topic:不是一个“消息列表”,而是事件的路由边界
2.1 Topic 的定义
Topic 是一组具有相同订阅语义、保留策略和路由规则的事件流。可以把它理解为:
其中:
name:Topic 名称;retention:事件保留多久;partitioning:事件如何分区;subscriptions:哪些消费者组订阅;delivery_policy:重试、确认、死信等投递规则。
例如:
agent.commands
agent.domain-events
agent.projections
agent.notifications
agent.audit
这些 Topic 不应只按技术组件命名,还应体现事件的生命周期和消费目的。
2.2 命令 Topic 与事实 Topic
在 Agent 系统中,至少要区分两类消息。
命令
命令表示:
希望某个组件尝试执行某个动作。
例如:
refund.review.requested
reply.generation.requested
notification.send.requested
命令通常有一个明确的目标消费者。命令可能失败,也可能被拒绝。
事实事件
事实事件表示:
某个动作已经发生,其他组件可以据此更新自己的状态。
例如:
refund.review.completed
refund.approved
refund.rejected
reply.generated
notification.sent
事实事件不应该假设所有消费者都必须成功执行同一种动作。审计消费者、统计消费者、通知消费者可能分别处理同一个事实。
一个常见错误是把所有消息都命名成“动作”:
process_order
handle_message
run_agent
这会让消费者无法判断消息是请求、结果、失败还是重试。更明确的事件名称通常包含状态变化:
conversation.message.received
research.plan.created
research.source.fetched
answer.generation.requested
answer.generated
answer.delivery.failed
2.3 Topic 不等于消费者实例
Topic 是发布端和消费端之间的逻辑边界;消费者是实际执行处理逻辑的进程或任务。
同一个 Topic 可以有多个消费者组:
agent.domain-events
├── support-agent-group
├── analytics-group
├── audit-group
└── notification-group
同一消费者组中的多个实例通常用于水平扩展,但一个事件在同一消费者组内通常只应由一个实例负责处理。不同消费者组则各自收到一份逻辑副本。
因此:
“这个事件被消费过了”
必须明确是:
“被 support-agent-group 消费过了”
因为审计组可能还没有消费,通知组也可能正在重试。
2.4 一个 Topic 是否应该承载所有 Agent 事件
不应把所有消息塞入一个 Topic:
agent-events
这种设计会造成:
- 无关消费者被迫读取大量消息;
- 不同消息需要不同保留时间,却无法独立配置;
- 高优先级消息可能被低优先级任务阻塞;
- 事件权限边界不清晰;
- 重试策略无法按业务风险区分。
但也不能为每一个事件类型创建一个 Topic。Topic 数量过多会使运维、监控和权限管理复杂化。
较合理的划分依据是:
- 消费者是否相同;
- 顺序键是否相同;
- 保留时间是否相同;
- 重试和死信策略是否相同;
- 数据敏感级别是否相同;
- 是否允许独立扩缩容。
3. 消费者:一次处理必须同时面对重复、崩溃和部分成功
3.1 消费者的最小生命周期
一个消费者处理事件通常经历以下阶段:
available
↓
claimed / leased
↓
processing
├── succeeded → acknowledged
├── retryable failure → delayed retry
├── permanent failure → dead-lettered
└── lease timeout → redelivered
其中 acknowledged 不应简单理解为“函数返回了”。它应该表示:
处理结果已经可靠持久化,或者该消息已经被明确判定为无需再次处理。
错误的确认顺序:
收到消息
↓
立即 ACK
↓
调用模型
↓
进程崩溃
这会导致消息永久丢失。
通常更安全的顺序是:
收到消息
↓
执行处理
↓
提交状态和副作用记录
↓
ACK
但这仍然无法自动解决外部副作用问题。假设消费者完成了支付退款请求,随后在 ACK 前崩溃,消息再次投递后就可能再次退款。
所以,可靠消费者不能只依赖 ACK,而需要幂等性。
3.2 至少一次投递与幂等性
在存在网络超时、消费者崩溃和租约过期时,消息系统通常更容易保证:
即“消息至少送达一次”,而不是严格只送达一次。
如果同一个事件可能被处理多次,消费者必须满足:
这里的等价不是指执行过程完全相同,而是指最终可观察状态和外部副作用不会被重复放大。
最简单的去重表:
CREATE TABLE inbox (
consumer_group TEXT NOT NULL,
event_id TEXT NOT NULL,
received_at TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP,
result TEXT,
PRIMARY KEY (consumer_group, event_id)
);
处理逻辑:
BEGIN;
INSERT INTO inbox(consumer_group, event_id, result)
VALUES ('reply-generator', 'evt_123', 'processing')
ON CONFLICT (consumer_group, event_id) DO NOTHING;
-- 如果插入行数为 0,说明这个消费者组已经处理过该事件
-- 此时直接提交并 ACK
-- 执行业务处理
UPDATE agent_tasks
SET status = 'reply_generated',
version = version + 1
WHERE task_id = 'task_123'
AND version = 4;
UPDATE inbox
SET result = 'succeeded'
WHERE consumer_group = 'reply-generator'
AND event_id = 'evt_123';
COMMIT;
这里有一个关键边界:inbox 去重只能防止同一个消费者重复执行同一段数据库事务,不能自动防止外部系统重复副作用。
对于发送短信、调用支付、创建工单等外部动作,应使用业务幂等键:
idempotency_key = task_id + ":" + action_type + ":" + logical_attempt
例如:
task_123:refund:1
外部服务或本地发送器应记录该键:
CREATE TABLE side_effects (
effect_key TEXT PRIMARY KEY,
effect_type TEXT NOT NULL,
status TEXT NOT NULL,
provider_ref TEXT,
created_at TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP
);
如果请求超时,不能根据“没有收到响应”推断外部动作没有发生。正确做法是使用相同的幂等键重试,或先查询外部系统。
3.3 Consumer 的错误分类
消费者不能把所有异常都当作可重试错误。
| 错误类型 | 示例 | 处理方式 |
|---|---|---|
| 临时基础设施错误 | 数据库连接失败、模型服务 503 | 延迟重试 |
| 限流 | 429、配额暂时不足 | 按 Retry-After 或指数退避 |
| 输入错误 | 缺少 task_id、事件版本非法 |
进入死信,人工或修复后重放 |
| 业务拒绝 | 用户无权限、订单已关闭 | 发布业务失败事件,不应无限重试 |
| 模型输出不合规 | JSON 无法解析、违反结构约束 | 受限重试或转人工 |
| 外部副作用不确定 | 请求超时 | 使用幂等键查询或重试 |
错误重试的核心不是“多试几次”,而是判断:
对于不可重试错误无限重试,只会造成队列堆积和成本放大。
4. 关联 ID:把一次任务从日志、事件和副作用中串起来
4.1 四种 ID 不应混为一谈
事件驱动 Agent 至少需要区分以下标识。
event_id
当前事件本身的唯一标识:
evt_123
用途是去重、审计和重放。
correlation_id
同一个业务流程或用户请求链路的标识:
corr_456
一次用户请求触发多个 Agent、多个 Topic 和多个 Worker 时,它们通常共享 correlation_id。
causation_id
直接导致当前事件产生的上一个事件:
causation_id = evt_122
它描述因果边:
evt_122 ──caused──> evt_123
task_id
系统中可恢复、可查询的业务任务标识:
task_789
它对应持久化任务状态,不应依赖日志系统才能找到。
有些系统还需要:
conversation_id:用户会话;run_id:某个 Agent 运行实例;attempt:当前处理尝试次数;parent_run_id:父 Agent 运行;aggregate_id:被修改的业务实体;trace_id、span_id:分布式追踪。
4.2 关联 ID 的传播规则
事件传播时,推荐采用以下规则:
event_id 每次生成新事件都重新生成
correlation_id 同一业务流程内保持不变
causation_id 等于直接触发者的 event_id
task_id 同一可恢复任务内保持不变
run_id 每次 Agent 运行或恢复运行重新生成
例如:
用户消息
event_id = evt_001
correlation_id = corr_001
task_id = task_001
规划完成
event_id = evt_002
correlation_id = corr_001
causation_id = evt_001
task_id = task_001
检索请求
event_id = evt_003
correlation_id = corr_001
causation_id = evt_002
task_id = task_001
回复生成
event_id = evt_004
correlation_id = corr_001
causation_id = evt_003
task_id = task_001
如果把 event_id 当作全链路 ID,每次生成新事件都复用它,结果会失去事件唯一性,去重逻辑也会误判。相反,如果每个组件都生成新的 correlation_id,一次任务就会被切成多个互不相干的调用链。
4.3 关联 ID 与安全边界
关联 ID 不是权限凭证,也不应包含敏感信息:
错误:
correlation_id = user_138_phone_13800138000
正确:
correlation_id = corr_01JABC...
日志中可以记录 task_id 和 correlation_id,但事件 payload 不应为了方便排查而携带完整身份证号、银行卡号或原始私密对话。关联信息和业务数据应分别治理。
5. 顺序:真正需要保证的是“哪个实体上的顺序”
5.1 全局顺序通常不是需求
假设系统收到两个会话的消息:
A1: 会话 A 的“取消订单”
B1: 会话 B 的“查询物流”
A2: 会话 A 的“确认取消”
通常只要求:
A1 < A2
并不要求:
A1 < B1 < A2
如果系统强行保证所有会话的全局顺序,单个慢任务就可能阻塞所有用户。
因此应先定义顺序域:
例如:
- 会话内顺序:
conversation_id; - 订单内顺序:
order_id; - 任务内顺序:
task_id; - 用户账户内顺序:
account_id。
5.2 分区顺序的充分条件
如果事件流按 order_key 分区,并且同一个分区内由单一逻辑序列处理,那么可以得到:
这里:
seq(e)是事件在该顺序域中的序号;apply(e_i) \prec apply(e_j)表示e_i先于e_j应用。
但这只有在以下条件成立时才有效:
- 同一
order_key始终路由到同一分区; - 消费者不会让同一分区中的事件并发越过前序事件;
- 事件序号由可靠的生产者或存储层生成;
- 重试不会让旧事件在新事件之后错误覆盖状态。
仅仅给事件加一个时间戳,不足以保证顺序。时钟可能漂移,事件可能异步发布,甚至后生成的事件可能携带更早的业务时间。
5.3 有序消费不等于有序完成
假设同一任务有两个事件:
seq=10 reply.generation.requested
seq=11 notification.send.requested
消费者虽然按顺序取出,但 seq=10 的模型调用耗时 20 秒,seq=11 可能在另一个 Worker 中先完成。此时:
完成顺序:seq=11 → seq=10
如果发送器只依据完成事件执行,就可能先发送一条尚未生成内容的通知。
解决方案有三类:
方案一:同一顺序域串行执行
同一个 task_id → 同一个 Worker 串行处理
简单但吞吐量低。
方案二:状态机校验前置状态
发送器只允许在状态满足条件时执行:
notification.send.requested
前置条件:reply.status = 'generated'
如果前置状态不满足,则延迟重试或暂存。
方案三:事件序号与版本门槛
事件携带:
{
"task_id": "task_123",
"required_version": 8
}
消费者读取当前状态版本,只有:
current_version >= required_version
时才继续处理。
5.4 并发更新的反例
初始状态:
order.status = "paid"
order.version = 5
两个 Agent 同时读取版本 5:
Agent A:决定取消订单
Agent B:决定修改收货地址
如果没有版本检查:
A 写入 status = "cancelled"
B 写入 address = "新地址"
最终数据库可能只保留 B 的旧快照,覆盖 A 的取消状态。
使用乐观并发控制:
UPDATE orders
SET status = :new_status,
version = version + 1
WHERE order_id = :order_id
AND version = :expected_version;
若影响行数为 0,说明版本已变化。此时不能盲目重试原操作,而应重新读取状态,让 Agent 或确定性业务逻辑重新判断。
6. 最终一致性:允许暂时不同,但不允许无条件错误
6.1 定义
在事件驱动系统中,通常存在:
- 命令服务的主状态;
- Agent 运行状态;
- 搜索索引;
- 会话投影;
- 通知状态;
- 审计日志;
- 统计数据。
这些状态不会在同一个数据库事务中同时更新。因此在时间 上,不同读取者可能看到不同版本:
最终一致性要求的是:如果后续事件能够持续送达、消费者能够最终成功处理,并且没有新的更新,那么存在某个有限时间 ,使得:
其中 是投影规则。
这不是“最终一定正确”的无条件承诺。它依赖于:
- 事件不会永久丢失;
- 消费者会重试或进入可恢复死信;
- 处理逻辑具备幂等性;
- 事件应用顺序满足业务需要;
- 投影逻辑能够处理重复和迟到事件;
- 主状态本身没有被非法并发覆盖。
6.2 最终一致性不能替代业务不变量
例如订单取消和支付退款存在以下不变量:
如果允许通知投影暂时落后,这通常没问题:
主库:订单已取消
通知投影:仍显示“处理中”
但如果退款服务先执行了退款,订单主状态却仍是“已支付”,就可能违反业务不变量。
因此要区分:
- 可暂时落后的读取模型:搜索索引、统计、通知列表;
- 必须原子保护的业务不变量:余额、库存、支付状态、权限;
- 必须由外部确认的事实:支付完成、短信发送、文件上传。
最终一致性适合传播状态,不适合掩盖原子性要求。
6.3 一个具体收敛算例
初始状态:
task.status = "received"
task.version = 1
projection.status = "received"
事件流:
E1: task.accepted
E2: plan.created
E3: answer.generated
E4: notification.sent
主任务状态的转换:
version=1, received
--E1-->
version=2, accepted
--E2-->
version=3, planned
--E3-->
version=4, generated
--E4-->
version=5, delivered
如果投影消费者先收到 E3,但没有收到 E1 和 E2,有三种处理方式:
- 根据事件 payload 直接构建状态;
- 发现版本间隙,暂存
E3; - 从事件存储补读缺失事件。
若事件包含 aggregate_version:
{
"event_type": "answer.generated",
"aggregate_id": "task_123",
"aggregate_version": 4
}
投影当前版本为 2 时,应识别:
expected = 3
received = 4
这不是普通业务失败,而是版本间隙。继续应用可能导致投影丢失中间状态。
7. 多 Agent 共享状态:所有权必须先于协作
7.1 共享状态的三层划分
多 Agent 系统经常把状态混成一个 JSON:
{
"messages": [],
"plan": {},
"order": {},
"draft": {},
"send_status": ""
}
然后所有 Agent 都可以读写。这会让共享状态变成无明确所有者的全局变量。
更可靠的划分是:
| 状态 | 所有者 | 其他 Agent 的访问方式 |
|---|---|---|
| 会话消息 | 会话服务 | 追加事件或读取投影 |
| 任务生命周期 | 任务协调器 | 发布命令、消费状态事件 |
| 检索结果 | 检索 Agent | 发布结果事件 |
| 回复草稿 | 生成 Agent | 通过版本化草稿接口更新 |
| 发送状态 | 发送器 | 由发送器独占写入 |
| 审计记录 | 审计服务 | 只追加 |
所有权表示谁有权改变状态,而不是谁可以读取状态。
7.2 状态变更的三种策略
单一所有者
一个 Agent 或服务独占某类状态:
reply-draft → generation-worker
send-status → sender
其他组件只能发布请求。这种方式最容易推理。
乐观并发控制
允许多个组件尝试更新,但必须携带期望版本:
read version=7
compute patch
UPDATE ... WHERE version=7
更新失败时重新读取并解决冲突。
锁
对短事务、强互斥资源可以使用数据库锁或分布式锁。但锁不应覆盖模型调用:
错误:
BEGIN
SELECT ... FOR UPDATE
调用模型 30 秒
UPDATE ...
COMMIT
这会长时间占用锁,并在消费者崩溃时放大故障。
更合理的是:
读取状态
提交“处理中”租约
释放数据库锁
调用模型
使用版本条件提交结果
锁解决的是“同一时刻谁能进入临界区”,不解决崩溃恢复、重复消息和外部副作用。
7.3 冲突不一定是异常
两个 Agent 的更新可能语义上互不冲突:
A 更新 reply.text
B 更新 reply.citations
如果状态按字段拆分,可以合并;但以下更新通常冲突:
A:reply.status = approved
B:reply.status = rejected
因此冲突策略应由领域定义:
- last-write-wins:适合低价值元数据,不适合支付和审批;
- field-level merge:适合独立字段;
- version rejection:让上层重新决策;
- state-machine transition:只允许合法状态转换;
- human review:无法自动判定时交给人。
8. 事件溯源:事件是事实记录,不是任意日志
8.1 事件溯源的核心
事件溯源将实体当前状态视为事件折叠结果:
例如:
State_0: received
E1: accepted → accepted
E2: plan.created → planned
E3: answer.generated → generated
E4: notification.sent → delivered
事件溯源的价值在于:
- 可以重建状态;
- 可以审计因果链;
- 可以重放到新投影;
- 可以定位某次 Agent 决策之前看到了什么。
但事件溯源要求事件具有稳定语义。以下内容不适合直接作为领域事件:
模型内部思考文本
临时 prompt
某次 HTTP 调用的完整响应
调试日志
这些内容可以作为运行记录或观测数据,但不应成为重建业务状态的唯一依据。
8.2 事件表与当前状态表可以并存
事件溯源不等于每次读取都扫描全部事件。常见结构是:
events 追加写事实
aggregate_state 保存当前状态快照
projection_* 面向查询的读取模型
写入流程:
事务开始
├── 校验 aggregate_version
├── 写入新事件
├── 更新 aggregate_state
└── 写入 outbox 事件
事务提交
读取流程:
优先读 aggregate_state 或 projection
必要时从 events 重建或校验
快照必须带版本:
snapshot.version = 120
如果从版本 120 继续重放事件,必须从 aggregate_version = 121 开始,而不是按客户端时间猜测。
9. 事务性 Outbox:解决“数据库更新成功但事件没发出去”
9.1 双写问题
下面的代码存在不可恢复窗口:
UPDATE task SET status = 'planned';
publish("plan.created");
如果数据库更新成功而进程在发布前崩溃:
主状态已更新
事件没有发布
下游 Agent 永远不会启动
反过来,如果事件发布成功而数据库提交失败,下游又会根据一个不存在的状态开始执行。
9.2 Outbox 的做法
将业务状态和待发布事件写入同一个数据库事务:
CREATE TABLE outbox (
id BIGSERIAL PRIMARY KEY,
event_id TEXT NOT NULL UNIQUE,
topic TEXT NOT NULL,
aggregate_id TEXT NOT NULL,
aggregate_version BIGINT NOT NULL,
payload JSONB NOT NULL,
published_at TIMESTAMP NULL,
attempts INTEGER NOT NULL DEFAULT 0,
next_attempt_at TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP
);
事务:
BEGIN;
UPDATE agent_tasks
SET status = 'planned',
version = version + 1
WHERE task_id = 'task_123'
AND version = 2;
INSERT INTO outbox (
event_id,
topic,
aggregate_id,
aggregate_version,
payload
)
VALUES (
'evt_plan_003',
'agent.domain-events',
'task_123',
3,
'{"event_type":"plan.created","task_id":"task_123"}'
);
COMMIT;
独立 Publisher 定期读取未发布记录:
SELECT id, event_id, topic, payload
FROM outbox
WHERE published_at IS NULL
AND next_attempt_at <= CURRENT_TIMESTAMP
ORDER BY id
LIMIT 100
FOR UPDATE SKIP LOCKED;
发布成功后标记:
UPDATE outbox
SET published_at = CURRENT_TIMESTAMP
WHERE id = :id;
这里仍然可能出现:
事件已成功发布
Publisher 在标记 published_at 前崩溃
于是事件会再次发布。因此 Outbox 通常提供的是:
数据库状态与“最终会发布事件”之间的一致性
而不是消息系统层面的严格只发一次。下游仍必须幂等。
10. 一个完整的事件驱动 Agent 架构
下面以“客服退款咨询”为例。系统包含:
- 会话服务:接收用户消息;
- 路由 Agent:判断问题类型;
- 订单 Agent:查询订单和退款状态;
- 回复生成 Worker:生成自然语言回复;
- 发送器:向用户发送消息;
- 审计消费者:保存事件链。
flowchart LR
U[用户] --> API[会话 API]
API --> DB[(任务状态库)]
API --> O[(Outbox)]
O --> P[事件发布器]
P --> T1[(agent.commands)]
P --> T2[(agent.domain-events)]
T1 --> R[路由 Agent]
T1 --> Q[订单 Agent]
T1 --> G[生成 Worker]
T1 --> S[发送器]
R --> T2
Q --> T2
G --> T2
S --> T2
T2 --> PR[会话/任务投影]
T2 --> AU[审计消费者]
T2 --> M[(指标与追踪)]
S --> EXT[消息渠道]
一次处理的事件链可能是:
conversation.message.received
↓
support.route.requested
↓
order.lookup.requested
↓
order.lookup.completed
↓
reply.generation.requested
↓
reply.generated
↓
notification.send.requested
↓
notification.sent
关键路径不是“每个 Agent 都直接调用下一个 Agent”,而是:
一个 Agent 发布事实或命令
另一个消费者根据事件决定是否继续
这样可以让每个组件独立重试、扩缩容和恢复。
10.1 事件信封
推荐把路由所需元数据放在固定信封中:
{
"event_id": "evt_01J8...",
"event_type": "reply.generation.requested",
"schema_version": 1,
"occurred_at": "2026-09-01T10:12:03.145Z",
"correlation_id": "corr_01J8...",
"causation_id": "evt_01J7...",
"task_id": "task_01J8...",
"conversation_id": "conv_01J8...",
"aggregate": {
"type": "agent_task",
"id": "task_01J8...",
"version": 8
},
"producer": "order-agent",
"attempt": 1,
"payload": {
"answer_context_ref": "ctx_01J8..."
}
}
payload 可以演进,但信封字段应尽量稳定。消费者应拒绝无法识别的 schema_version,或通过兼容解析器处理旧版本。
11. Agent 队列与背压:会话任务、生成 Worker、发送器和过期丢弃
11.1 背压是什么
当事件产生速度高于消费者处理速度时,队列长度增长:
其中:
- :时刻 的待处理任务数;
- :生产速率;
- :消费速率。
当长期满足:
队列就会持续增长,最终表现为延迟升高、内存耗尽、重试风暴或用户收到过时回复。
背压不是简单地“加更多 Worker”。如果瓶颈是模型 API 限流或发送渠道限流,盲目扩容只会增加失败率。
11.2 会话任务队列
会话消息通常不能无限并行。用户连续发送:
M1:帮我查订单
M2:订单号是 123
M3:我想知道什么时候退款
如果 M1、M2、M3 同时进入模型,Agent 可能在没有订单号时先做出错误判断。
会话任务队列应定义合并或覆盖规则:
同一 conversation_id:
- 旧任务尚未开始:可合并
- 旧任务正在模型调用:通常不强行取消
- 新消息到达:标记旧回复可能过期
- 只允许最新任务向用户发送最终回复
可使用会话版本:
conversation_version = 17
生成任务携带:
{
"conversation_id": "conv_1",
"input_version": 17
}
发送前检查当前版本:
SELECT version
FROM conversations
WHERE conversation_id = 'conv_1';
如果当前版本已经是 18,则版本 17 的回复可能已经过期,不应发送。
11.3 生成 Worker
生成 Worker 的职责是:
- 获取结构化上下文;
- 调用模型;
- 验证输出;
- 持久化草稿;
- 发布
reply.generated。
它不应直接向用户发送消息。否则生成和发送无法分别重试,发送失败时可能重复调用昂贵的模型。
一个较清晰的状态机:
queued
↓
generating
├── generated
├── retryable_failed
├── permanently_failed
└── expired
模型返回结构化结果时,应验证:
{
"reply_text": "退款已于 2026-08-31 提交,预计原路返回。",
"confidence": 0.91,
"requires_human_review": false,
"source_refs": ["order:123", "refund:456"]
}
不能只检查 JSON 是否可解析,还要检查:
reply_text是否存在;source_refs是否对应实际查询结果;requires_human_review是否触发审批流程;- 回复是否仍对应当前
conversation_version; - 是否泄露了内部工具结果。
11.4 发送器
发送器面对的是外部副作用,因此必须独立实现幂等和状态查询:
send.requested
↓
sender.claimed
↓
调用渠道
├── sent
├── retryable_failed
├── permanently_failed
└── unknown
unknown 很重要。网络超时后,系统并不知道渠道是否已经发送成功。
正确处理:
请求发送,幂等键 = task_123:reply:1
↓
超时
↓
使用同一幂等键查询渠道状态
├── 已发送 → 记录 sent
├── 未发送 → 使用同一幂等键重试
└── 无法确认 → 转人工或进入延迟队列
不应在每次超时后生成新的幂等键:
错误:
task_123:reply:1
task_123:reply:2
task_123:reply:3
这会把“重试同一个动作”错误地变成“执行三个不同动作”。
11.5 过期丢弃
事件驱动 Agent 经常产生“到达时已经没有价值”的任务。例如:
- 用户已经发送了新消息;
- 旧回复已经被更新版本替代;
- 订单状态已经变化;
- 预约时间已经过去;
- 通知内容已被人工修改。
应在事件中携带:
{
"created_at": "2026-09-01T10:00:00Z",
"expires_at": "2026-09-01T10:00:30Z",
"conversation_version": 17
}
消费者处理前检查:
now > expires_at
→ 标记 expired,不再调用模型或外部系统
current_conversation_version > event.conversation_version
→ 标记 superseded,不再发送
过期丢弃不是消息丢失,而是业务上明确判定:
继续执行该事件已经不能产生有效结果。
必须保留丢弃原因:
expired
superseded
cancelled
invalidated_by_state_change
否则监控中只看到“队列变短”,却不知道是成功消费还是大量丢弃。
12. 事件顺序与最终一致性的结合
12.1 不是所有消费者都需要同样的顺序
对同一事件流:
| 消费者 | 顺序要求 |
|---|---|
| 审计消费者 | 需要保留因果顺序 |
| 统计消费者 | 通常可乱序,按事件时间修正 |
| 搜索投影 | 需要实体版本检查 |
| 发送器 | 需要会话或任务顺序 |
| 推荐消费者 | 可能只关心最新状态 |
因此顺序是消费者的业务约束,而不是 Topic 自动提供的普遍属性。
12.2 “最新状态覆盖旧状态”的安全条件
如果投影采用 last-write-wins,至少需要:
只接受 aggregate_version 更大的事件
伪代码:
def apply_event(projection, event):
if event.aggregate_version <= projection.version:
return "ignored_as_old"
if event.aggregate_version != projection.version + 1:
return "version_gap"
projection = transition(projection, event)
projection.version = event.aggregate_version
return "applied"
如果只比较 occurred_at,可能出现如下反例:
E1:业务版本 10,时间 10:00:01
E2:业务版本 11,时间 09:59:59
由于时钟、批处理或跨服务时间来源不同,E2 可能被错误当成旧事件。实体版本比跨机器时间更适合作为状态应用依据。
12.3 可交换事件可以并行
若两个事件满足:
则它们对状态更新是可交换的,可以并行处理。
例如两个独立的统计计数增量:
message.received +1
tool.called +1
通常可以并行。
但以下事件不可交换:
refund.approved
refund.cancelled
它们必须按订单或退款单串行化,或由状态机拒绝非法转换。
13. OpenAI Agents SDK 与事件驱动外部编排的边界
OpenAI Agents SDK 负责 Agent 运行循环、工具调用、Agent 间 handoff、guardrails、追踪和可恢复运行状态;而 Responses API 更适合由应用自行控制模型交互、工具调用、循环和分支。也就是说,SDK 内部的 Agent loop 与系统外部的事件总线是两个不同层次:前者管理一次 Agent run,后者管理跨进程、跨服务的可靠执行。(developers.openai.com)
可以这样划分:
事件总线:
任务排队、重试、租约、路由、跨服务状态传播
Agents SDK:
当前 Agent 的模型循环、工具调用、handoff、guardrail
业务数据库:
任务状态、版本、幂等记录、Outbox、审计
发送器:
外部消息渠道的幂等副作用
不应把 SDK 的一次运行结果当作全系统事务提交。一个 Agent run 成功,只说明该运行完成了自己的执行路径;它不自动保证:
- 数据库投影已经更新;
- 下游消费者已经消费;
- 用户已经收到消息;
- 外部支付或工单系统已经完成动作。
Anthropic 也建议从简单实现开始,只在任务确实需要动态规划、工具使用或多步自主执行时引入更复杂的 Agent 结构;复杂度会带来额外延迟、成本和错误累积。(anthropic.com)
14. 失败路径:生产问题通常发生在边界,而不是模型调用内部
14.1 发布前崩溃
数据库更新成功
Outbox 未写入
如果业务状态和 Outbox 不在同一事务中,系统无法知道是否需要补发。修复方式是事务性 Outbox,或通过定期对账任务发现状态与事件的不一致。
14.2 发布后崩溃
事件已发布
published_at 尚未更新
结果是重复发布。下游必须使用 event_id 去重,而 Publisher 可以在重试时继续使用同一个 event_id。
14.3 消费成功但 ACK 丢失
数据库已提交
ACK 未送达
消息再次投递
这正是 Inbox 去重和业务幂等存在的原因。
14.4 模型调用完成但结果提交失败
模型已经产生结果
数据库提交失败
如果模型调用是纯计算,可以重试;如果模型调用触发了外部工具,则需要按工具类型分别处理。工具调用不能只靠“重新运行整个 Agent”解决,因为前面的工具副作用可能已经成功。
14.5 状态版本落后
Agent 读取 version=7
其他消费者提交 version=8
Agent 用 version=7 覆盖写入
必须通过条件更新拒绝旧版本,并重新加载状态。对于需要人工判断的冲突,不能让模型无上下文地重复执行。
14.6 旧回复晚到
用户消息 M1 → 生成慢
用户消息 M2 → 生成快并发送
M1 的回复随后完成
如果没有 conversation_version 检查,用户会先看到 M2 回复,再收到基于旧上下文的 M1 回复。
15. 诊断:先沿因果链查,不要只看队列长度
一次任务的排查路径应从 task_id 和 correlation_id 开始:
task_id
↓
事件列表
↓
每个消费者组的处理记录
↓
Agent run / tool call
↓
数据库版本变化
↓
Outbox 发布状态
↓
外部发送器幂等记录
建议至少记录以下字段:
event_id
event_type
consumer_group
task_id
correlation_id
causation_id
aggregate_id
aggregate_version
attempt
status
error_class
latency_ms
model_run_id
tool_call_id
需要区分几个常被混淆的指标:
- 队列深度:还有多少未处理消息;
- 消费延迟:事件产生到开始处理的时间;
- 处理耗时:开始处理到完成的时间;
- 重试率:处理失败后再次尝试的比例;
- 版本间隙数:投影遇到的缺失序号;
- 过期丢弃率:任务到达时已经无效的比例;
- 未知副作用数:外部请求超时且无法确认结果的数量。
例如队列深度下降并不一定代表系统恢复,也可能是消费者大量把消息标记为 expired 或 dead-lettered。监控必须把成功、跳过、过期和死信分开。
16. 常见误解与边界
误解一:有了 Topic 就有了可靠性
Topic 只提供传输和订阅边界。可靠性还需要:
持久化事件
+ 明确投递语义
+ 幂等处理
+ 重试与死信
+ 状态版本
+ 观测和对账
误解二:关联 ID 能保证顺序
关联 ID 只能帮助定位同一链路,不能让消息自动按顺序处理。顺序需要分区键、序号、串行消费或版本门槛。
误解三:ACK 就等于业务成功
ACK 只表示消费者认为消息不需要再次投递。业务成功必须由状态提交和副作用确认定义。
误解四:最终一致性就是“稍后会正确”
如果事件丢失、消费者永久失败、旧事件覆盖新状态,系统可能永远不收敛。最终一致性是带前提的收敛性质,不是对错误设计的免责。
误解五:给所有事件加锁即可解决并发
锁不能覆盖长时间模型调用,也不能解决进程崩溃后的重复执行。对于 Agent 系统,版本、幂等键、状态机和补偿逻辑往往比长锁更重要。
误解六:多 Agent 越多越可靠
多 Agent 适合职责确实不同、可以并行或需要不同工具和策略的任务。Anthropic 将路由、并行化、orchestrator-workers 和 evaluator-optimizer 视为不同的编排模式,而不是默认都应使用的复杂结构。(anthropic.com)
17. 一套可落地的最小协议
如果要为一个新的事件驱动 Agent 系统定义基线,可以先固定以下协议。
事件必须包含
event_id
event_type
schema_version
occurred_at
correlation_id
causation_id
task_id
aggregate_id
aggregate_version
payload
生产者必须保证
业务状态更新与 Outbox 写入在同一事务中
同一逻辑事件重试时不随意更换 event_id
事件类型和 schema_version 可演进
消费者必须保证
处理成功后再 ACK
按 consumer_group + event_id 去重
对可重试和不可重试错误分类
对状态更新使用版本条件
对外部副作用使用业务幂等键
顺序必须明确
顺序域是什么:conversation_id、task_id 还是 aggregate_id
哪些事件必须串行
哪些事件可以并行
迟到事件如何处理
版本间隙如何修复
队列必须具备
最大并发
租约或可见性超时
重试次数和退避
死信队列
过期时间
取消和 superseded 状态
状态必须具备
明确所有者
版本号
合法状态转换
冲突策略
恢复入口
事件驱动 Agent 的核心不是把同步调用改成异步消息,而是把一次不可见的长流程拆成一组可观察、可重试、可恢复、可验证的状态转换。Topic 负责隔离事件流,消费者负责执行和确认,关联 ID 负责重建因果链,顺序机制负责保护同一实体上的状态演进,最终一致性则负责让不同读取模型在事件持续处理后收敛。
当这些概念被分别定义后,多 Agent 协作就不再依赖“希望消息按预期到达”,而是建立在明确的所有权、版本、幂等和恢复协议之上。
系列导航与关联阅读
- 系列入口:Agent 工程完整路线:从运行循环、记忆与协议到安全、评测和生产交付
- 上一篇:多 Agent 共享状态:所有权、版本、冲突、锁和事件溯源
- 下一篇:Agent 并发与调度:会话隔离、任务池、优先级、公平和资源配额
- 延伸:Agent 队列与背压:会话任务、生成 Worker、发送器和过期丢弃
官方资料
本文依据 Agent、模型、协议与框架官方资料重新梳理;正文、示例与生产清单由 WR BLOG 编写。

评论
0 条讨论