RocketMQ 对比 Kafka 消费失败的重试
先说结论:
- Kafka:原生没有消息级别的重试机制。就是那套 partition + offset 模型。所以 Kafka 原生做不到优雅的指数退避重投。
- RocketMQ:有独立的重试队列 + 延迟等级机制,重试和 offset 提交是解耦的。所以它天生就能"隔一段时间再投递"。
一、Kafka:靠 offset,“不能"优雅重试
Kafka 的消费模型非常"薄”。它的设计哲学是:broker 只管存储和顺序读取,重试逻辑是消费者自己的事。
offset 的本质
一个 partition 里的消息是一条严格有序的日志,consumer 用一个 offset 指针记录"我消费到哪了"。
partition-0: [msg0][msg1][msg2][msg3][msg4]...
↑
committed offset = 2 (表示 0、1 已确认处理完)
消费失败会怎样?
关键点:Kafka 的 offset 是"连续推进"的水位线,不是"逐条 ack"。 这跟 RocketMQ 那种逐条确认完全不同。
所以当 msg2 处理失败时,我们其实只有两个选择:
选择 1:不提交 offset(卡在 2)
下次 poll 会从 msg2 重新拉。如果 msg2 一直失败,就永远卡在这里,整个 partition 阻塞,后面的消息全堵死。
选择 2:提交 offset(推进到 3)
msg2 就等于被"跳过 / 丢弃"了,不会再拉到。
那 Kafka 生态里的"重试"怎么实现?
全靠消费者应用层自己搭,broker 不参与。常见方案有三种:
1. 本地内存重试
Spring Kafka 的 RetryTemplate、DefaultErrorHandler,在消费者进程内做几次重试。但这会阻塞后续消息,且进程一挂,重试状态就丢了。
2. 记到 DB
失败消息记录到 DB,再启动任务扫描重试。
3. 外挂一个 DelayServer
比如美团的 Mafka(Kafka 改造)。消费失败重试的实现:
- 业务消费普通 topic 消息,抛出异常;
- Mafka SDK 捕获异常,封装原消息,header 带上:原 topic、partition、原 offset、
groupId、重试次数、到期时间、异常栈; - SDK 把封装后的消息发送给 DelayServer,必须确认 DelayServer 写入成功;
- 确认托管成功,立刻 commit 原始消息 offset,原分区放行后续消息,不会阻塞;
- DelayServer 内部存储这条延迟任务,按到期时间调度;
- 到期后:DelayServer 把消息投递到【该消费组专属的内部重试队列】,不是原始业务 topic;
- 当前 group 的消费者拉取这条重试消息,
retryCount++;header 里面携带原始 topic 信息,业务感知上仿佛是原 topic 过来的消息; - 再次执行业务逻辑:
- 消费成功:正常处理结束;
- 继续失败:再次交给 DelayServer,继续延迟退避;
- 超过
maxRetry:投递用户配置的死信 topic,停止重试。
业务视角:消费者代码写的是监听原始 topic,但是底层 Mafka SDK 同时监听业务 topic + 当前 group 的内部重试队列;业务代码完全无感知,注解还是写原始 topic,SDK 做了屏蔽。
为什么不能直接写回原始 topic
Kafka 的 topic 是多消费组共享,不能因为一个消费组的消费失败,让别的消费组也做重试。
二、RocketMQ:内建重试队列,offset 与重试解耦
RocketMQ 是"消息队列"思维(不是纯日志),它原生支持逐条重试和延迟投递。
消费失败会怎样?
消费者的监听器返回一个状态:
public ConsumeConcurrentlyStatus consumeMessage(...) {
try {
// 业务处理
return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
} catch (Exception e) {
// 告诉 broker:待会儿重投
return ConsumeConcurrentlyStatus.RECONSUME_LATER;
}
}
当返回 RECONSUME_LATER 时,发生的事情是:
- 消费者照常把 offset 提交 / 推进了(所以不会阻塞后面的消息,也不会把
msg3、msg4拖下水)。 - 消费者主动把这条失败的消息发回给 broker,broker 把它放进一个特殊的重试队列:
%RETRY%{consumerGroup}。 - broker 根据这条消息已经重试了几次,赋予它一个延迟等级(delayLevel),到点后才重新投递给消费者。
这就是"阶梯 / 指数退避"的来源
RocketMQ 有 18 个固定延迟等级:
level: 1 2 3 4 5 6 7 8 9 10 ... 18
delay: 1s 5s 10s 30s 1m 2m 3m 4m 5m 6m ... 2h
消费重试时按重试次数逐级递增延迟(第 1 次重试用 level 3 = 10s,第 2 次 level 4 = 30s…… 阶梯式变长),默认最多重试 16 次,全部失败后进入死信队列 %DLQ%{consumerGroup}。
完整时序:注意 SCHEDULE 与 %RETRY% 的先后
这里有个容易搞反的点:消息不是先进 %RETRY% 再进 SCHEDULE,而是先被 SCHEDULE_TOPIC_XXXX 藏起来延迟,到期后才落到 %RETRY%。 SCHEDULE 是"等待区",%RETRY% 是"等完之后消费者来取的取件柜"。
① 消费者消费失败,返回 RECONSUME_LATER
│
▼
② 消费者发送SEND_MSG_BACK请求把这条消息发回 broker,请求目标topic填写%RETRY%{consumerGroup},同时根据reconsumeTimes计算本次delayLevel
│
▼
③ broker收到这条消息,**并不会真正写入%RETRY%物理队列**;检测delayLevel>0,内存做偷梁换柱:
- 将真实目标 REAL_TOPIC=%RETRY%{groupName}、REAL_QID保存到消息属性
- 将消息topic改写为 SCHEDULE_TOPIC_XXXX,queueId设置为delayLevel对应的队列
│
▼
④ 消息写入SCHEDULE_TOPIC_XXXX内部队列,此时消息处于等待状态,业务、消费者组都看不到
│
▼
⑤ ScheduleMessageService定时扫描SCHEDULE_TOPIC_XXXX队列,消息到期
│
▼
⑥ 读取消息属性中的REAL_TOPIC=%RETRY%{group},真正把消息写入%RETRY%{group}的commitLog与consumequeue
│
▼
⑦ 消费者组启动时自动订阅自己的%RETRY%{group},拉取重试消息,reconsumeTimes+1,再次执行业务
│
├── 消费成功 → 结束
└── 消费失败 → 回到②,继续循环
为什么用 %RETRY%{消费组} 而不是直接投回原 topic?
和Mafka同样道理,不能因为某一个消费组失败,污染topic对应的所有消费组。
延迟到底靠什么实现?
SCHEDULE_TOPIC_XXXX 是全局唯一的一个内部调度 topic(不是每个业务 topic 配一个),内部按 18 个延迟等级分成 18 个队列。
ScheduleMessageService 用"定时任务 + 顺序扫描“实现——因为只有 18 个固定等级,同一等级队列里的消息延迟时间完全相同,所以"到期顺序 = 入队顺序”,天然有序、根本不用排序,从队头往后扫,遇到第一条没到期的就停。代价是只能延迟这 18 个固定档位。
三、三方对比:Kafka / Mafka / RocketMQ
| 维度 | Kafka(原生) | Mafka(美团,Kafka 改造) | RocketMQ(原生) |
|---|---|---|---|
| 底层血统 | 日志系统 | Kafka 系 | 消息队列系 |
| 确认模型 | offset 连续水位线 | offset 连续水位线 | 逐条 ack(consumeStatus) |
| 谁负责"等" | 业务应用层自己 | 外挂的 DelayServer 组件 | broker 内部(ScheduleMessageService) |
| 延迟载体 | 业务自建 | DelayServer 托管 | 全局唯一 SCHEDULE_TOPIC_XXXX |
| 重试消息去哪 | 重投到延迟 topic(业务自己发) | 到期投回该消费组专属的内部重试队列 | 投到 %RETRY%{消费组}(按组隔离) |
| offset 角色 | 照常提交,不参与重试 | 照常 commit,不参与重试 | 照常提交,不参与重试 |
| 退避档位 | 业务自建 | 支持任意延迟 | 固定 18 级(5.x 时间轮支持任意) |
| 死信 | 需自建 DLT | DLQ | 内建 %DLQ% |
| 谁写的脚手架 | 业务自己写全套 | 中间件团队封装进 SDK + 组件 | broker 原生自带 |
| 阻塞后续消息 | 卡 offset 就阻塞整个 partition | 不阻塞(失败消息挪走了) | 不阻塞(失败消息挪走了) |