跳到内容

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 发送后如果超时:

text
broker 未收到 -> 重试是必要的
broker 已持久化但确认丢失 -> 重试会产生重复

Kafka idempotent producer 用 producer identity、partition sequence 等机制在会话范围去除特定重试重复;RabbitMQ publisher confirm 后断线也可能让 producer 不确定确认是否已经到达。业务事件仍应有稳定 event_id 或幂等键。

Consumer Ack 必须晚于业务落地

text
receive
validate schema
perform idempotent business transaction
record event_id / resulting state
ack or commit offset

先 ack 再写数据库,进程在两者之间崩溃会永久丢失业务处理;先写数据库再 ack,崩溃会导致重投,但幂等表可以识别重复。

sql
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 把业务状态与待发布事件写进同一本地事务:

sql
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,并携带聚合版本:

json
{
  "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 和代表性旧消息。

参考

下一章不再比较产品宣传语,而是建立一套能记录事实、约束、成本和退出方案的选型流程。

Built with VitePress | Software Systems Atlas