RabbitMQ在业务系统中承担异步解耦、削峰填谷的职责。当消息无法被正常消费或需要延迟处理时,死信队列(DLX)和延迟队列是两个核心解决方案。本文以订单超时取消、消息重试和延迟任务调度三个业务场景为例,讲解RabbitMQ DLX配置、TTL延迟队列和插件延迟队列的完整实现。
死信队列基础概念与交换器配置
死信(Dead Letter)指被消费者拒绝(nack/reject且不重新入队)、消息TTL过期或队列达到最大长度时,RabbitMQ将这些消息转发到绑定的死信交换器。死信队列本质是普通队列绑定了DLX参数。
通过Spring AMQP声明队列时配置死信路由:
@Configuration
public class RabbitMQConfig {
// 业务交换器
@Bean
public DirectExchange orderExchange() {
return ExchangeBuilder.directExchange("order.exchange")
.durable(true).build();
}
// 死信交换器
@Bean
public DirectExchange orderDlxExchange() {
return ExchangeBuilder.directExchange("order.dlx.exchange")
.durable(true).build();
}
// 业务队列(绑定死信交换器)
@Bean
public Queue orderQueue() {
return QueueBuilder.durable("order.queue")
.withArgument("x-dead-letter-exchange", "order.dlx.exchange")
.withArgument("x-dead-letter-routing-key", "order.dlx")
.withArgument("x-max-priority", 10) // 支持优先级队列
.build();
}
// 死信队列
@Bean
public Queue orderDlxQueue() {
return QueueBuilder.durable("order.dlx.queue").build();
}
// 绑定关系
@Bean
public Binding orderBinding() {
return BindingBuilder.bind(orderQueue())
.to(orderExchange())
.with("order.create");
}
@Bean
public Binding orderDlxBinding() {
return BindingBuilder.bind(orderDlxQueue())
.to(orderDlxExchange())
.with("order.dlx");
}
}
订单超时取消:TTL + DLX延迟方案
电商场景中,用户下单后30分钟未支付需自动取消。利用消息TTL过期触发死信路由实现延迟消费:
@Service
public class OrderService {
@Autowired
private RabbitTemplate rabbitTemplate;
// 发送订单创建消息,设置30分钟TTL
public void sendOrderMessage(String orderId) {
MessagePostProcessor postProcessor = message -> {
message.getMessageProperties().setExpiration("1800000"); // 30分钟=1800000ms
return message;
};
rabbitTemplate.convertAndSend(
"order.exchange",
"order.create",
orderId,
postProcessor
);
}
// 死信队列消费者:处理超时订单
@RabbitListener(queues = "order.dlx.queue")
public void handleTimeoutOrder(String orderId, Message message, Channel channel)
throws IOException {
try {
// 检查订单状态,未支付则取消
Order order = orderMapper.selectById(orderId);
if (order != null && order.getStatus() == OrderStatus.PENDING) {
order.setStatus(OrderStatus.CANCELLED);
order.setCancelReason("超时未支付自动取消");
orderMapper.updateById(order);
// 释放库存
inventoryService.releaseStock(order.getItems());
}
channel.basicAck(
message.getMessageProperties().getDeliveryTag(), false);
} catch (Exception e) {
// 拒绝消息,不再重新入队(避免死循环)
channel.basicReject(
message.getMessageProperties().getDeliveryTag(), false);
log.error("处理超时订单失败: orderId={}", orderId, e);
}
}
}
TTL方案的局限:RabbitMQ对队列中消息的TTL检查是从队头开始的。如果队头消息TTL尚未过期,后面的消息即使TTL已过也不会被移除,导致延迟时间不准确。队列级TTL可解决此问题,但会对所有消息统一设置相同的过期时间。
rabbitmq-delayed-message-exchange插件方案
对于延迟时间不固定的场景,安装官方延迟消息插件实现精确的每消息延迟:
# 安装延迟消息插件
wget https://github.com/rabbitmq/rabbitmq-delayed-message-exchange/releases/download/3.12.0/rabbitmq_delayed_message_exchange-3.12.0.ez
cp rabbitmq_delayed_message_exchange-3.12.0.ez /usr/lib/rabbitmq/lib/rabbitmq_server-3.12.0/plugins/
rabbitmq-plugins enable rabbitmq_delayed_message_exchange
systemctl restart rabbitmq-server
使用延迟交换器发送消息,精确控制每条消息的延迟时间:
@Configuration
public class DelayedQueueConfig {
// 自定义延迟交换器(x-delayed-message类型)
@Bean
public CustomExchange delayedExchange() {
Map<String, Object> args = new HashMap<>();
args.put("x-delayed-type", "direct"); // 底层路由模式
return new CustomExchange(
"delayed.exchange",
"x-delayed-message",
true, false, args);
}
@Bean
public Queue delayedQueue() {
return QueueBuilder.durable("delayed.queue").build();
}
@Bean
public Binding delayedBinding() {
return BindingBuilder.bind(delayedQueue())
.to(delayedExchange())
.with("delayed.routing.key")
.noargs();
}
}
@Service
public class DelayedMessageService {
@Autowired
private RabbitTemplate rabbitTemplate;
// 发送延迟消息,延迟时间由参数控制
public void sendDelayedMessage(String content, long delayMillis) {
rabbitTemplate.convertAndSend(
"delayed.exchange",
"delayed.routing.key",
content,
message -> {
message.getMessageProperties().setHeader("x-delay", delayMillis);
return message;
}
);
}
@RabbitListener(queues = "delayed.queue")
public void handleDelayedMessage(String content) {
log.info("收到延迟消息: {}, 时间: {}", content, LocalDateTime.now());
}
}
插件方案通过在消息头设置x-delay参数实现毫秒级精确延迟,消息存储在Mnesia表中而非普通队列,不受FIFO顺序限制。
消息重试与指数退避策略
消费失败时需要自动重试,但直接nack重新入队会导致快速重试风暴。通过DLX + TTL实现指数退避重试:
@Configuration
public class RetryQueueConfig {
// 三级重试队列,TTL分别为10s、30s、60s
@Bean
public Queue retryQueue1() {
return QueueBuilder.durable("retry.queue.1")
.withArgument("x-dead-letter-exchange", "order.exchange")
.withArgument("x-dead-letter-routing-key", "order.create")
.withArgument("x-message-ttl", 10000)
.build();
}
@Bean
public Queue retryQueue2() {
return QueueBuilder.durable("retry.queue.2")
.withArgument("x-dead-letter-exchange", "order.exchange")
.withArgument("x-dead-letter-routing-key", "order.create")
.withArgument("x-message-ttl", 30000)
.build();
}
}
@Component
public class OrderConsumer {
@Autowired
private RabbitTemplate rabbitTemplate;
@Value("${rabbitmq.max-retry-count:3}")
private int maxRetryCount;
@RabbitListener(queues = "order.queue")
public void process(OrderMessage msg, Message message, Channel channel)
throws IOException {
long deliveryTag = message.getMessageProperties().getDeliveryTag();
// 获取已重试次数
Integer retryCount = (Integer) message.getMessageProperties()
.getHeader("x-retry-count");
retryCount = retryCount == null ? 0 : retryCount;
try {
// 业务处理
doProcess(msg);
channel.basicAck(deliveryTag, false);
} catch (Exception e) {
if (retryCount >= maxRetryCount) {
// 超过重试上限,路由到最终死信队列人工处理
channel.basicReject(deliveryTag, false);
log.error("消息重试{}次后仍失败,转入死信队列: {}", retryCount, msg, e);
} else {
// 路由到对应级别的重试队列
String retryQueue = retryCount == 0 ? "retry.queue.1" : "retry.queue.2";
rabbitTemplate.convertAndSend("", retryQueue, msg, m -> {
m.getMessageProperties().setHeader("x-retry-count", retryCount + 1);
return m;
});
channel.basicAck(deliveryTag, false);
log.warn("处理失败,第{}次重试,消息转入{}", retryCount + 1, retryQueue, e);
}
}
}
}
消息可靠性与监控运维
确保消息不丢失的三个环节配置:生产端、broker端和消费端。
# 生产端:开启发布确认
spring.rabbitmq.publisher-confirm-type=correlated
spring.rabbitmq.publisher-returns=true
# broker端:队列和消息持久化(durable=true + persistent消息)
# 消费端:手动ACK
spring.rabbitmq.listener.simple.acknowledge-mode=manual
通过RabbitMQ Management API监控队列积压和消费者状态:
# 查询队列深度和消费者数量
curl -u admin:password http://localhost:15672/api/queues/%2F/order.queue | \
python -m json.tool | grep -E '"messages"|"consumers"'
# 关键指标告警阈值:
# messages_ready > 1000 队列积压告警
# consumers == 0 消费者离线告警
# messages_unacknowledged > 500 未确认消息过多告警
死信队列中的消息需要人工介入或定期清理,建议配置队列最大长度x-max-length防止单队列消息堆积溢出。
原创文章,作者:小编,如若转载,请注明出处:https://www.yunthe.com/rabbitmq-si-xin-dui-lie-yu-yan-chi-xiao-xi-jia-gou-she-ji/