

下面是一个生产级别的 Spring Boot 整合 RocketMQ 死信队列(DLQ)的完整示例,涵盖:配置、生产者、正常消费者(模拟失败触发重试)、死信队列消费者、监控告警、幂等与补偿设计。 一、死信队列核心概念速览 在 RocketMQ 中,死信队列(Dead Letter Queue, DLQ) 用于存储那些消费失败且达到最大重试次数的消息。 关键特性: 死信队列命名规则:%DLQ%{consumerGroup} 死信队列按消费者组创建,同一组内所有消费者产生的死信汇聚到同一个 DLQ 死信消息不会被原消费者组再次消费,需要新建消费者组单独订阅 DLQ 才能处理 死信消息默认保留 3 天,超时自动清理 二、项目依赖与版本
org.springframework.boot spring-boot-starter-parent 3.2.0 org.apache.rocketmq rocketmq-spring-boot-starter 2.2.3 org.apache.rocketmq rocketmq-client 5.1.0 org.projectlombok lombok true org.springframework.boot spring-boot-starter-logging 三、配置文件(application.yml) spring: application: name: rocketmq-dlq-demo rocketmq: name-server: 127.0.0.1:9876 # 生产环境建议配置多个,分号分隔 # ========== 生产者配置 ========== producer: group: order-producer-group send-message-timeout: 3000 retry-times-when-send-failed: 2 retry-next-server: true # ========== 消费者配置(正常业务消费者) ========== consumer: group: order-consumer-group # ⚠️ 这个 Group 将产生死信队列 %DLQ%order-consumer-group consume-mode: CLUSTERING # 集群模式(重试只对集群模式生效) consume-thread-min: 5 consume-thread-max: 20 pull-batch-size: 32 # ========== 死信队列的消费者组单独配置 ========== dlq: rocketmq: consumer: group: dlq-handler-group # 专门处理死信的消费者组,独立于原 Group ⚠️ 关键点:死信队列的消费者组必须独立于原消费者组,否则无法消费 DLQ 中的消息。 四、生产者代码 package com.example.rocketmq.producer; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; import org.apache.rocketmq.client.producer.SendResult; import org.apache.rocketmq.client.producer.SendStatus; import org.apache.rocketmq.spring.core.RocketMQTemplate; import org.apache.rocketmq.spring.support.RocketMQHeaders; import org.springframework.messaging.Message; import org.springframework.messaging.support.MessageBuilder; import org.springframework.stereotype.Service; @Slf4j @Service @RequiredArgsConstructor public class OrderProducer { private final RocketMQTemplate rocketMQTemplate; /** * 发送订单消息(同步发送,必须检查结果) */ public void sendOrder(String orderId, String content) { String destination = "order-topic:order-create"; Message
message = MessageBuilder .withPayload(content) .setHeader(RocketMQHeaders.KEYS, orderId) // 业务 Key,用于查询和幂等 .build(); try { SendResult result = rocketMQTemplate.syncSend(destination, message); // ✅ 生产级别:必须检查发送状态 if (result.getSendStatus() != SendStatus.SEND_OK) { log.error("消息发送失败:orderId={}, status={}, msgId={}", orderId, result.getSendStatus(), result.getMsgId()); // 走补偿逻辑(落本地消息表 or 告警) return; } log.info("消息发送成功:orderId={}, msgId={}", orderId, result.getMsgId()); } catch (Exception e) { log.error("消息发送异常:orderId={}", orderId, e); // 走补偿逻辑 } } } 五、正常消费者(模拟失败触发重试与死信) package com.example.rocketmq.consumer; import lombok.extern.slf4j.Slf4j; import org.apache.rocketmq.spring.annotation.ConsumeMode; import org.apache.rocketmq.spring.annotation.RocketMQMessageListener; import org.apache.rocketmq.spring.core.RocketMQListener; import org.springframework.stereotype.Component; @Slf4j @Component @RocketMQMessageListener( topic = "order-topic", consumerGroup = "order-consumer-group", // ⚠️ 这个 Group 会产生死信队列 selectorExpression = "order-create", consumeMode = ConsumeMode.CONCURRENTLY, maxReconsumeTimes = 3 // 🔥 为了演示效果,重试3次就进DLQ(生产环境默认16次) ) public class OrderConsumer implements RocketMQListener { private int retryCount = 0; @Override public void onMessage(String message) { retryCount++; log.info("收到消息(第{}次尝试):{}", retryCount, message); // 🔥 模拟业务处理失败 —— 触发重试,最终进入死信队列 // 生产环境中,这里可能是:调用下游失败、数据校验不通过等 throw new RuntimeException("模拟消费失败,触发重试!"); } } 💡 生产环境建议:maxReconsumeTimes 默认 16 次,可根据业务重要性调整。演示环境设为 3 次方便快速看到效果。 六、死信队列消费者(核心) package com.example.rocketmq.consumer.dlq; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; import org.apache.rocketmq.common.message.MessageExt; import org.apache.rocketmq.spring.annotation.RocketMQMessageListener; import org.apache.rocketmq.spring.core.RocketMQListener; import org.springframework.stereotype.Component; /** * 死信队列消费者 * * 专门订阅 %DLQ%order-consumer-group,处理所有进入死信的消息 */ @Slf4j @Component @RocketMQMessageListener( topic = "%DLQ%order-consumer-group", // 🔥 死信队列 Topic 命名规则 consumerGroup = "dlq-handler-group", // 🔥 必须使用独立的消费者组 selectorExpression = "*" // 消费所有 Tag ) public class DeadLetterQueueConsumer implements RocketMQListener { private final DlqMessageService dlqMessageService; public DeadLetterQueueConsumer(DlqMessageService dlqMessageService) { this.dlqMessageService = dlqMessageService; } @Override public void onMessage(MessageExt message) { String msgId = message.getMsgId(); String body = new String(message.getBody()); String keys = message.getKeys(); String originalTopic = message.getProperty("RETRY_TOPIC"); // 获取原始 Topic int reconsumeTimes = Integer.parseInt( message.getProperty("RECONSUME_TIMES") // 已经重试的次数 ); log.warn("========== 死信消息到达 =========="); log.warn("MsgId: {}", msgId); log.warn("Keys: {}", keys); log.warn("原始 Topic: {}", originalTopic); log.warn("已重试次数: {}", reconsumeTimes); log.warn("消息体: {}", body); log.warn("================================="); // 🔥 生产级处理:根据业务场景选择策略 dlqMessageService.handleDeadLetter(message); } } 七、死信处理服务(核心业务逻辑) package com.example.rocketmq.service; import lombok.extern.slf4j.Slf4j; import org.apache.rocketmq.common.message.MessageExt; import org.springframework.stereotype.Service; import java.nio.charset.StandardCharsets; @Slf4j @Service public class DlqMessageService { /** * 处理死信消息 —— 生产级核心逻辑 * * 策略选择(根据业务场景): * 1. 自动重试:将消息重新发送到原 Topic * 2. 落库补偿:存入数据库,人工审核后补偿 * 3. 直接告警:通知开发人工介入 * 4. 忽略跳过:非核心业务可跳过 */ public void handleDeadLetter(MessageExt message) { String body = new String(message.getBody(), StandardCharsets.UTF_8); String keys = message.getKeys(); String originalTopic = message.getProperty("RETRY_TOPIC"); // ---------- 策略一:存入数据库,等待人工补偿(推荐) ---------- // saveToCompensationTable(message); // 发送告警通知 // ---------- 策略二:自动重新投递(需谨慎,防止死循环) ---------- // 只有明确知道问题已修复(如下游恢复)才可启用 // resendToOriginalTopic(message); // ---------- 策略三:记录日志 + 告警(本 Demo 采用) ---------- log.error("【死信告警】消息进入死信队列,需人工介入!orderId={}, body={}", keys, body); // 生产环境:发送钉钉/企业微信/邮件告警 // sendAlert("死信消息需要人工处理", keys, body); } /** * 将死信消息重新发送到原始 Topic(谨慎使用) */ private void resendToOriginalTopic(MessageExt message) { // 实现重新发送逻辑 // ⚠️ 注意:需要防止无限循环,建议增加重投次数限制 } } 八、幂等消费设计(防止重复消费) RocketMQ 保证 “至少一次(At Least Once)” 语义,消费端必须实现幂等。 方案一:Redis SETNX 快速去重 @Component public class IdempotentHelper { @Autowired private StringRedisTemplate redisTemplate; public boolean tryProcess(String msgId, String eventType) { String key = "processed:" + eventType + ":" + msgId; Boolean success = redisTemplate.opsForValue() .setIfAbsent(key, "1", Duration.ofDays(7)); return Boolean.TRUE.equals(success); } } 方案二:数据库唯一索引(强一致) CREATE TABLE order_process_record ( id BIGINT PRIMARY KEY AUTO_INCREMENT, msg_key VARCHAR(64) NOT NULL, event_type VARCHAR(32) NOT NULL, status TINYINT DEFAULT 0, created_at DATETIME, UNIQUE KEY uk_msg_key_event (msg_key, event_type) ); 九、监控与告警(生产必备) 9.1 查看死信队列 # 查看所有 Topic,筛选 DLQ ./mqadmin topicList -n 127.0.0.1:9876 | grep "%DLQ%" # 查看死信队列的消费进度 ./mqadmin consumerProgress -n 127.0.0.1:9876 -g dlq-handler-group 9.2 Prometheus 告警规则 groups: - name: rocketmq_dlq rules: - alert: RocketMQDLQMessage expr: sum(rocketmq_dlq_messages) > 0 for: 5m annotations: summary: "RocketMQ 死信队列有消息,需要人工处理" 死信队列的监控可以通过 rocketmq-exporter 采集 DLQ 的 Offset 变化来实现告警。 十、完整的消息流转图 flowchart TB subgraph Producer[生产者] P[OrderProducer] -->|发送| Topic[order-topic] end subgraph Normal[正常消费链路] Topic --> Consumer[OrderConsumer
order-consumer-group] Consumer -->|消费失败| Retry[自动重试
maxReconsumeTimes=3] Retry -->|3次均失败| DLQ[死信队列
%DLQ%order-consumer-group] end subgraph DLQHandler[死信处理链路] DLQ --> DLQConsumer[DeadLetterQueueConsumer
dlq-handler-group] DLQConsumer --> Handler[DlqMessageService] Handler -->|策略一| DB[(存入数据库
人工补偿)] Handler -->|策略二| Alert[📱 发送告警] Handler -->|策略三| RePush[重新投递
到原 Topic] end style DLQ fill:#ffcdd2 style Handler fill:#fff3e0