15.2 Outbox、幂等消费与可恢复消息
一次进程崩溃恰好发生在“数据库已提交、消息未发出”之间,议事厅收到了一笔永远不会通知下游的订单。
服务常需要同时更新数据库并发布事件。先提交数据库再发消息,进程可能在两步之间崩溃;先发消息再提交数据库,消费者可能看到最终并不存在的事实。
Transactional Outbox 把业务变更和待发布消息写入同一个本地事务,再由独立 relay 发布消息。
Outbox 关闭原子缺口
sql
BEGIN;
UPDATE registration
SET status = 'CONFIRMED'
WHERE id = :registration_id
AND status = 'PENDING_PAYMENT';
INSERT INTO outbox (
event_id, aggregate_id, aggregate_version,
event_type, payload, occurred_at, publish_state
) VALUES (
:event_id, :registration_id, :version,
'RegistrationConfirmed', :payload, CURRENT_TIMESTAMP, 'PENDING'
);
COMMIT;relay 可以轮询 Outbox,也可以通过数据库日志 CDC 捕获。无论哪种方式,都可能在“消息已发送、发布状态尚未记录”时崩溃,因此重复发布仍然可能发生。
Outbox 提供的是“业务变更与发布意图原子记录”,不是自动的 exactly-once 端到端处理。
消费者必须幂等
消费者可以用事件 ID 建立 Inbox/去重表:
sql
BEGIN;
INSERT INTO consumed_message (consumer, event_id, consumed_at)
VALUES (:consumer, :event_id, CURRENT_TIMESTAMP)
ON CONFLICT DO NOTHING;
-- 只有插入成功时才执行业务更新
UPDATE reward_account
SET points = points + :delta
WHERE player_id = :player_id;
COMMIT;实际实现必须检查插入影响行数,并让“记录已消费”和业务更新处于同一事务。若只在内存 Set 中去重,重启后记录会丢失;若先标记已消费再更新业务,崩溃会永久漏处理。
有些操作天然幂等,例如“把状态设置为版本 7 的值”;有些操作如“积分加 10”必须依赖事件 ID 或业务操作 ID 去重。
顺序应按业务键定义
消息系统通常只能在一个分区内保证顺序。把 aggregate_id 作为分区键,可以让同一报名或比赛的事件进入同一有序流;不同聚合之间不应假设全局顺序。
消费者还应检查聚合版本:
text
当前版本 5,收到版本 6 → 应用
当前版本 6,收到版本 6 → 重复,忽略
当前版本 5,收到版本 8 → 缺少 6、7,暂缓或修复只依赖到达时间无法可靠处理重放和跨分区延迟。
重试与死信队列
先区分错误:
- 瞬态错误:数据库暂时不可用,可以退避重试;
- 永久错误:schema 不兼容、必填字段缺失,重复不会自愈;
- 业务拒绝:不应伪装成基础设施重试;
- 未知错误:需要保留上下文并限制重试次数。
死信队列不是垃圾桶。进入死信的消息需要:
- 原始 payload、headers、事件 ID 和失败原因;
- 告警和明确所有者;
- 修复、重放或丢弃的审计操作;
- 重放时继续遵守幂等和顺序规则。
无限重试会阻塞分区并制造成本,直接跳过又可能破坏业务不变量。
运营指标
至少监控:
- Outbox 最老待发布消息的年龄;
- 发布和消费延迟;
- 重复消息、幂等命中和版本缺口;
- 重试次数、死信数量与最老死信年龄;
- 业务对账差异,而不只是 broker 是否在线。
下一课区分 CQRS 与 Event Sourcing,说明它们何时能解决读写模型和历史重建问题,何时只是增加复杂度。
参考资料
- Chris Richardson, Transactional Outbox
- Debezium, Outbox Event Router
- Apache Kafka, Design: Delivery Semantics