Skip to content
Kafka 可靠性与故障治理 · 第 2 篇 / 共 4 篇
领域数据与中间件
专题Kafka 专题
当前序列Kafka 可靠性与故障治理
阅读位置第 2 篇 / 共 4 篇当前专题第 2 个序列 / 共 6 个序列

Kafka 重试与死信怎么设计:失败消息该立即重试、延迟重试还是转死信

Kafka 真正在生产里跑起来后,失败消息处理很快就会变成一个绕不过去的问题。

很多团队最开始的实现都很简单:

  • 消费失败就 catch
  • 然后原地重试几次
  • 还不行就继续抛异常

这种做法在压测时可能还能跑,但一到线上就很容易出现:

  • 某条坏消息把一个分区长期卡死
  • 下游数据库或接口已经异常,消费者还在高频重试持续放大流量
  • 偏移量迟迟不推进,消费堆积越来越高
  • 失败消息没有统一落点,排查时只能翻日志
  • 业务方问“这笔数据到底丢了没”,团队一时答不上来

先说结论

Kafka 失败处理不要只问“要不要重试”,而要先分清失败属于哪一类,再决定放在哪条链路里处理。

最常见也最实用的处理分层是:

  • 主链路有限立即重试
  • 独立链路延迟重试
  • 不再适合继续消费的死信隔离
  • 失败后的人工修复与补偿回放

真正能长期跑稳的设计,通常都有这几个共性:

  • 主消费链路要尽量短,不能被少量异常消息拖死
  • 立即重试只处理“几秒内可能恢复”的临时故障
  • 延迟重试要脱离主 Topic,避免反复占用同一消费线程
  • 死信不只是存起来,还要能定位、修复、回放和审计
  • 所有消费逻辑默认接受“重复消费可能发生”,幂等能力必须提前补好

一、先看清线上到底在失败什么

失败消息看起来都叫“消费失败”,但处理策略完全不一样。更适合先分成下面四类。

1. 短暂性故障

例如:

  • 下游接口偶发超时
  • Redis 短暂抖动
  • 数据库连接池瞬时打满
  • 网络闪断

这类问题往往有较高概率在短时间内恢复,适合做有限次数的立即重试。

2. 资源性故障

例如:

  • 下游数据库整体负载已经很高
  • 一个外部服务熔断或降级
  • 限流阈值被持续打满

这时候继续立即重试,通常只会把故障放大,应该尽快退到延迟重试链路。

3. 数据性故障

例如:

  • JSON 结构不合法
  • 字段缺失或类型错误
  • 业务状态不满足前置条件
  • 幂等键冲突,且无法自动修复

这类错误大概率不是“再试一次就好了”,继续重试只是在浪费吞吐,应该尽快转死信。

4. 程序性故障

例如:

  • 某个新版本消费逻辑有 bug
  • 配置中心下发了错误配置
  • 依赖升级后序列化逻辑变更

这类问题是否先重试,取决于修复时效。如果预计很快热修,可以先少量延迟重试;如果故障会持续一段时间,就应该尽早隔离并等修复后统一回放。

二、为什么不能把所有失败都放在主链路无限重试

因为 Kafka 消费线程和分区绑定得非常紧。

如果某条消息一直失败,却仍占着主分区消费位,就会出现几个直接后果:

  • 后面的正常消息一起被堵住
  • lag 持续升高,看起来像“整个系统消费不过来”
  • 消费者频繁重试导致 CPU、连接池、外部接口一起抖动
  • 日志里全是同一条报错,真正的新问题反而被淹没

所以失败治理的核心不是“绝不放过一条消息”,而是:

  • 在可靠性和主链路吞吐之间做分层处理
  • 让能恢复的错误有机会恢复
  • 让恢复不了的错误尽快离开主战场

三、立即重试适合放在哪一层

立即重试适合解决“短时间内高度可恢复”的错误,但它一定要很克制。

更稳妥的经验是:

  • 次数要少,一般 1 到 3 次足够
  • 间隔要短,更多像快速探测,不像真正兜底
  • 只针对明确的临时性错误码或异常类型
  • 立即重试期间不要阻塞过久,否则会拖慢 poll 节奏并引发更大的消费抖动

如果业务一次处理要调用多个外部系统,立即重试还要注意一件事:

  • 前半段成功、后半段失败时,重复执行是否安全

这就是为什么消息消费侧一定要把幂等设计放在前面,而不是放在故障后再补。

四、延迟重试真正解决的是什么问题

很多失败不是不能恢复,而是“现在恢复不了”。

典型场景包括:

  • 下游限流 5 分钟
  • 数据同步链路还没追平
  • 某个依赖服务正在发布
  • 业务前置状态稍后才会补齐

如果这种场景仍在主 Topic 内原地重试,本质上只是反复占用消费资源。更合理的做法是:

  • 记录失败原因和重试次数
  • 将消息投递到单独的重试 Topic
  • 按重试层级或时间窗口由独立消费者再处理

这样做的价值很直接:

  • 主消费链路继续推进
  • 不同延迟级别的失败消息可以分开控速
  • 可以对重试链路单独监控,不和主链路混在一起

五、死信队列不是垃圾桶,而是异常消息的治理入口

很多系统虽然建了死信 Topic,但后面没人看、没人清、没人回放。这样其实只是把问题“转移出去”,并没有真正解决。

死信更合理的定位应该是:

  • 失败消息的隔离区
  • 数据修复和人工介入的入口
  • 故障归因和规则优化的观察样本

哪些消息更适合直接进死信:

  • 反序列化失败
  • 关键字段缺失
  • 业务校验永远不通过
  • 超过最大重试次数仍失败
  • 当前版本逻辑明确不支持该类消息

一旦进了死信,至少要保留这些信息:

  • 原始消息体
  • 业务主键或幂等键
  • 来源 Topic、分区、offset
  • 首次失败时间与最近失败时间
  • 失败原因摘要和异常堆栈
  • 已重试次数

没有这些字段,后面就算想回放,也很难判断“该不该回、怎么回、回到哪里”。

六、一个更能落地的失败处理顺序

更推荐把失败处理顺序设计成下面这样:

1. 先判断是不是可恢复错误

如果是超时、连接失败、服务临时不可用,可以先进入立即重试。

2. 立即重试失败后,再判断是否值得等待

如果问题预计稍后可恢复,比如下游限流、依赖链路延迟,可以转入延迟重试。

3. 明显不可恢复或超过窗口后,转死信

不要让这类消息继续占住主链路资源。

4. 死信后必须有补偿路径

例如:

  • 修复数据后重新投递
  • 通过后台工具进行单条或批量回放
  • 按错误类型聚合处理

七、设计重试链路时最容易忽略的几个边界

1. 重试次数不等于可靠性更高

如果错误本身不可恢复,重试越多只是放大损耗。

2. 延迟重试也需要限速

很多团队把失败消息转到重试 Topic 后就放心了,但一旦重试 Topic 在某个时间点集中回放,还是可能再次把下游打爆。

3. 死信回放必须防重复

无论是人工回放还是程序批量补偿,都要把幂等键、业务状态检查、回放去重考虑进去。

4. 不要让“消费成功”和“业务成功”混为一谈

消息从 Kafka 里取出来只是开始。真正决定是否提交 offset 的,是业务成功边界,而不是代码有没有执行到最后一行。

八、线上排查时更实用的观察顺序

如果你发现某个 Topic 消费持续异常,可以按下面顺序排查:

  1. 看 lag 是全局上涨还是集中在少数分区。
  2. 看错误是偶发超时还是某类固定数据持续失败。
  3. 看是否存在某条坏消息长期卡住同一分区。
  4. 看下游资源是否已经过载,重试是否在持续放大压力。
  5. 看死信和重试 Topic 是否有堆积,以及失败类型是否集中。
  6. 最后再决定是扩容消费者、调重试参数,还是先做消息隔离。

很多时候,真正的问题不是消费者不够,而是错误消息没有分层治理。

九、实战里更稳的一套设计思路

如果业务链路比较重要,通常可以这样落:

  1. 主 Topic 消费只保留很轻的立即重试。
  2. 单独建设延迟重试 Topic 或多级重试 Topic。
  3. 为死信消息建立后台查询、筛选和回放能力。
  4. 所有消费逻辑以业务主键做幂等。
  5. 监控拆分为主消费成功率、重试成功率、死信量和回放成功率四层。

这样一来,失败消息不再是“吞掉还是报错”二选一,而是有完整治理路径。

一句话总结

Kafka 失败治理真正要解决的,不是“失败后要不要再试一次”,而是如何把临时错误、资源性错误、数据错误和程序错误拆开处理。

主链路只负责快速消费和轻量兜底,延迟重试负责给系统恢复时间,死信负责隔离不可恢复消息,补偿回放负责把异常闭环补齐。把这四层理顺后,Kafka 才能真正既稳又可维护。

延伸阅读相关文章优先当前专题,再补跨专题关联。
同一序列 · 顺着当前主线继续读Kafka 重复消费治理案例适合把顺序性、幂等、重试、死信和消息积压放在一条线上连续看。Kafka 专题 · Kafka 可靠性与故障治理同一序列 · 顺着当前主线继续读Kafka 消息积压排查适合把顺序性、幂等、重试、死信和消息积压放在一条线上连续看。Kafka 专题 · Kafka 可靠性与故障治理同专题其他序列 · 共享标签:RabbitMQ重试主题和死信主题怎么设计层级适合把 Rebalance、offset 提交、重试层级和 Broker 磁盘状态放在同一条 Kafka 消费治理主线上看。Kafka 专题 · Kafka 消费治理与运维观察同专题其他序列 · 共享标签:RabbitMQRabbitMQ、Kafka、RocketMQ 怎么选适合先理解为什么要用 MQ,再衔接 Kafka 的核心架构、消费组和生产端基础。Kafka 专题 · 消息系统选型与 Kafka 基础跨专题关联 · 同场景:基础学习BlockingQueue 和生产者消费者模型适合沿着 JUC、volatile、CAS、AQS 和锁实现这条主线往下看。Java 专题 · Java 并发与锁机制跨专题关联 · 同场景:基础学习RocketMQ 的 Topic、Queue 和消费模型怎么理解适合把 Topic、Queue、顺序消息和事务消息放在一起理解。RocketMQ 专题 · RocketMQ 核心模型与业务消息
继续阅读Kafka 可靠性与故障治理当前序列第 2 篇 / 共 4 篇当前专题第 2 个序列 / 共 6 个序列
往前看
上一篇Kafka 顺序性、幂等生产者和精确一次回到当前序列上一章上一序列消息系统选型与 Kafka 基础从第 1 篇开始:MQ 为什么要用,以及怎么保证消息可靠
往后看
下一篇Kafka 重复消费治理案例继续当前序列下一章下一序列Kafka 高级投递与消费治理从第 1 篇开始:Kafka 分区键怎么设计更稳

把零散经验整理成可查、可复用、可持续更新的企业级知识门户