RocketMQ 死信队列:16 次重试后的最后归宿
一句话结论(30s)
死信队列是 RocketMQ 消费失败的「最后防线」,本质是延迟递增重试 + 16 次上限 + 隔离三重兜底。因为延迟递增让前几次间隔短(网络抖动恢复后快速重试)、后几次间隔长(下游长时间宕机给足恢复时间),16 次上限防止无限重试阻塞队列,失败消息最终隔离到 %DLQ% 保证正常消息流不受影响。
核心原理(2min)
Consumer 返回 RECONSUME_LATER → 消息进入 %RETRY%{ConsumerGroup} 重试 Topic,按 1s/5s/10s/30s…2h 递增间隔重试,16 次仍失败则转入 %DLQ%{ConsumerGroup} 永久不重试。关键机制:运维订阅 DLQ 告警 → 排查修复 → 用 admin.resendMessageByHandleOffset() 手动回放到原始 Topic(先修 bug 再回放);DLQ 消息默认保留 48 小时作排查缓冲。
底层深入(5-10min)
重试的延迟递增
Consumer 消费失败返回 RECONSUME_LATER → 消息进入 %RETRY%{ConsumerGroup} 重试 Topic:
第 1 次: 1s 后重试
第 2 次: 5s 后重试
第 3 次: 10s 后重试
第 4 次: 30s 后重试
... (间隔递增)
第 16 次: 2h 后重试
重试次数的上限配置在消费组订阅配置 SubscriptionGroupConfig 里,默认就是 16 次:
private int retryQueueNums = 1;
private int retryMaxTimes = 16;
private GroupRetryPolicy groupRetryPolicy = new GroupRetryPolicy();
重试 Topic 的名字不是魔法字符串,而是由前缀常量拼出来的(MixAll.java):
public static final String RETRY_GROUP_TOPIC_PREFIX = "%RETRY%";
public static final String DLQ_GROUP_TOPIC_PREFIX = "%DLQ%";
public static String getRetryTopic(final String consumerGroup) {
return RETRY_GROUP_TOPIC_PREFIX + consumerGroup;
}
分析:retryMaxTimes 默认 16,是「16 次上限」的直接来源;%RETRY% + groupName 让每条消费组拥有独立的重试 Topic,重试进度(reconsumeTimes)随消息属性传递,避免不同组之间互相干扰。
💭 思考:为什么重试和死信要拆成
%RETRY%和%DLQ%两个 Topic,而不是都在一个 Topic 里靠属性区分?—— 因为两者的语义和处理方式完全不同:重试消息还要继续投递、继续计数,死信消息要「永久不再碰」。拆成两个 Topic,让「还在重试的」和「已经放弃的」物理隔离,正常消费流只盯 %RETRY%,运维只告警 %DLQ%,各自职责清晰,也不会互相拖累。
16 次全部失败 → 消息转入 %DLQ%{ConsumerGroup} 死信队列,永久不重试。
为什么用延迟递增?
如果每次都等间隔重试(比如全部 10s),消费问题的根因(如 bug、下游宕机、资源不足)不会在短时间内修复,等间隔重试会在下游恢复前反复失败浪费资源。
延迟递增让重试逐步退避——前几次间隔短(网络抖动临时恢复后快速重试),后几次间隔长(下游长时间宕机则给足够时间恢复)。16 次上限保证不会无限重试阻塞消息队列。
🤔 思考穿插:为什么重试要用「延迟递增」而不是固定间隔?—— 因为消费失败的根因(bug、下游宕机)往往要一段时间才恢复,固定间隔会在恢复前反复失败、白费资源;延迟递增让前几次快速重试「网络抖动」这类瞬时故障、后几次留足恢复时间。所以它是「退避重试」思想的落地。
进 DLQ 的判断:%DLQ% + group
getDLQTopic 与 getRetryTopic 是一对镜像方法,死信 Topic 同样是「前缀 + 消费组」:
public static String getDLQTopic(final String consumerGroup) {
return DLQ_GROUP_TOPIC_PREFIX + consumerGroup;
}
真正做「进不进 DLQ」判断的是 Broker 侧 SendMessageProcessor.handleRetryAndDLQ,它同时处理重试与死信两条路径:
private boolean handleRetryAndDLQ(SendMessageRequestHeader requestHeader, RemotingCommand response,
RemotingCommand request,
MessageExt msg, TopicConfig topicConfig, Map<String, String> properties) {
String newTopic = requestHeader.getTopic();
if (null != newTopic && newTopic.startsWith(MixAll.RETRY_GROUP_TOPIC_PREFIX)) {
String groupName = KeyBuilder.parseGroup(newTopic);
SubscriptionGroupConfig subscriptionGroupConfig =
this.brokerController.getSubscriptionGroupManager().findSubscriptionGroupConfig(groupName);
if (null == subscriptionGroupConfig) {
response.setCode(ResponseCode.SUBSCRIPTION_GROUP_NOT_EXIST);
response.setRemark(
"subscription group not exist, " + groupName + " " + FAQUrl.suggestTodo(FAQUrl.SUBSCRIPTION_GROUP_NOT_EXIST));
return false;
}
int maxReconsumeTimes = subscriptionGroupConfig.getRetryMaxTimes();
if (request.getVersion() >= MQVersion.Version.V3_4_9.ordinal() && requestHeader.getMaxReconsumeTimes() != null) {
maxReconsumeTimes = requestHeader.getMaxReconsumeTimes();
}
int reconsumeTimes = requestHeader.getReconsumeTimes() == null ? 0 : requestHeader.getReconsumeTimes();
boolean sendRetryMessageToDeadLetterQueueDirectly = false;
if (!brokerController.getRebalanceLockManager().isLockAllExpired(groupName)) {
LOGGER.info("Group has unexpired lock record, which show it is ordered message, send it to DLQ "
+ "right now group={}, topic={}, reconsumeTimes={}, maxReconsumeTimes={}.", groupName,
newTopic, reconsumeTimes, maxReconsumeTimes);
sendRetryMessageToDeadLetterQueueDirectly = true;
}
if (reconsumeTimes > maxReconsumeTimes || sendRetryMessageToDeadLetterQueueDirectly) {
Attributes attributes = this.brokerController.getBrokerMetricsManager().newAttributesBuilder()
.put(LABEL_CONSUMER_GROUP, requestHeader.getProducerGroup())
.put(LABEL_TOPIC, requestHeader.getTopic())
.put(LABEL_IS_SYSTEM, BrokerMetricsManager.isSystem(requestHeader.getTopic(), requestHeader.getProducerGroup()))
.build();
this.brokerController.getBrokerMetricsManager().getSendToDlqMessages().add(1, attributes);
properties.put(MessageConst.PROPERTY_DELAY_TIME_LEVEL, "-1");
newTopic = MixAll.getDLQTopic(groupName);
int queueIdInt = randomQueueId(DLQ_NUMS_PER_GROUP);
topicConfig = this.brokerController.getTopicConfigManager().createTopicInSendMessageBackMethod(newTopic,
DLQ_NUMS_PER_GROUP,
PermName.PERM_WRITE | PermName.PERM_READ, 0
);
msg.setTopic(newTopic);
msg.setQueueId(queueIdInt);
msg.setDelayTimeLevel(0);
if (null == topicConfig) {
response.setCode(ResponseCode.SYSTEM_ERROR);
response.setRemark("topic[" + newTopic + "] not exist");
return false;
}
}
}
int sysFlag = requestHeader.getSysFlag();
if (TopicFilterType.MULTI_TAG == topicConfig.getTopicFilterType()) {
sysFlag |= MessageSysFlag.MULTI_TAGS_FLAG;
}
msg.setSysFlag(sysFlag);
return true;
}
分析:判断入口是 newTopic.startsWith("%RETRY%"),用 KeyBuilder.parseGroup 从重试 Topic 名反解出消费组 groupName;拿到该组的 retryMaxTimes(默认 16)后,与消息携带的 reconsumeTimes 比较,一旦 reconsumeTimes > maxReconsumeTimes(或顺序消费的锁未过期需要直接进 DLQ),就把 Topic 换成 MixAll.getDLQTopic(groupName) 即 %DLQ% + group,同时清空延迟等级、随机落到 DLQ 队列并落盘。这就是「16 次后进 DLQ、按组隔离」的代码依据。
💭 思考:为什么顺序消费一旦检测到锁未过期,就直接进 DLQ、不等满 16 次?—— 因为顺序消费里一条消息失败会卡住整队:若还按 16 次重试慢慢来,这条毒消息会把后面所有消息都堵死。直接进 DLQ 是「牺牲这一条、保全整队」的取舍——尽快把毒消息移出队列,让后续消息恢复消费。
🤔 思考穿插:为什么要用「重试次数 > 上限」来判断进 DLQ,而不是「重试了 N 秒」?—— 因为消费失败的开销是按「次数」计算的(每次都要重新投递、重新消费),不是按时间;用次数上限能精确控制「最多折腾几次」,避免无限重试拖垮队列。所以 reconsumeTimes 才是进 DLQ 的核心闸门。
DLQ 处理流程
1. 告警:订阅 DLQ Topic → 收到消息 → 触发钉钉/企微/邮件告警
2. 排查:业务日志 + 消息内容 → 分析失败原因
3. 修复:修 bug 或补偿下游数据
4. 回放:admin.resendMessageByHandleOffset() → 将 DLQ 消息重新投递到原始 Topic
DLQ 消息不要在原 ConsumerGroup 直接重消费——先把 bug 修好,再手动回放。
🤔 思考穿插:DLQ 消息为什么不直接让原 ConsumerGroup 再消费一次?—— 因为 bug 没修、下游没恢复,直接重消费只会再次失败、再次进 DLQ,纯属浪费;正确顺序是「先修 bug → 再手动回放」。所以 DLQ 是「人工介入的信号」,而不是自动重试的第二通道。
总结
死信队列是 RocketMQ 的”最后防线”——正常消息流不受问题消息影响(隔离到 DLQ),问题消息不会无限重试浪费资源(16 次上限),运维有足够缓冲排查和修复(DLQ 消息默认保留 48 小时)。
章末提问
追问 1:一条消息消费失败后,RocketMQ 是怎么一步步把它送进死信队列的?
结论先行:消费失败返回 RECONSUME_LATER → 进入 %RETRY% 重试 Topic 按递增间隔重试,累计 16 次仍失败 → 转 %DLQ% 永久不再重试。
因为:重试 Topic 按消费组隔离,reconsumeTimes 随消息属性传递;Broker 在 handleRetryAndDLQ 里判断 reconsumeTimes > maxReconsumeTimes 时,把 Topic 换成 %DLQ% + group 并落盘。16 次上限保证了不会无限重试,死信则把毒消息隔离出正常流。
追问 2:为什么重试间隔要递增(退避),而不是固定间隔?
结论先行:因为消费失败的根因往往需要一段时间才能恢复,递增退避能在恢复前避免反复无效重试。
因为:前几次间隔短,能快速重试「网络抖动」这类瞬时故障;后几次间隔长,给下游长时间宕机留足恢复时间。固定间隔会在根因恢复前反复失败、白费资源,所以用「逐步退避」更合理。
追问 3:顺序消费和普通并发消费,进 DLQ 的判断有什么区别?
结论先行:普通消费按 reconsumeTimes > maxReconsumeTimes 判断;顺序消费一旦检测到该队列的锁未过期(说明是有序消息),会直接进 DLQ,不再等满 16 次。
因为:顺序消费里一条失败会卡住整队,若还等 16 次重试会长时间阻塞后续所有消息;直接进 DLQ 能尽快把毒消息移出、恢复整队消费。代码里 sendRetryMessageToDeadLetterQueueDirectly = true 正是这条路径。