尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

RocketMq消息重试机制:消费失败重试、重试队列、死信队列原理与实战

RocketMq消息重试机制:消费失败重试、重试队列、死信队列原理与实战 消息重试机制消费失败重试、重试队列、死信队列原理与实战作者黒漂技术佬适用场景无人售货柜、智慧农业、工控物联网一、为什么需要消息重试先想一个场景你在无人售货柜上扫码开门拿了一瓶可乐柜子里的主板收到出货指令后——网络突然抖了一下指令没执行成功。怎么办用户已经拿货走人了订单却没完成。这就是消息重试要解决的问题。在分布式系统中以下三种情况几乎无法避免失败类型举例说明网络抖动消费者和Broker之间网络闪断消息没收到服务不可用下游数据库挂了、第三方支付接口超时业务异常库存不足、余额不够、数据校验不通过如果没有重试机制这些消息就丢了业务数据就不一致了。RocketMQ在设计之初就内置了完善的重试机制分生产者重试和消费者重试两个层面。二、生产者重试发送失败自动重试生产者发送消息到Broker时如果遇到网络异常或Broker响应超时RocketMQ会自动重试。DefaultMQProducerproducernewDefaultMQProducer(vending_producer_group);// 发送失败重试次数默认2次加上首次发送共3次producer.setRetryTimesWhenSendFailed(3);// 发送超时时间默认3000msproducer.setSendMsgTimeout(3000);producer.start();MessagemsgnewMessage(VendingTopic,出货指令.getBytes());SendResultresultproducer.send(msg);关键参数说明retryTimesWhenSendFailed同步发送失败时的重试次数默认2retryTimesWhenSendAsyncFailed异步发送失败时的重试次数默认2sendMsgTimeout发送超时时间默认3000ms小白提示生产者重试是发送阶段的保障确保消息成功到达Broker。但消息到达Broker后消费者消费失败就是另一回事了。三、消费者重试消费失败后延迟重试这才是重头戏。消费者拿到消息后如果业务处理抛异常或返回RECONSUME_LATERRocketMQ不会直接丢弃而是把消息放进重试队列过一会儿再投递给消费者。3.1 默认重试策略RocketMQ消费者默认最多重试16次间隔逐步增大重试次数间隔时间重试次数间隔时间第1次10s第9次7m第2次30s第10次8m第3次1m第11次9m第4次2m第12次10m第5次3m第13次20m第6次4m第14次30m第7次5m第15次1h第8次6m第16次2h注意上面的间隔是开源版RocketMQ的默认延迟级别。实际间隔取决于延迟消息的实现不同版本可能略有差异。16次重试全部失败后总耗时约4.6小时消息会被转入死信队列。3.2 重试队列的本质RocketMQ内部有一个特殊的Topic%RETRY%ConsumerGroup。消费失败的消息会被存到这个Topic里按照延迟级别投递。DefaultMQPushConsumerconsumernewDefaultMQPushConsumer(vending_consumer_group);// 最大重试次数默认16consumer.setMaxReconsumeTimes(5);// 设置为5次适合快速失败场景consumer.registerMessageListener(newMessageListenerConcurrently(){OverridepublicConsumeConcurrentlyStatusconsumeMessage(ListMessageExtmsgs,ConsumeConcurrentlyContextcontext){for(MessageExtmsg:msgs){try{// 业务处理执行出货指令processVendingCommand(msg);}catch(Exceptione){// 消费失败稍后重试intretryTimesmsg.getReconsumeTimes();log.warn(出货指令消费失败第{}次重试msgId{},retryTimes,msg.getMsgId());returnConsumeConcurrentlyStatus.RECONSUME_LATER;}}returnConsumeConcurrentlyStatus.CONSUME_SUCCESS;}});关键点msg.getReconsumeTimes()获取当前消息已重试次数返回RECONSUME_LATER触发重试setMaxReconsumeTimes可以自定义最大重试次数四、死信队列Dead Letter Queue4.1 什么是死信队列当一条消息重试16次或自定义次数仍然失败RocketMQ认为这条消息没救了把它转入死信队列。死信队列的Topic命名规则%DLQ%ConsumerGroup。死信消息的特征不会再被自动消费死信队列中的消息不会自动投递给消费者保留原始消息内容消息体、属性不变额外记录了死信原因有效期72小时默认保留3天过期自动清理4.2 查看死信消息通过RocketMQ控制台或命令行工具查看# 查看死信队列主题mqadmin topicList-n127.0.0.1:9876|grepDLQ# 查看死信消息内容mqadmin queryMsgByKey-n127.0.0.1:9876\-t%DLQ%vending_consumer_group\-kORDER_20250805_0014.3 死信处理方案死信不能放着不管必须有处理机制/** * 死信消息监控与处理 * 建议单独起一个消费者订阅死信Topic */ComponentpublicclassDeadLetterHandler{AutowiredprivateAlertServicealertService;AutowiredprivateRedisTemplateString,StringredisTemplate;publicvoidstartDeadLetterConsumer()throwsMQClientException{DefaultMQPushConsumerdlqConsumernewDefaultMQPushConsumer(dlq_monitor_group);// 订阅死信队列dlqConsumer.subscribe(%DLQ%vending_consumer_group,*);dlqConsumer.registerMessageListener((MessageListenerConcurrently)(msgs,context)-{for(MessageExtmsg:msgs){// 1. 记录死信日志log.error(收到死信消息! msgId{}, topic{}, 重试次数{},msg.getMsgId(),msg.getTopic(),msg.getReconsumeTimes());// 2. 触发告警通知钉钉/企业微信/短信alertService.sendUrgentAlert(售货柜出货指令消费失败已进入死信队列\n消息ID: msg.getMsgId()\n消息内容: newString(msg.getBody())请人工介入处理);// 3. 存入死信记录表方便后续重新投递saveDeadLetterRecord(msg);}returnConsumeConcurrentlyStatus.CONSUME_SUCCESS;});dlqConsumer.start();}}三种常见处理策略策略适用场景实现方式人工介入关键业务、不可丢失告警通知运维人员手动处理告警通知所有死信消息钉钉/企业微信机器人推送重新投递临时故障已恢复从死信记录表读取重新发送到原Topic五、无人售货柜场景实战5.1 出货指令消费失败的重试和告警方案完整流程出货指令发送 → 消费者收到 → 尝试执行 ↓ ↓ 消息持久化 成功→ 是 → ACK完成 ↓ 否 进入重试队列(%RETRY%) ↓ 重试1~5次间隔递增 ↓ 仍然失败→ 进入死信队列(%DLQ%) ↓ 钉钉告警 人工介入核心代码ComponentRocketMQMessageListener(topicvending_command_topic,consumerGroupvending_command_consumer_group,maxReconsumeTimes5// 最多重试5次)publicclassVendingCommandConsumerimplementsRocketMQListenerMessageExt{OverridepublicvoidonMessage(MessageExtmessage){StringcommandnewString(message.getBody());intretryCountmessage.getReconsumeTimes();try{// 解析出货指令VendingCommandcmdJSON.parseObject(command,VendingCommand.class);// 调用售货柜硬件接口执行出货booleansuccessvendingMachineService.executeCommand(cmd);if(!success){// 硬件返回失败触发重试thrownewRuntimeException(售货柜出货失败设备ID: cmd.getDeviceId());}log.info(出货指令执行成功, deviceId{}, retryCount{},cmd.getDeviceId(),retryCount);}catch(Exceptione){log.error(出货指令消费失败, retryCount{},retryCount,e);// 第3次重试开始升级告警if(retryCount3){alertService.sendWarning(售货柜指令重试retryCount次仍未成功\n设备ID: extractDeviceId(message)\n异常: e.getMessage());}// 抛异常触发重试thrownewRuntimeException(e);}}}5.2 分级重试策略不是所有失败都值得重试16次。建议按业务场景设置不同策略// 核心交易出货指令快速失败5次重试consumer.setMaxReconsumeTimes(5);// 非核心状态上报容忍度高16次重试consumer.setMaxReconsumeTimes(16);// 临时性操作广告推送不重试consumer.setMaxReconsumeTimes(1);六、重试的副作用消息可能重复消费这是重试机制带来的最大副作用——消息重复。想象这个流程消费者收到消息 → 业务执行成功 → 准备发送ACK ↓ 网络断了ACK没发出去 ↓ Broker没收到ACK → 认为消费失败 ↓ 重新投递消息 → 消费者再次消费结果同一条消息被消费了两次。如果这条消息是扣款10元用户就被扣了20元。所以重试机制必须配合幂等设计使用。幂等性设计的具体方案我们在下一篇详细展开。七、小结机制作用关键配置生产者重试保证消息到达BrokerretryTimesWhenSendFailed消费者重试消费失败后延迟重试maxReconsumeTimes重试队列临时存储重试消息%RETRY%ConsumerGroup死信队列存储最终失败消息%DLQ%ConsumerGroup核心原则重试解决暂时性故障死信兜底永久性故障幂等防护重复消费。
返回列表