Appearance
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 消费持续异常,可以按下面顺序排查:
- 看 lag 是全局上涨还是集中在少数分区。
- 看错误是偶发超时还是某类固定数据持续失败。
- 看是否存在某条坏消息长期卡住同一分区。
- 看下游资源是否已经过载,重试是否在持续放大压力。
- 看死信和重试 Topic 是否有堆积,以及失败类型是否集中。
- 最后再决定是扩容消费者、调重试参数,还是先做消息隔离。
很多时候,真正的问题不是消费者不够,而是错误消息没有分层治理。
九、实战里更稳的一套设计思路
如果业务链路比较重要,通常可以这样落:
- 主 Topic 消费只保留很轻的立即重试。
- 单独建设延迟重试 Topic 或多级重试 Topic。
- 为死信消息建立后台查询、筛选和回放能力。
- 所有消费逻辑以业务主键做幂等。
- 监控拆分为主消费成功率、重试成功率、死信量和回放成功率四层。
这样一来,失败消息不再是“吞掉还是报错”二选一,而是有完整治理路径。
一句话总结
Kafka 失败治理真正要解决的,不是“失败后要不要再试一次”,而是如何把临时错误、资源性错误、数据错误和程序错误拆开处理。
主链路只负责快速消费和轻量兜底,延迟重试负责给系统恢复时间,死信负责隔离不可恢复消息,补偿回放负责把异常闭环补齐。把这四层理顺后,Kafka 才能真正既稳又可维护。