Skip to content
Go back

RocketMQ死信队列——16次重试后的最终兜底

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

getDLQTopicgetRetryTopic 是一对镜像方法,死信 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 正是这条路径。


Share this post on:

Previous Post
RocketMQ 消息幂等性:为什么 MessageId 不可靠,以及如何真正去重
Next Post
RocketMQ架构——为什么CommitLog所有Topic共用一份文件