消息消费失败后如何设计重试与死信机制?
消费失败后立即重试,看起来最直接,却可能把一次短暂故障放大成重试风暴。数据库已经过载时,成千上万条消息以毫秒级频率再次访问,只会让恢复更慢。
生产级重试系统的目标不是“尽量多试几次”,而是区分错误、控制速率、保持幂等,并为最终失败提供可观测的人工闭环。
一、先给失败分类
可重试错误
网络超时、连接池暂时耗尽、下游 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 无法反序列化时,是否能以受限字节形式进入隔离区而不再次触发同一解析器?
- 是否有积压告警、责任人和受控重放工具?
- 重放前是否确认代码和下游故障已经修复?
总结
重试是一种故障恢复策略,也是一种额外流量源。可靠设计必须先分类错误,再用有限次数、指数退避和抖动控制节奏,以幂等保证重复执行安全,最终失败进入可运营的死信闭环。最危险的方案不是不重试,而是无边界、不可观察地持续重试。