17.2 确认边界、幂等消费与积压恢复
“消息只处理一次”听起来像 broker 的配置,实际上横跨 producer、broker、consumer 和业务数据库。只要确认消息可能丢失,发送方就必须在“可能已成功”和“可能未成功”之间做决定。
三种 delivery 语义先限定范围
- at-most-once:可能丢失,但不主动重投;
- at-least-once:不轻易丢失,但故障恢复会重复交付;
- exactly-once:在明确边界内,最终可观察效果等价于一次。
“Exactly-once”必须补全作用域。Kafka transaction 可以原子写多个 Kafka partition,并让 read_committed 消费者隐藏未提交记录;这不自动把一封邮件、一个 HTTP 调用和任意外部数据库也纳入同一事务。
Producer 确认仍会留下不确定性
Producer 发送后如果超时:
broker 未收到 -> 重试是必要的
broker 已持久化但确认丢失 -> 重试会产生重复Kafka idempotent producer 用 producer identity、partition sequence 等机制在会话范围去除特定重试重复;RabbitMQ publisher confirm 后断线也可能让 producer 不确定确认是否已经到达。业务事件仍应有稳定 event_id 或幂等键。
Consumer Ack 必须晚于业务落地
receive
validate schema
perform idempotent business transaction
record event_id / resulting state
ack or commit offset先 ack 再写数据库,进程在两者之间崩溃会永久丢失业务处理;先写数据库再 ack,崩溃会导致重投,但幂等表可以识别重复。
CREATE TABLE consumed_event (
consumer_name text NOT NULL,
event_id uuid NOT NULL,
consumed_at timestamptz NOT NULL DEFAULT now(),
PRIMARY KEY (consumer_name, event_id)
);
BEGIN;
INSERT INTO consumed_event (consumer_name, event_id)
VALUES ('inventory-projector', :event_id)
ON CONFLICT DO NOTHING;
-- 只有上一条确实插入时,才执行同一事务内的业务更新
COMMIT;实现时必须读取插入影响行数,重复事件不能再次更新业务表。幂等记录的保留期至少覆盖 broker 可能重放的最长窗口和人工恢复窗口。
Outbox 解决数据库提交与发布的双写
业务事务若先改数据库再发消息,进程可能在两步之间崩溃;若先发消息再改数据库,消费者可能看到尚未成立的事实。
Outbox 把业务状态与待发布事件写进同一本地事务:
BEGIN;
UPDATE orders
SET status = 'paid'
WHERE order_id = :order_id
AND status = 'pending';
INSERT INTO outbox_event (
event_id, aggregate_type, aggregate_id, event_type, payload
) VALUES (
:event_id, 'order', :order_id, 'OrderPaid', :payload
);
COMMIT;独立 relay 或 CDC 读取 outbox 并发布。Relay 也可能重复发布,所以消费者仍要幂等。Outbox 解决原子产生事实,不等于全链路自动 exactly-once。
顺序只在选定键与通道内成立
若同一订单事件必须保持顺序,应让它们使用同一 partition key/queue,并携带聚合版本:
{
"event_id": "...",
"aggregate_id": "order-9182",
"aggregate_version": 7,
"event_type": "OrderPaid"
}消费者可拒绝倒退版本,或把缺口暂存等待。跨 partition 没有天然全局顺序;消费者并行处理、重试队列和死信回放也可能改变完成次序。
Poison Message 与重试预算
反序列化失败、永久业务约束错误和暂时依赖故障不能用同一无限重试循环处理。
建议分类:
- transient:指数退避、抖动和有限次数;
- rate-limited:尊重服务端重试时间并限流;
- permanent/schema:送隔离队列,保留原消息、错误与版本;
- operator action:修复后可按原 key 顺序受控回放。
死信队列不是垃圾桶。必须有告警、负责人、修复流程、重放工具和再次失败策略。
Backpressure 与积压恢复
队列能把突发流量变成积压,但不能凭空增加下游容量。应同时监控:
- ingress rate 与 sustainable processing rate;
- consumer lag / backlog age,而不只消息数量;
- 每条处理延迟与失败率;
- partition/queue 间偏斜;
- broker 磁盘、复制和保留期余量;
- 从峰值积压恢复到正常所需时间。
若每秒进入 10,000 条、稳定只能处理 8,000 条,积压会无限增长。削峰成立的前提是长期平均消费能力高于长期平均生产速度,或系统能丢弃/聚合部分事件。
RabbitMQ consumer prefetch 控制未确认 delivery 数量;Kafka max.poll.interval.ms、批大小和处理线程模型共同影响 rebalance 与吞吐。参数应围绕单条处理成本和内存预算压测,不能照抄固定推荐值。
Schema 演进
消息一旦进入长保留日志,旧消费者和历史重放都会遇到它。事件应有版本与兼容策略:
- 新增可选字段通常向后兼容;
- 改字段含义比改字段名更危险;
- 删除字段前确认所有消费者和历史回放;
- 事件表达已发生事实,不复用成远程过程调用命令包;
- 在 CI 中验证 schema compatibility 和代表性旧消息。
参考
- KafkaProducer:Idempotence and Transactions
- RabbitMQ:Consumer Acknowledgements and Publisher Confirms
- Debezium:Outbox Event Router
下一章不再比较产品宣传语,而是建立一套能记录事实、约束、成本和退出方案的选型流程。