多 Agent 系统最脆弱的一环不是 LLM——是消息。主 Agent 发一条"请帮我调研中国市场"给子 Agent,子 Agent 没收到,或者收到后崩溃了,或者处理完了但主 Agent 不知道。这些失败模式在没有持久化信箱的系统里是不可诊断的——你能看到的只是"子 Agent 一直没回复"。

本文介绍我们的消息总线从"Redis Pub/Sub 至多一次"到"MySQL 持久化信箱至少一次"的升级过程。


1. 起点:Redis 收件箱的契约与边界

升级前的消息总线设计:

send:  LPUSH mailbox:{sid}:{agent} + PUBLISH agent:notify:{sid}
read:   Lua LRANGE + DEL(原子消费)

这其实是最优的 at-most-once 方案——Lua 原子读取+删除,Redis 内无重复问题。它在系统初始阶段足够好:

  • 消息量小(几十条/会话),丢失概率低
  • Redis 队列操作微秒级,延迟可忽略
  • Lua 原子消费自然去重

但在多 Agent 协作场景下暴露了两个硬伤:

硬伤一:消费即删除 = 崩溃即丢失。 子 Agent 从 Redis 取出消息、开始处理、崩了——消息被 DEL 了,永不再来。主 Agent 只能等超时,然后收到一个 cold 的 "子 Agent 未响应",无法重试。

硬伤二:Pub/Sub 没有送达保证。 notify 是唤醒信号,消费者的 Redis 连接断开期间(网络闪断、Redis 主从切换),唤醒信号丢失。如果消费者没在主动轮询期(子 Agent idle loop 每 1s 轮询一次),消息就会积在 List 里直到 TTL 过期,然后永久丢失。


2. 设计:持久化信箱 + 投递状态机

2.1 分层架构

写入路径:  SendMessageExecutor → MessageBus.send()
              ├── 1) INSERT agent_message_inbox (status=pending)  ← 唯一事实源
              └── 2) PUBLISH agent:notify:{sid}                   ← 纯唤醒,可丢

消费路径:  SubAgentRunner / MainAgent → MessageBus.readInbox()
              ├── 1) UPDATE status='delivered' (单事务)            ← 标记已投
              └── 2) SELECT delivered 消息

确认路径:  处理成功后 → MessageBus.ack()
              └── UPDATE status='acked'                           ← 终端态

MySQL 是唯一事实源;Redis 降级为纯唤醒信号——丢了不影响可靠性,补发机制自愈。

2.2 消息状态机

pending ──poll──▶ delivered ──ack──▶ acked
   ▲                 │
   │                 │ delivered_at 超时 (60s) + deliver_count < 5
   └─── redeliver ◀──┘
                     │ deliver_count ≥ 5
                     ▼
                   dead(+ 告警日志)
  • pending:已落库,等待投递
  • delivered:已投递但消费者未确认——如果消费者崩溃,消息停在此状态,超时后重投
  • acked:终端态,消费者确认已处理
  • dead:毒消息隔离(投递 5 次均失败),人工介入

2.3 幂等:同一条消息多次投递无害

投递和确认之间如果消费者崩溃,重投会发同一条消息两次。幂等靠两样:

  1. msgId 去重:每个 Agent 进程内维护已处理 msgId 集合,同一次运行内重投的直接跳过
  2. msgId 唯一键send 方如果收到 DuplicateKeyException,说明同消息已入信箱(发送重试),直接返回原 msgId

3. 实现细节

3.1 Transactional Poll

readInbox 是一个事务:

-- Step 1: 标记投递(同时取 pending 和超时 delivered)
UPDATE agent_message_inbox
   SET status='delivered', delivered_at=NOW(3), deliver_count=deliver_count+1
 WHERE to_agent_id=? AND session_id=?
   AND (status='pending' OR (status='delivered' AND delivered_at < NOW() - INTERVAL 60 SECOND))
   AND deliver_count < 5
 ORDER BY id LIMIT 50;

-- Step 2: 取回已标记的消息
SELECT * FROM agent_message_inbox
 WHERE to_agent_id=? AND session_id=? AND status='delivered' AND id > ?;

先 UPDATE 后 SELECT 在同一个事务中,保证"标记且取回"的原子性。即使多个消费者并发 poll 同一信箱,MySQL 的行锁也保证每条消息只会被一个消费者标记。

3.2 定时重投清扫器

独立的后台任务(每 15s 执行一次):

  1. 重投status='delivered' AND delivered_at < now-60s AND deliver_count<5 的消息自然会在下次 poll 中被重新标记和取回
  2. 补发 notifystatus='pending' 超过 30s 的消息,补发一次 Pub/Sub 唤醒(每个消息最多 3 次),复活因 Redis 问题丢失的唤醒信号
  3. 死信标记deliver_count>=5 的消息 → status=dead + WorkLogEvent 告警

3.3 ack 时机

ack 不在"收到消息时"调用——而在"消息被成功处理后"。具体:

消费方 ack 时机
子 Agent 本轮 LLM 推理完成 + result 发送成功后
主 Agent(回合开始) inbox 注入 memory 后
主 Agent(尾随循环) SSE 事件发送成功后

这个设计保证了:任何时候消费者崩溃,未 ack 的消息都会在 60s 后重投;而已 ack 的消息不会重复消费。


4. 为什么不用现成的消息队列

方案 为什么不选
Kafka 引入新组件——当前零 MQ 运维负担;Agent 消息量(会话粒度、几十到几百条)不需要分区日志架构
RabbitMQ ACK/持久化等语义成熟,但同样引入新组件;我们的写入/读写模式全是会话内短链路,MQ 的 exchange/binding 模型过重
继续用 Redis Redis List 的消费式删除无法实现 at-least-once;Redis Stream 支持 ACK 但需要 Redis 5.0+ 且消费者组的管理复杂度不低于 MySQL
MySQL + 定时器 选了。 SQL 事务语义天然支撑 poll + ack;sweeper 用 Spring @Scheduled 15s 执行一次,代码量 < 100 行;现有基础设施零新增运维

5. 效果

指标 升级前 升级后
消息丢失 Pub/Sub 丢失 + 崩溃丢失 0(重投自愈)
消费语义 at-most-once at-least-once
子 Agent 崩溃恢复 消息永丢 60s 内重投
唤醒信号丢失 消息积压至 TTL 过期 30s 内补发 notify
运维复杂度 纯 Redis Redis + MySQL(MySQL 已有)

6. 经验总结

  1. at-most-once 不是 bug,是设计选择。 初期用 Redis List + Lua 是最简实现,性能最优。升级是在消息可靠性成为瓶颈后才做的,不是过早工程。

  2. 事务是消息投递的最小原子单元。 poll 中的 UPDATE+SELECT 必须同事务;如果分两步做,两个消费者可能取到同一条消息。

  3. ack 在"处理完成后"不在"收到时"。 这是 at-least-once 的精髓——收到消息不等同于消费成功。崩溃的任意时间点上,未 ack 的消息都会重投。

  4. 定时器是零运维的消息重试引擎。 15s 执行的 sweeper 不到 100 行代码,覆盖了重投、补发唤醒、死信三个关键职责。不需要额外的 scheduler 组件。

  5. 降级 Redis 是正确的决策。 Redis Pub/Sub 从"消息通道"降为"唤醒信号",丢不丢不影响可靠性——只有一种存储持有事实状态,其余都是加速缓存。