跳到主要内容

消息确认、重试与死信

消息可靠性需要分别确认生产端写入、broker 保存、消费者处理和消费进度提交。任何两个阶段之间发生故障,都可能产生消息丢失、重复或长时间积压。

1. 生产者确认消息是否进入 broker

RabbitMQ publisher confirm 用于通知生产者消息是否由 broker 接受。它与消费者 acknowledgement 是两套独立协议:publisher confirm 成功,只说明发布阶段完成,不代表任何消费者已经处理。

Kafka producer 的 acks 决定 broker 在什么条件下确认写入。acks=all 表示当前同步副本集合满足确认条件,还应结合 min.insync.replicas 和副本数量配置。它能降低副本故障时的数据丢失风险,但网络超时后生产者仍可能无法确定写入是否已经成功。

生产者因此需要稳定消息 ID、有限重试和发送结果监控。业务数据库提交后再发送消息还会形成双写窗口,可通过 Transactional Outbox 把业务变化和待发送事件写入同一本地事务。

2. 消费者在副作用完成后确认

消费者若先 ack 或提交 offset,再写数据库,进程在两者之间崩溃会丢失业务处理。通常应先完成可持久化副作用,再确认消息进度。

反过来,业务已提交但确认尚未成功时崩溃,broker 会再次投递或消费者会从旧 offset 重读。因此 at-least-once 消费必然要求幂等。

对于批量 offset 提交,还要明确批内失败时从哪里恢复。不能因为最后一条成功,就越过前面尚未完成的消息提交进度。

3. 短暂故障使用有限即时重试

数据库瞬时连接失败、下游短暂超时可以在消费者内进行少量重试。重试应使用退避和随机抖动,并服从单条消息的总处理时限。

立即无限重试会占住消费线程,持续攻击尚未恢复的下游。错误如果来自参数、schema 或业务约束,再多次执行也不会成功,应尽快转入可观察的失败路径。

4. 延迟重试避免阻塞主消费队列

需要等待数十秒或数分钟后再试的消息,可以进入专用 retry topic 或 retry queue,并记录尝试次数、下次执行时间和原始消息标识。

重试层级应有限,例如 10 秒、1 分钟、10 分钟后各尝试一次。达到上限后转入死信;不能让消息在主队列和重试队列之间无限循环。

Kafka 常用多个 retry topic 表达延迟层级。RabbitMQ 可以结合 TTL、dead-letter exchange 或延迟消息能力实现,但要检查实际投递保证和队列积压行为。

5. 死信保存无法自动完成的消息

消息达到最大重试次数、被拒绝不再重入队、过期或超过队列限制时,可以路由到 dead-letter queue。死信中应保留:

  • 原 topic、queue 和 routing key。
  • 消息 ID、业务键和 schema 版本。
  • 首次与最后失败时间。
  • 尝试次数和最近错误类别。
  • 追踪标识及必要的处理上下文。

死信不是永久堆放区。团队需要告警、查看、修复、重放和归档流程,并限制敏感数据和保存时间。

6. 死信转发本身也可能失败

RabbitMQ 的 dead-lettering 是一次新的发布过程。目标 exchange 配置错误、目标 queue 不存在或集群故障时,死信可能无法到达预期位置。不同队列类型和配置能提供的转发保证也不同。

因此要监控 dead-letter 数量、转发失败和目标队列增量,并做真实故障演练。仅声明一个 DLX 不等于已经获得无损失败存储。

7. 重放要避免制造第二次事故

修复消费者后,死信或历史消息可能需要重放。重放前应确认:

  1. 消费逻辑已经具备幂等能力。
  2. 下游容量能够承受重放与实时流量之和。
  3. 顺序敏感消息按业务键处理。
  4. 可以暂停、限速并记录重放进度。
  5. 已有成功结果不会被旧消息覆盖。

重放工具要使用新的执行批次标识,同时保留原始消息 ID,便于去重和审计。

8. 可靠性要通过端到端指标验证

至少观察生产确认失败、发送重试、broker 副本状态、消费 lag、未确认消息、处理失败率、各重试层积压、死信数量和最老消息年龄。

单看队列长度无法区分流量增长、消费者变慢、分区倾斜和毒消息。应把消息 ID 与业务结果关联,定期核对已产生事件和已完成副作用。

9. 常见问题

9.1 RabbitMQ 的 publisher confirm 成功是否表示消息不会丢

它只确认 broker 接受了发布,具体持久性还取决于 queue、消息和集群配置。之后的路由、消费、ack 和死信是其他阶段,需要分别监控和处理。

9.2 消费失败后直接 requeue 有什么问题

同一条无法处理的消息可能立即再次投递,形成高频循环并占满消费者。应区分可重试与永久错误,限制即时次数,使用延迟重试,并在上限后进入死信。

10. 面试题

10.1 如何设计消息确认、重试和死信流程

出现公司:Shopee、京东

考察重点

  • 生产确认与消费确认的独立边界。
  • 副作用提交和 ack、offset 提交的顺序。
  • 延迟重试、死信、重放和端到端核对。

相关内容:第 1 节“生产者确认消息是否进入 broker”至第 8 节“可靠性要通过端到端指标验证”。

参考回答

生产端先确认消息是否进入 broker,RabbitMQ 的 publisher confirm 与 consumer ack 相互独立;Kafka 的 producer acks 还要结合副本和 min.insync.replicas。消费者通常先完成数据库等持久化副作用,再 ack 或提交 offset,这会在确认失败时产生重复,所以消费逻辑必须幂等。

短暂故障可以在总时限内少量退避重试,较长等待进入有限层级的 retry topic 或 queue,永久错误和达到上限的消息进入死信。死信需要告警、诊断、限速重放和转发失败监控。最后还要用业务消息 ID 核对事件与处理结果,不能只看 broker 的发送成功率。