一、先从一次线上事故说起
那天晚上十一点多,我正在家里准备睡觉,手机突然连续震动。打开一看,监控群里炸了锅:订单状态同步服务堆积了几十万条消息,消费端一直报错重试,最后消息全部进了死信队列。更麻烦的是,因为当时没有给死信队列配任何自动处理机制,这些消息就静静躺在那里,第二天业务方发现大量订单状态没更新,才手动去补救。
这件事让我意识到,死信队列不是终点,而是起点。很多人只知道RocketMQ有死信队列这个概念,但真正遇到生产问题时会发现,如果只重试不恢复,死信队列很快变成垃圾堆;如果只告警不自动处理,半夜爬起来手动捞消息就成了家常便饭。所以这篇文章,我不讲太高深的理论,就站在一个普通后端开发的角度,聊聊消费端死信队列的完整处理策略,重点放在“重试次数耗尽之后怎么自动化恢复”、“怎么设计告警才能既及时又不吵人”。
二、RocketMQ消费重试是怎么运作的
2.1 消息重试的核心机制
在RocketMQ里,消费者拉取到消息后会执行业务逻辑。如果抛出异常,RocketMQ并不会立刻把消息丢掉,而是按照预设的重试次数和延迟级别,把消息重新放回队列等待再次消费。这个重试对开发者来说几乎是透明的,你只需要在消费者代码里抛出异常就行。
举个最常见的例子,用Spring Boot接入RocketMQ,消费订单消息后调用远程接口更新库存,远程接口偶尔超时,我们希望失败后能重试几次。下面是一个最基础的消费端代码。
// 技术栈:Java + Spring Boot + RocketMQ Client
@Component
@RocketMQMessageListener(
topic = "ORDER_TOPIC",
consumerGroup = "ORDER_CONSUMER_GROUP"
)
public class OrderMessageListener implements RocketMQListener<OrderMessage> {
private static final Logger log = LoggerFactory.getLogger(OrderMessageListener.class);
@Override
public void onMessage(OrderMessage message) {
try {
// 这里是真实的业务处理:调用库存服务的远程接口
boolean success = inventoryClient.deductStock(message.getOrderId());
if (!success) {
// 返回false也认为消费失败,会触发重试
throw new RuntimeException("库存扣减失败,orderId=" + message.getOrderId());
}
log.info("订单消息消费成功:{}", message.getOrderId());
} catch (Exception e) {
log.error("订单消息消费异常,orderId={}", message.getOrderId(), e);
// 注意:这里没有重试代码,RocketMQ会自动进行重试
throw e;
}
}
}
上面的代码里,你不需要写任何循环重试的逻辑,RocketMQ会按照默认的延迟级别进行重试。默认情况下,一条消息最多重试16次,第一次重试延迟10秒,第二次30秒,之后逐级递增,最长时间为2小时。当重试次数达到上限后,消息就会被放入死信队列。
2.2 死信队列长什么样
每个消费组都有一个自己的死信队列,命名规则是“%DLQ%消费组名称”。比如你的消费者组叫ORDER_CONSUMER_GROUP,那么它的死信队列就是%DLQ%ORDER_CONSUMER_GROUP。
死信队列不是一个单独独立的物理队列,逻辑上它是一个特殊的Topic,里面存的就是重试耗尽的消息。当你打开RocketMQ控制台,就能看到这个队列里堆积了多少消息,消息体还是原来的内容,只不过额外带了重试次数等系统属性。
理解这两个概念后,下一步就是考虑:死信队列里的消息总不能一直躺着吧?怎么让它们自动恢复呢?
三、死信队列的自动化恢复方案
3.1 方案一:定时任务扫描死信队列并重新投递
最简单、也最容易上手的方式,就是写一个定时任务,周期性来扫描死信队列,把里面的消息取出来,重新发到正常的Topic里,让消费者再次处理。这个方案的好处是不需要改RocketMQ本身的任何配置,只靠普通业务代码就能搞定。
有人可能会问:直接重新发回原Topic,会不会又重试16次然后又进死信?这是必然的。所以我们在重新投递前要做一个判断,比如记录这条消息已经重投过几次了,如果重投超过3次,就说明这是一个真正的“坏消息”,不能再投了,直接转人工或者落库。
来看一个具体的代码实现。这里我们使用RocketMQ的官方Java客户端,写一个定时任务,从死信队列拉取消息,先尝试解析业务数据,然后重新投递。
// 技术栈:Java + Spring Boot + RocketMQ Client 4.9.x
@Component
public class DldRecoveryTask {
private static final Logger log = LoggerFactory.getLogger(DldRecoveryTask.class);
// 正常业务Topic
private static final String ORIGIN_TOPIC = "ORDER_TOPIC";
// 死信队列名称:%DLQ%消费组名
private static final String DLQ_TOPIC = "%DLQ%ORDER_CONSUMER_GROUP";
// 最大重投次数,超过这个次数就不再自动恢复
private static final int MAX_RESEND_TIMES = 3;
@Autowired
private RocketMQTemplate rocketMQTemplate;
/**
* 每5分钟执行一次恢复任务
* 注意:这里使用了Spring的定时任务,需要开启@EnableScheduling
*/
@Scheduled(fixedDelay = 5 * 60 * 1000, initialDelay = 10 * 1000)
public void recoverDeadMessages() {
log.info("开始扫描死信队列:{}", DLQ_TOPIC);
// 创建一个消费者实例,只消费死信队列中的消息
DefaultMQPullConsumer pullConsumer = new DefaultMQPullConsumer("DLQ_RECOVERY_CONSUMER");
pullConsumer.setNamesrvAddr("127.0.0.1:9876");
try {
pullConsumer.start();
// 从每个队列拉取消息,这里简化处理,只拉取第一个队列
// 实际上RocketMQ的死信队列可能有多个队列,可以循环处理
MessageQueue messageQueue = new MessageQueue(DLQ_TOPIC, "broker-a", 0);
pullConsumer.setStartDeliverTime(System.currentTimeMillis() - 10 * 60 * 1000);
PullResult pullResult = pullConsumer.pullBlocking(messageQueue, "*", 0, 32);
if (pullResult.getPullStatus() != PullStatus.FOUND) {
log.info("死信队列中没有消息");
return;
}
// 遍历拉取到的消息
for (MessageExt msg : pullResult.getMsgFoundList()) {
handleDeadMessage(msg);
}
} catch (Exception e) {
log.error("扫描死信队列发生异常", e);
} finally {
pullConsumer.shutdown();
}
}
/**
* 处理单条死信消息
*/
private void handleDeadMessage(MessageExt msg) {
// 从用户属性里获取重投次数,如果没有,默认0
String resendTimesStr = msg.getUserProperty("resendTimes");
int resendTimes = resendTimesStr == null ? 0 : Integer.parseInt(resendTimesStr);
log.info("接收到死信消息,orderId={},当前重投次数={}", msg.getKeys(), resendTimes);
// 如果重投次数达到上限,不再自动恢复,记录报警日志
if (resendTimes >= MAX_RESEND_TIMES) {
log.error("消息已经重投{}次,放弃自动恢复,消息ID={}", resendTimes, msg.getMsgId());
// 这里可以调用人工处理接口,或者入库存表,等待人工介入
return;
}
try {
// 重新构造一条新消息,发送到原始业务Topic
Message newMsg = new Message(
ORIGIN_TOPIC,
msg.getBody()
);
// 保留原始消息的业务Key,方便追踪
newMsg.setKeys(msg.getKeys());
// 给新消息打上重投次数的标记,下次恢复时可以判断
newMsg.putUserProperty("resendTimes", String.valueOf(resendTimes + 1));
// 发送到原始Topic,消费者会再次消费
SendResult sendResult = rocketMQTemplate.getProducer().send(newMsg);
if (sendResult.getSendStatus() == SendStatus.SEND_OK) {
log.info("死信消息重投成功,orderId={},新MsgId={}", msg.getKeys(), sendResult.getMsgId());
} else {
log.error("死信消息重投失败,orderId={}", msg.getKeys());
}
} catch (Exception e) {
log.error("死信消息重投异常,orderId={}", msg.getKeys(), e);
}
}
}
这段代码里有几个值得注意的地方:
第一,我用的是DefaultMQPullConsumer,也就是手动拉取模式。为什么不直接用PushConsumer?因为Push模式会自动把消息推给处理器,而死信队列里的消息不适合让消费者直接去消费,我们是想把它们捞出来重新投递,所以手动拉取更可控。
第二,把重投次数放在消息的自定义属性resendTimes里,这是非常关键的设计。如果不加这个标记,每次扫描死信队列都会发现有消息,然后无限循环重投,永远结束不了。
第三,每次重投时都重新构造一条新消息,这样做的原因是,原消息是从死信队列拉取出来的,如果直接发送这条消息本身,它的Topic属性还是死信队列,发不到正确的业务Topic去。构造新消息时把body和业务key复制过来就行了。
3.2 方案二:利用延迟消息做“二次消费”的优雅替代
如果你觉得上面定时扫描的方案有点笨重,RocketMQ本身也支持在消费失败时,把消息投递到一个自定义的延迟Topic,延迟一段时间后再由同样的消费者来处理。这种方式可以绕开默认的重试机制,让你自己掌握重试节奏。
比如,消费订单消息失败后,我们不是抛出异常让它自动重试,而是把消息发送到一个ORDER_BIZ_DELAY的Topic,延迟10秒后再消费。在这个延迟Topic的消费逻辑里,同样判断业务处理是否成功,如果继续失败,再投递到延迟更长的Topic。当自定义重试次数达到上限后,直接发送到死信队列。
// 技术栈:Java + Spring Boot + RocketMQ Client
@Component
public class BizRetryConsumer {
private static final Logger log = LoggerFactory.getLogger(BizRetryConsumer.class);
// 业务延迟重试Topic
private static final String RETRY_TOPIC = "ORDER_BIZ_RETRY";
@Autowired
private RocketMQTemplate rocketMQTemplate;
/**
* 这是主消费者,处理订单消息
*/
@RocketMQMessageListener(
topic = "ORDER_TOPIC",
consumerGroup = "ORDER_CONSUMER_GROUP"
)
public static class MainListener implements RocketMQListener<OrderMessage> {
@Autowired
private RocketMQTemplate rocketMQTemplate;
@Override
public void onMessage(OrderMessage orderMessage) {
// 业务处理结果由你自己决定
boolean success = doBusinessProcess(orderMessage);
if (!success) {
// 不抛异常,而是把消息发到自定义重试Topic,延迟10秒
Message retryMsg = MessageBuilder.withPayload(orderMessage)
.setHeader("retryCount", 0)
.build();
rocketMQTemplate.syncSend("ORDER_BIZ_RETRY_Topic", retryMsg, 3000, 10);
log.info("发送业务重试消息,orderId={}", orderMessage.getOrderId());
}
}
}
/**
* 这是重试消费者,专门消费重试Topic的消息
*/
@RocketMQMessageListener(
topic = RETRY_TOPIC,
consumerGroup = "ORDER_RETRY_GROUP"
)
public static class RetryListener implements RocketMQListener<OrderMessage> {
private static final int MAX_RETRY = 3;
@Override
public void onMessage(OrderMessage orderMessage) {
// 获取这次重试是第几次
int retryCount = MessageUtil.getRetryCount(orderMessage);
boolean success = doBusinessProcess(orderMessage);
if (success) {
log.info("重试处理成功,orderId={}", orderMessage.getOrderId());
return;
}
// 重试次数达到上限,投递到死信队列
if (retryCount >= MAX_RETRY) {
log.error("自定义重试达到上限,发送到死信队列,orderId={}", orderMessage.getOrderId());
rocketMQTemplate.syncSend("%DLQ%ORDER_CONSUMER_GROUP", orderMessage);
return;
}
// 未达到上限,继续发送到重试Topic,把次数+1
orderMessage.setRetryCount(retryCount + 1);
rocketMQTemplate.syncSend(RETRY_TOPIC, orderMessage);
}
}
}
上面代码中的MessageUtil.getRetryCount是个抽象演示,实际实现时可以从消息的header或自定义属性中读取。这种方案最明显的好处是,重试间隔和次数完全由业务决定,不受RocketMQ默认16次或固定延迟级别的约束。缺点是需要自己管理额外Topic和消费组,代码上多了一层。
3.3 方案三:基于RocketMQ的定时消息功能做回放
如果你用的RocketMQ版本是4.9.x以上,或者用的是阿里云商业版RocketMQ,就支持定时消息。这时候可以把死信队列里的消息取出来,发送一条带延迟的定时消息,延迟比如10分钟后再投递到业务Topic。这个方案和方案一类似,但是比定时任务更实时,因为每条消息的投递时间都是独立计算的。
这里就不展开写完整代码了,思路是:你写一个专门监听死信队列的消费者,每当有死信消息到达时,不要马上处理,而是构造一条新的延迟消息,设置延迟级别,发到原业务Topic。这样不用轮询扫描,死信消息一到就能被重新计划投递。
四、告警设计:让报警恰到好处
自动化恢复不是万能的,有些消息可能反复失败,最终需要人工介入。如果没有任何告警,死信队列里的消息可能会悄悄积累几天。如果告警太频繁,比如每发一条死信就报警,运维和开发又会被骚扰到崩溃。所以告警设计要分级别、分节奏。
4.1 死信消息积压数量告警
这是最基础的告警,每隔一段时间统计死信队列的消息数量,超过阈值就报警。推荐使用Prometheus + Alertmanager来监控RocketMQ指标,RocketMQ的broker会暴露出死信队列的消息堆积数量。
# prometheus规则配置示例
groups:
- name: rocketmq_alerts
rules:
- alert: 死信队列消息积压过多
expr: rocketmq_message_accumulation{queue="%DLQ%ORDER_CONSUMER_GROUP"} > 100
for: 10m
labels:
severity: warning
annotations:
summary: "死信队列 {{ $labels.queue }} 积压 {{ $value }} 条消息"
description: "订单消费者组死信队列消息超过100条,请检查消费端是否出现持续异常"
这里的告警规则是:死信队列数量持续10分钟超过100条才报警,避免临时的波动引起误报。
4.2 单条消息重试次数告警
上面的积压告警解决了量的问题,但有些场景是死信队列数量不多,可某一类关键消息一直失败。比如用户下单后的余额扣减,如果连续失败,一条消息就足够引发严重问题。这时候需要监控单条消息的重试次数。
怎么实现呢?在自己的消费代码里,每次消费异常时,把重试次数作为日志或者指标打出来。这样可以通过日志系统聚合出“重试次数超过N次的消息”并告警。
// 技术栈:Java + SLF4J + Metrics(可以接入Prometheus)
try {
processMessage(message);
} catch (Exception e) {
// 获取当前重试次数,RocketMQ在重试时会把重试次数放到消息系统属性RECONSUME_TIMES中
String reconsumeTimes = message.getProperty("RECONSUME_TIMES");
int times = reconsumeTimes == null ? 0 : Integer.parseInt(reconsumeTimes);
log.warn("消息消费失败,orderId={},已重试{}次", message.getOrderId(), times);
// 如果重试次数超过5次,就上报一个业务指标,用于告警
if (times >= 5) {
metricsService.incrementCounter("order_consume_fail_over_5times", message.getOrderId());
}
throw e;
}
这样当order_consume_fail_over_5times这个指标在短时间内增长明显时,告警系统就能发现是某个订单或者某类订单在持续异常。
4.3 恢复任务执行结果告警
自动化恢复任务本身也可能失败,比如死信队列里的消息投递回原Topic时发送超时、RocketMQ服务端故障等。所以要给恢复任务本身加一个状态监控。在恢复任务结束后,把本次处理的消息数量、成功数量、失败数量记录下来,如果成功率低于某个阈值,就触发告警。
// 技术栈:Java + Spring Boot
public void recoverDeadMessagesWithCheck() {
int total = 0;
int success = 0;
int fail = 0;
// 执行恢复逻辑...
List<MessageExt> deadMessages = fetchDeadMessages();
for (MessageExt msg : deadMessages) {
total++;
boolean isOk = resend(msg);
if (isOk) {
success++;
} else {
fail++;
}
}
// 上报恢复执行指标
reportMetrics(total, success, fail);
// 如果成功率低于80%并且总数超过10条,发出告警
double successRate = total > 0 ? (double) success / total : 1.0;
if (total > 10 && successRate < 0.8) {
alertService.sendAlert("死信恢复成功率过低", String.format("本次共处理%d条,成功%d条,失败%d条", total, success, fail));
}
}
4.4 告警频率的冷静设计
告警不能像机关枪一样发射,否则大家看多了就麻木了。我比较推荐每次告警后设置一个冷却时间,比如同一个死信队列,如果已经发过告警,那么后续2小时内不再重复发同类告警,除非级别升级。另外,告警应该分级:
- 死信队列积压超过100条但少于1000条:一般告警,发到群,工作时间内处理。
- 超过1000条且持续增长:严重告警,电话通知。
- 关键业务死信恢复成功率低于50%:严重告警,立刻暂停恢复任务,防止消息雪崩。
你可能注意到,恢复任务失败时要暂停恢复任务,这个很重要。如果连续重投失败,死信恢复任务本身可能把RocketMQ发送线程打爆,导致正常消息也发不出去。所以当检测到连续失败时,应该自动熔断,停止恢复,让消息继续留在死信队列里,等待人工排查。
五、应用场景分析
死信队列自动化恢复策略并不是所有的项目都需要,如果你们的消费者逻辑非常简单,比如只做日志写入,几乎不会失败,那完全不需要搞这些。但下面这些场景强烈建议配备:
- 订单支付、库存扣减、余额变动等核心交易链路:一条消息失败影响很大,必须自动恢复。
- 对接第三方接口的消费逻辑:第三方接口不稳定是常态,自动重试恢复能省去大量人工补签。
- 数据同步场景:比如将订单数据同步到大数据平台,偶尔会出现字段格式错误,这类消息重试也没用,不需要自动恢复,只要告警提示人工修数据就行。
- 高并发场景:积压消息不能一直留在死信队列中,否则会占用大量存储空间,也导致数据延迟。
六、技术优缺点整理
上面几种恢复方案各有优劣,我列出来方便你按需选择。
定时任务扫描重投的优点是实现简单,不依赖RocketMQ的高级特性,任何版本都能用。缺点是存在实时性问题,定时任务最短间隔也至少1分钟,对时效性要求高的业务不太友好。同时每次扫描都会拉取死信队列,频繁拉取会占用少量网络开销。
自定义延迟Topic重试的优点是可以灵活控制重试次数和延迟时间,逻辑完全掌握在自己手里。缺点是代码复杂度提高,而且需要维护额外的Topic、消费组、以及消息流转链路的监控。
基于定时消息的方案最优雅,实时性高,但前提是RocketMQ版本支持。如果你们用的是老版本4.5.x,就享受不到这个特性。
再来说告警设计方面,积压数量告警只能看到堆积量,不能定位具体问题;单条消息重试次数告警时效性好,但需要侵入业务代码埋点。所以建议两者都做,一个是宏观监控,一个是微观感知。
七、关于这个主题,我特别想提醒你的几件事
第一,死信队列的消息要保留原始业务信息,尤其是消息key。后面不管是自动恢复还是人工排查,都要靠消息key找到对应订单或业务记录。
第二,重投次数标记一定要设计好。我看到过很多系统把死信消息反复重投,每次重投都是原样发回,最终无限循环,直到把整个RocketMQ集群拖垮。可以用消息自定义属性、消息体里的字段,或者Redis记录消息ID来去重。
第三,恢复任务必须考虑“幂等性”。即使你不是无限循环重投,重投后的消息也有可能会被消费两次。比如第一次恢复任务重投成功但超时了,实际上消费者已经处理成功了,你又重投了第二次,就会造成重复处理。所以你的业务消费代码本身最好具备幂等性,或者在恢复任务里记录每条消息的已投递状态。
第四,不要忽视死信队列本身的存储清理。死信队列里的消息不会自动消失,如果一直不消费,会占用broker磁盘。可以保留几天的死信消息,之后允许过期清理,同时把关键内容异步转储到数据库。
八、总结
RocketMQ的死信队列不是用来收藏问题的,而是用来给开发者一个缓冲和恢复的机会。重试次数耗尽不是终点,而是新问题处理的起点。我们既可以通过定时任务把消息捞回来重新投递,也可以利用自定义延迟Topic或定时消息来灵活控制重试节奏。重要的是,不管哪种恢复方式,都要给消息打上重投次数标记,避免死循环。
告警设计方面,要结合积压数量、单条重试次数、恢复任务成功率三个维度去建立监控。同时设置冷却时间和分级通知,让告警真正有效,而不是沦为背景噪音。
最后,希望你在设计自己的消费端死信处理方案时,不要只盯着“把消息消费掉”,而要多想一想:消息为什么会死?有没有办法让它复活?复活之后怎么避免再死?这才是工程上真正需要打磨的地方。
九、附录:一个完整的死信处理流程参考
为了方便你理解上面所有策略怎么组合起来,我画了一张伪代码流程图(这里不能贴图,就用文字描述一下)。
消费者抛出异常后,RocketMQ默认重试16次,重试期间如果成功则结束。失败进入死信队列后,恢复任务每5分钟扫描一次死信队列,读取消息并判断重投次数。如果重投次数小于3,则把消息重新发送到原业务Topic。如果重投次数已经等于3,则把消息写入数据库的“人工处理表”,同时触发高优告警。数据库里记录的消息会被专门的后台管理页面展示,运维人员可以手工查询原始内容,修改数据后手动重新投递。同时Prometheus监控死信队列的消息数量和恢复任务的失败率,超过阈值就发出不同级别的告警。
这样一套组合下来,既实现了自动化恢复,又保证了在极端场景下人工能够兜底,整个消费链路才算真正健壮。
Comments