消息中间件在高并发架构中承担异步解耦和削峰填谷职责。RabbitMQ作为主流消息中间件,在业务中台建设中广泛应用。消息消费失败时的处理策略直接影响系统可靠性。死信队列(Dead Letter Queue)将消费失败的消息路由到备用队列,配合重试机制实现消息可靠投递。本文记录Spring Boot集成RabbitMQ死信队列的完整配置和代码实现,涵盖服务治理和消息中间件运维要点。
RabbitMQ死信队列原理与架构设计
消息在以下三种情况下会成为死信:
- 消息被拒绝(basic.reject/basic.nack),且requeue参数为false
- 消息TTL过期:消息在队列中存活时间超过设定的TTL值
- 队列达到最大长度:队列消息数量超过x-max-length限制
架构设计:业务交换机将消息路由到业务队列,业务队列绑定死信交换机。消息消费失败后进入死信交换机,路由到死信队列。死信队列消费者进行告警通知和人工补偿处理。
业务交换机(exchange.business) --> 业务队列(queue.business)
|
消费失败/TTL过期
|
死信交换机(exchange.dlx)
|
死信队列(queue.dlx)
Spring Boot项目初始化与依赖配置
Maven依赖配置:
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-amqp</artifactId>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-web</artifactId>
</dependency>
application.yml配置RabbitMQ连接参数:
spring:
rabbitmq:
host: 192.168.1.50
port: 5672
username: admin
password: ${RABBITMQ_PASSWORD}
virtual-host: /production
publisher-confirm-type: correlated
publisher-returns: true
listener:
simple:
acknowledge-mode: manual
prefetch: 10
retry:
enabled: true
max-attempts: 3
initial-interval: 1000
multiplier: 2.0
队列与交换机声明配置
使用Java Config声明业务队列、死信交换机和死信队列的绑定关系:
import org.springframework.amqp.core.*;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
@Configuration
public class RabbitMQConfig {
// === 业务队列配置 ===
public static final String BUSINESS_EXCHANGE = "exchange.business";
public static final String BUSINESS_QUEUE = "queue.business";
public static final String BUSINESS_ROUTING_KEY = "routing.business";
// === 死信队列配置 ===
public static final String DLX_EXCHANGE = "exchange.dlx";
public static final String DLX_QUEUE = "queue.dlx";
public static final String DLX_ROUTING_KEY = "routing.dlx";
// 业务交换机
@Bean
public DirectExchange businessExchange() {
return ExchangeBuilder.directExchange(BUSINESS_EXCHANGE)
.durable(true)
.build();
}
// 死信交换机
@Bean
public DirectExchange dlxExchange() {
return ExchangeBuilder.directExchange(DLX_EXCHANGE)
.durable(true)
.build();
}
// 业务队列(绑定死信交换机)
@Bean
public Queue businessQueue() {
return QueueBuilder.durable(BUSINESS_QUEUE)
.withArgument("x-dead-letter-exchange", DLX_EXCHANGE)
.withArgument("x-dead-letter-routing-key", DLX_ROUTING_KEY)
.withArgument("x-message-ttl", 60000) // 消息TTL 60秒
.withArgument("x-max-length", 10000) // 队列最大长度
.build();
}
// 死信队列
@Bean
public Queue dlxQueue() {
return QueueBuilder.durable(DLX_QUEUE).build();
}
// 绑定关系
@Bean
public Binding businessBinding() {
return BindingBuilder.bind(businessQueue())
.to(businessExchange())
.with(BUSINESS_ROUTING_KEY);
}
@Bean
public Binding dlxBinding() {
return BindingBuilder.bind(dlxQueue())
.to(dlxExchange())
.with(DLX_ROUTING_KEY);
}
}
消息生产者实现与确认机制
生产者发送消息时启用Confirm回调,确保消息到达交换机和队列:
import org.springframework.amqp.rabbit.core.RabbitTemplate;
import org.springframework.amqp.rabbit.connection.CorrelationData;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Component;
import lombok.extern.slf4j.Slf4j;
import java.util.UUID;
@Slf4j
@Component
public class OrderMessageProducer {
@Autowired
private RabbitTemplate rabbitTemplate;
public void sendOrderMessage(OrderDTO order) {
CorrelationData correlationId = new CorrelationData(UUID.randomUUID().toString());
rabbitTemplate.convertAndSend(
RabbitMQConfig.BUSINESS_EXCHANGE,
RabbitMQConfig.BUSINESS_ROUTING_KEY,
order,
message -> {
message.getMessageProperties().setDeliveryTag(1);
message.getMessageProperties().setExpiration("30000"); // 单条消息TTL
return message;
},
correlationId
);
log.info("订单消息已发送: orderId={}, correlationId={}",
order.getOrderId(), correlationId.getId());
}
}
配置Confirm回调处理消息投递结果:
@Component
@Slf4j
public class RabbitConfirmCallback implements RabbitTemplate.ConfirmCallback {
@Override
public void confirm(CorrelationData correlationData, boolean ack, String cause) {
if (ack) {
log.info("消息投递成功: correlationId={}", correlationData.getId());
} else {
log.error("消息投递失败: correlationId={}, cause={}",
correlationData.getId(), cause);
// 投递失败处理:记录到数据库,定时任务重发
}
}
}
消息消费者实现与异常处理
消费者手动确认消息,处理失败时reject到死信队列:
import org.springframework.amqp.rabbit.annotation.RabbitListener;
import org.springframework.amqp.core.Message;
import org.springframework.amqp.rabbit.core.ChannelAwareMessageListener;
import com.rabbitmq.client.Channel;
import org.springframework.stereotype.Component;
import lombok.extern.slf4j.Slf4j;
@Slf4j
@Component
public class OrderMessageConsumer {
@RabbitListener(queues = RabbitMQConfig.BUSINESS_QUEUE)
public void handleMessage(OrderDTO order, Message message, Channel channel) throws Exception {
long deliveryTag = message.getMessageProperties().getDeliveryTag();
try {
// 业务处理
processOrder(order);
// 处理成功,手动ACK
channel.basicAck(deliveryTag, false);
log.info("订单处理成功: orderId={}", order.getOrderId());
} catch (BusinessException e) {
// 业务异常,判断重试次数
Integer retryCount = getRetryCount(message);
if (retryCount >= 3) {
// 超过重试次数,reject到死信队列
log.error("订单处理失败,超过重试次数,进入死信队列: orderId={}, retry={}",
order.getOrderId(), retryCount);
channel.basicReject(deliveryTag, false);
} else {
// 未超过重试次数,重新入队重试
log.warn("订单处理失败,重试中: orderId={}, retry={}",
order.getOrderId(), retryCount);
channel.basicNack(deliveryTag, false, true);
}
} catch (Exception e) {
// 未知异常,直接进入死信队列
log.error("订单处理未知异常,进入死信队列: orderId={}", order.getOrderId(), e);
channel.basicReject(deliveryTag, false);
}
}
private void processOrder(OrderDTO order) {
// 扣减库存、创建物流单、发送通知等业务逻辑
inventoryService.deduct(order.getProductId(), order.getQuantity());
logisticsService.createShipment(order);
}
private Integer getRetryCount(Message message) {
Object count = message.getMessageProperties()
.getHeader("x-retry-count");
return count != null ? (Integer) count : 0;
}
}
死信队列消费者与告警补偿
死信队列消费者负责告警通知和人工补偿处理:
@Slf4j
@Component
public class DeadLetterConsumer {
@Autowired
private AlertService alertService;
@Autowired
private OrderFailedRepository orderFailedRepo;
@RabbitListener(queues = RabbitMQConfig.DLX_QUEUE)
public void handleDeadLetter(Message message, Channel channel) throws Exception {
long deliveryTag = message.getMessageProperties().getDeliveryTag();
try {
OrderDTO order = (OrderDTO) message.getMessageProperties()
.getConverter()
.fromMessage(message);
// 持久化失败订单到数据库,供人工排查
OrderFailedRecord record = new OrderFailedRecord();
record.setOrderId(order.getOrderId());
record.setPayload(JSON.toJSONString(order));
record.setReason("消费失败,超过最大重试次数");
record.setCreateTime(LocalDateTime.now());
orderFailedRepo.save(record);
// 发送告警通知
alertService.sendAlert(String.format(
"订单消息进入死信队列,需人工处理。订单号: %s",
order.getOrderId()
));
channel.basicAck(deliveryTag, false);
log.warn("死信消息已处理: orderId={}", order.getOrderId());
} catch (Exception e) {
log.error("死信消息处理异常", e);
// 死信处理失败不reject,避免死信队列消息丢失
channel.basicNack(deliveryTag, false, false);
}
}
}
API接口规范与监控治理
对外提供消息发送的REST API接口:
@RestController
@RequestMapping("/api/v1/orders")
public class OrderController {
@Autowired
private OrderMessageProducer producer;
@PostMapping
public Result<String> createOrder(@RequestBody @Valid OrderDTO order) {
// 保存订单到数据库
orderService.save(order);
// 异步发送消息
producer.sendOrderMessage(order);
return Result.success(order.getOrderId());
}
}
RabbitMQ管理端监控关键指标:
# 通过RabbitMQ HTTP API获取队列状态
curl -u admin:password http://192.168.1.50:15672/api/queues/production/queue.dlx
# 关注指标:
# messages:死信队列积压消息数(应为0,持续增长说明消费端异常)
# messages_ready:待消费消息数
# consumers:消费者数量(应>=1)
# messages_unacknowledged:未确认消息数
配置Prometheus采集RabbitMQ指标,设置告警规则:死信队列消息数大于0时触发告警;业务队列消息积压超过5000条触发告警;消费者连接数降为0触发紧急告警。配合Grafana面板可视化消息投递成功率、消费延迟、重试次数等指标,构建完整的消息中间件监控体系。
原创文章,作者:小编,如若转载,请注明出处:https://www.yunthe.com/springboot-ji-cheng-rabbitmq-si-xin-dui-lie-xiao-xi-ke-kao/