多 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 幂等:同一条消息多次投递无害
投递和确认之间如果消费者崩溃,重投会发同一条消息两次。幂等靠两样:
- msgId 去重:每个 Agent 进程内维护已处理
msgId集合,同一次运行内重投的直接跳过 - 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 执行一次):
- 重投:
status='delivered' AND delivered_at < now-60s AND deliver_count<5的消息自然会在下次 poll 中被重新标记和取回 - 补发 notify:
status='pending'超过 30s 的消息,补发一次 Pub/Sub 唤醒(每个消息最多 3 次),复活因 Redis 问题丢失的唤醒信号 - 死信标记:
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. 经验总结
at-most-once 不是 bug,是设计选择。 初期用 Redis List + Lua 是最简实现,性能最优。升级是在消息可靠性成为瓶颈后才做的,不是过早工程。
事务是消息投递的最小原子单元。 poll 中的 UPDATE+SELECT 必须同事务;如果分两步做,两个消费者可能取到同一条消息。
ack 在"处理完成后"不在"收到时"。 这是 at-least-once 的精髓——收到消息不等同于消费成功。崩溃的任意时间点上,未 ack 的消息都会重投。
定时器是零运维的消息重试引擎。 15s 执行的 sweeper 不到 100 行代码,覆盖了重投、补发唤醒、死信三个关键职责。不需要额外的 scheduler 组件。
降级 Redis 是正确的决策。 Redis Pub/Sub 从"消息通道"降为"唤醒信号",丢不丢不影响可靠性——只有一种存储持有事实状态,其余都是加速缓存。