跳过导航

消息消费失败后如何设计重试与死信机制?

约 6 分钟...次浏览
专栏分布式系统与消息队列第 3 篇

消费失败后立即重试,看起来最直接,却可能把一次短暂故障放大成重试风暴。数据库已经过载时,成千上万条消息以毫秒级频率再次访问,只会让恢复更慢。

生产级重试系统的目标不是“尽量多试几次”,而是区分错误、控制速率、保持幂等,并为最终失败提供可观测的人工闭环。

一、先给失败分类

可重试错误

网络超时、连接池暂时耗尽、下游 503、限流等通常具有瞬时性,可以在延迟后重试。

不可重试错误

参数缺失、Schema 不兼容、业务规则拒绝、目标账号不存在等,重复执行不会自行恢复,应快速进入死信或人工处理。

未知错误

代码 Bug、空指针等不能简单永久重试。应有限重试以覆盖偶发环境因素,随后告警并进入死信。

RetryDecision classify(Throwable e) {
    if (e instanceof TimeoutException) return RetryDecision.RETRY;
    if (e instanceof ValidationException) return RetryDecision.DEAD_LETTER;
    return RetryDecision.RETRY_LIMITED;
}

二、不要在消费线程里无限循环

while (true) {
    try {
        handle(message);
        break;
    } catch (Exception ignored) {}
}

这种写法会占满消费线程,阻塞同一分区后续消息,持续攻击故障下游,并失去统一监控。重试必须有次数、间隔、总时间预算和退出路径。

三、指数退避与随机抖动

常见退避公式:

delay = min(maxDelay, baseDelay * 2^attempt) + randomJitter

例如 1 秒、2 秒、4 秒、8 秒,最多 5 分钟。随机抖动能避免大量消息在同一时刻再次冲击下游。

Spring Kafka 的错误处理器可以配置退避:

var recoverer = new DeadLetterPublishingRecoverer(kafkaTemplate);
var handler = new DefaultErrorHandler(
    recoverer,
    new ExponentialBackOffWithMaxRetries(4));
handler.addNotRetryableExceptions(ValidationException.class);

这个示例提供指数退避,但不自动代表已经加入随机抖动;若大量实例会同时失败,应按当前 Spring Kafka 版本提供的扩展点实现或配置 jitter,并通过时间分布测试确认它确实生效。具体 API 随版本变化,升级时应以当前版本文档为准并做集成测试。

四、阻塞重试与非阻塞重试

阻塞重试在当前消费调用中等待后再次执行,简单但会占用线程并阻塞分区,适合次数少、间隔短的瞬时失败。

非阻塞重试把失败消息投递到延迟级别不同的重试 Topic:

orders
 -> orders-retry-10s
 -> orders-retry-1m
 -> orders-retry-10m
 -> orders-dlt

它释放主消费线程,适合长间隔重试,但会增加 Topic、路由、监控与消息顺序管理复杂度。

五、死信队列不是垃圾桶

死信消息至少应保留:

  • 原 Topic、Partition、Offset。
  • 原始 key、payload 和 headers。
  • eventId、业务主键与消费者名称。
  • 异常类、错误摘要、首次和最后失败时间。
  • 已重试次数与应用版本。

敏感数据需要脱敏,异常堆栈可存到日志平台并通过 traceId 关联,避免消息头过大。

死信还需要:积压告警、查询页面、责任人、修复后重放工具、重放审批和审计记录。没有处理流程的 DLT 只是把故障藏得更深。

六、重试必须以幂等为前提

一次处理可能完成了数据库写入,却在返回成功前超时。重试时业务不能再次扣款:

INSERT INTO consumed_event(event_id, consumer)
VALUES (?, ?);
-- UNIQUE(event_id, consumer)

在同一本地事务中插入去重记录并完成业务更新。不同数据库对唯一键冲突的事务语义不同:例如 PostgreSQL 普通唯一冲突会使当前事务进入失败状态,不能简单 catch 后继续。应使用数据库支持的 ON CONFLICT DO NOTHING、等价原子语句或在事务边界外识别冲突,并且只有确认已有记录是已提交成功结果时才跳过。

外部 HTTP 调用也应携带稳定幂等键,让下游识别重复请求。仅在消费者本地去重,无法覆盖“下游成功但本地超时”的窗口。

七、重试与顺序性的冲突

订单的 CREATED 失败后被移到重试 Topic,原分区的 PAID 可能继续消费,导致状态乱序。

可选策略:

  • 严格顺序业务阻塞该 key 或整个分区,等待前序成功。
  • 按聚合 ID 做状态版本校验,晚到事件拒绝或暂存。
  • 将同一 key 的后续消息一起转移到有序重试通道。
  • 使用可重放的状态机,让事件乱序也不会非法迁移。

不能同时承诺“失败不阻塞任何后续消息”和“严格保持全顺序”,必须做取舍。

八、重放工具的安全设计

重放 DLT 时不要直接把所有消息一键灌回主 Topic。应支持按 eventId、时间范围、错误类型筛选,限制速率,预览影响,并保留操作者与批次号。

修复代码上线前重放只会再次失败;下游尚未恢复时批量重放则会制造第二次事故。

九、监控哪些指标

  • 主 Topic 与各重试 Topic 的积压量和最老消息年龄。
  • 按异常类型统计的失败率。
  • 首次成功率、各重试次数成功率。
  • DLT 写入和待处理数量。
  • 单条消息端到端处理时长。
  • 下游限流、超时和熔断状态。

最老消息年龄通常比单纯消息条数更能反映业务影响。

十、设计检查清单

  • 错误是否区分可重试与不可重试?
  • 是否有最大次数、最大间隔和总重试预算?
  • 是否使用指数退避与随机抖动?
  • 消费和下游调用是否都支持幂等?
  • 非阻塞重试是否破坏业务顺序?
  • DLT 是否保留足够上下文且做好脱敏?
  • 原始 payload 无法反序列化时,是否能以受限字节形式进入隔离区而不再次触发同一解析器?
  • 是否有积压告警、责任人和受控重放工具?
  • 重放前是否确认代码和下游故障已经修复?

总结

重试是一种故障恢复策略,也是一种额外流量源。可靠设计必须先分类错误,再用有限次数、指数退避和抖动控制节奏,以幂等保证重复执行安全,最终失败进入可运营的死信闭环。最危险的方案不是不重试,而是无边界、不可观察地持续重试。

分享:
文章作者:狼码纪
版权声明:本博客所有文章除特别声明外,均采用 CC BY-NC-SA 4.0 许可协议。文章可能参考了其他优秀文章,如有侵权请联系删除。