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 的一次运行不是一个函数调用

设一次用户任务为:

T=(tenant_id,user_id,conversation_id,task_id)T = (tenant\_id, user\_id, conversation\_id, task\_id)

它可能经历以下步骤:

  1. 接收用户请求;
  2. 判断任务类型;
  3. 查询订单、知识库或内部系统;
  4. 让一个或多个专门 Agent 分析;
  5. 生成回复或执行外部动作;
  6. 发送结果;
  7. 将执行过程和最终状态写入持久化存储。

在同步程序中,这些步骤可能都发生在一次 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 是一组具有相同订阅语义、保留策略和路由规则的事件流。可以把它理解为:

Topic=(name,retention,partitioning,subscriptions,delivery_policy)Topic = (name, retention, partitioning, subscriptions, delivery\_policy)

其中:

  • 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 数量过多会使运维、监控和权限管理复杂化。

较合理的划分依据是:

  1. 消费者是否相同;
  2. 顺序键是否相同;
  3. 保留时间是否相同;
  4. 重试和死信策略是否相同;
  5. 数据敏感级别是否相同;
  6. 是否允许独立扩缩容。

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 至少一次投递与幂等性

在存在网络超时、消费者崩溃和租约过期时,消息系统通常更容易保证:

delivery1delivery \geq 1

即“消息至少送达一次”,而不是严格只送达一次。

如果同一个事件可能被处理多次,消费者必须满足:

handle(e);handle(e)handle(e)handle(e); handle(e) \equiv handle(e)

这里的等价不是指执行过程完全相同,而是指最终可观察状态和外部副作用不会被重复放大。

最简单的去重表:

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 无法解析、违反结构约束 受限重试或转人工
外部副作用不确定 请求超时 使用幂等键查询或重试

错误重试的核心不是“多试几次”,而是判断:

retryable(error){true,false}retryable(error) \in \{true, false\}

对于不可重试错误无限重试,只会造成队列堆积和成本放大。


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_idspan_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_idcorrelation_id,但事件 payload 不应为了方便排查而携带完整身份证号、银行卡号或原始私密对话。关联信息和业务数据应分别治理。


5. 顺序:真正需要保证的是“哪个实体上的顺序”

5.1 全局顺序通常不是需求

假设系统收到两个会话的消息:

A1: 会话 A 的“取消订单”
B1: 会话 B 的“查询物流”
A2: 会话 A 的“确认取消”

通常只要求:

A1 < A2

并不要求:

A1 < B1 < A2

如果系统强行保证所有会话的全局顺序,单个慢任务就可能阻塞所有用户。

因此应先定义顺序域:

order_key(e)=aggregate_idorder\_key(e) = aggregate\_id

例如:

  • 会话内顺序:conversation_id
  • 订单内顺序:order_id
  • 任务内顺序:task_id
  • 用户账户内顺序:account_id

5.2 分区顺序的充分条件

如果事件流按 order_key 分区,并且同一个分区内由单一逻辑序列处理,那么可以得到:

ei.order_key=ej.order_keyseq(ei)<seq(ej)apply(ei)apply(ej)e_i.order\_key = e_j.order\_key \land seq(e_i) < seq(e_j) \Rightarrow apply(e_i) \prec apply(e_j)

这里:

  • seq(e) 是事件在该顺序域中的序号;
  • apply(e_i) \prec apply(e_j) 表示 e_i 先于 e_j 应用。

但这只有在以下条件成立时才有效:

  1. 同一 order_key 始终路由到同一分区;
  2. 消费者不会让同一分区中的事件并发越过前序事件;
  3. 事件序号由可靠的生产者或存储层生成;
  4. 重试不会让旧事件在新事件之后错误覆盖状态。

仅仅给事件加一个时间戳,不足以保证顺序。时钟可能漂移,事件可能异步发布,甚至后生成的事件可能携带更早的业务时间。

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 运行状态;
  • 搜索索引;
  • 会话投影;
  • 通知状态;
  • 审计日志;
  • 统计数据。

这些状态不会在同一个数据库事务中同时更新。因此在时间 tt 上,不同读取者可能看到不同版本:

Smain(t)Sprojection(t)S_{main}(t) \neq S_{projection}(t)

最终一致性要求的是:如果后续事件能够持续送达、消费者能够最终成功处理,并且没有新的更新,那么存在某个有限时间 tt',使得:

Sprojection(t)=f(Smain)S_{projection}(t') = f(S_{main})

其中 ff 是投影规则。

这不是“最终一定正确”的无条件承诺。它依赖于:

  1. 事件不会永久丢失;
  2. 消费者会重试或进入可恢复死信;
  3. 处理逻辑具备幂等性;
  4. 事件应用顺序满足业务需要;
  5. 投影逻辑能够处理重复和迟到事件;
  6. 主状态本身没有被非法并发覆盖。

6.2 最终一致性不能替代业务不变量

例如订单取消和支付退款存在以下不变量:

refunded(order)cancelled(order)refunded(order) \Rightarrow cancelled(order)

如果允许通知投影暂时落后,这通常没问题:

主库:订单已取消
通知投影:仍显示“处理中”

但如果退款服务先执行了退款,订单主状态却仍是“已支付”,就可能违反业务不变量。

因此要区分:

  • 可暂时落后的读取模型:搜索索引、统计、通知列表;
  • 必须原子保护的业务不变量:余额、库存、支付状态、权限;
  • 必须由外部确认的事实:支付完成、短信发送、文件上传。

最终一致性适合传播状态,不适合掩盖原子性要求。

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,但没有收到 E1E2,有三种处理方式:

  1. 根据事件 payload 直接构建状态;
  2. 发现版本间隙,暂存 E3
  3. 从事件存储补读缺失事件。

若事件包含 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 事件溯源的核心

事件溯源将实体当前状态视为事件折叠结果:

Staten=fold(apply,State0,[E1,E2,,En])State_n = fold(apply, State_0, [E_1, E_2, \dots, E_n])

例如:

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 背压是什么

当事件产生速度高于消费者处理速度时,队列长度增长:

Qt+1=max(0,Qt+λtμt)Q_{t+1} = \max(0, Q_t + \lambda_t - \mu_t)

其中:

  • QtQ_t:时刻 tt 的待处理任务数;
  • λt\lambda_t:生产速率;
  • μt\mu_t:消费速率。

当长期满足:

λ>μ\lambda > \mu

队列就会持续增长,最终表现为延迟升高、内存耗尽、重试风暴或用户收到过时回复。

背压不是简单地“加更多 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 的职责是:

  1. 获取结构化上下文;
  2. 调用模型;
  3. 验证输出;
  4. 持久化草稿;
  5. 发布 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 可交换事件可以并行

若两个事件满足:

apply(Ea,apply(Eb,S))=apply(Eb,apply(Ea,S))apply(E_a, apply(E_b, S)) = apply(E_b, apply(E_a, S))

则它们对状态更新是可交换的,可以并行处理。

例如两个独立的统计计数增量:

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_idcorrelation_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

需要区分几个常被混淆的指标:

  • 队列深度:还有多少未处理消息;
  • 消费延迟:事件产生到开始处理的时间;
  • 处理耗时:开始处理到完成的时间;
  • 重试率:处理失败后再次尝试的比例;
  • 版本间隙数:投影遇到的缺失序号;
  • 过期丢弃率:任务到达时已经无效的比例;
  • 未知副作用数:外部请求超时且无法确认结果的数量。

例如队列深度下降并不一定代表系统恢复,也可能是消费者大量把消息标记为 expireddead-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、模型、协议与框架官方资料重新梳理;正文、示例与生产清单由 WR BLOG 编写。