死信队列:消息消费失败的安全网
RabbitMQ中消息被拒绝(basic.reject/basic.nack且requeue=false)、消息TTL过期、队列达到最大长度时,消息变为”死信”。默认情况下死信被丢弃,这对生产系统不可接受——订单支付超时、库存扣减失败这类消息丢一条就是业务事故。
死信队列(DLX)的原理:给队列绑定一个exchange,消息变成死信后自动路由到DLX,再由DLX投递到死信存储队列。运维可以事后排查、重试或补偿。
死信队列配置实战
用Spring Boot + RabbitMQ的Java配置为例:
@Configuration
public class RabbitDLXConfig {
// ===== 死信exchange和队列 =====
@Bean
public DirectExchange deadLetterExchange() {
return new DirectExchange("order.dlx");
}
@Bean
public Queue deadLetterQueue() {
return QueueBuilder.durable("order.dlq").build();
}
@Bean
public Binding dlqBinding() {
return BindingBuilder.bind(deadLetterQueue())
.to(deadLetterExchange())
.with("order.dead");
}
// ===== 业务队列(绑定DLX) =====
@Bean
public Queue orderQueue() {
return QueueBuilder.durable("order.queue")
.withArgument("x-dead-letter-exchange", "order.dlx")
.withArgument("x-dead-letter-routing-key", "order.dead")
.build();
}
@Bean
public DirectExchange orderExchange() {
return new DirectExchange("order.exchange");
}
@Bean
public Binding orderBinding() {
return BindingBuilder.bind(orderQueue())
.to(orderExchange())
.with("order.create");
}
}
x-dead-letter-exchange和x-dead-letter-routing-key是队列级别的参数,声明时指定。消息变成死信后,RabbitMQ自动将消息投递到order.dlx,路由键为order.dead,最终进入order.dlq。
消费者异常处理与重试策略
消费者代码中,不要吞掉异常,也不要无条件requeue:
@Component
@RabbitListener(queues = "order.queue")
public class OrderConsumer {
@RabbitHandler
public void handle(Message message, Channel channel,
@Header(AmqpHeaders.DELIVERY_TAG) long tag) throws IOException {
try {
// 业务处理
processOrder(new String(message.getBody()));
channel.basicAck(tag, false);
} catch (BusinessException e) {
// 业务异常,拒绝进死信
log.error("订单处理失败: {}", e.getMessage());
channel.basicNack(tag, false, false);
} catch (Exception e) {
// 未知异常,requeue重试一次
log.warn("订单处理异常,重试", e);
channel.basicNack(tag, false, true);
}
}
}
这里有个问题:requeue=true会把消息重新放回队列头部,导致无限重试阻塞后续消息。生产环境推荐用Spring Retry做有限次数重试:
@Configuration
public class RetryConfig {
@Bean
public RetryOperationsInterceptor retryInterceptor() {
return RetryInterceptorBuilder.stateless()
.maxAttempts(3)
.backOffOptions(1000, 2.0, 10000) // 初始1s,倍数2,最大10s
.recoverer(new RejectRecoveryer())
.build();
}
}
backOffOptions设置指数退避:第1次1秒后重试,第2次2秒,第3次4秒,避免雪崩式重试压垮下游。
延时消息:TTL + DLX方案
业务场景:订单30分钟未支付自动取消。常见错误做法是轮询数据库,表数据量大时数据库和API接口双双扛不住。RabbitMQ的延时消息方案是:消息设置TTL,过期后进入DLX,DLX路由到实际消费队列。
@Configuration
public class DelayMessageConfig {
@Bean
public DirectExchange delayExchange() {
return new DirectExchange("delay.exchange");
}
@Bean
public Queue delayQueue() {
return QueueBuilder.durable("delay.queue")
.withArgument("x-dead-letter-exchange", "order.exchange")
.withArgument("x-dead-letter-routing-key", "order.timeout")
.build();
}
@Bean
public Binding delayBinding() {
return BindingBuilder.bind(delayQueue())
.to(delayExchange())
.with("order.delay");
}
}
发送延时消息时设置per-message TTL:
@Service
public class OrderService {
@Autowired
private RabbitTemplate rabbitTemplate;
public void createOrder(Order order) {
// 1. 创建订单
orderMapper.insert(order);
// 2. 发送延时消息(30分钟TTL)
rabbitTemplate.convertAndSend("delay.exchange", "order.delay",
order.getId(), message -> {
message.getMessageProperties().setExpiration("1800000");
return message;
});
}
}
30分钟后消息过期 → 进入order.exchange → 路由到order.timeout队列 → 消费者检查订单状态 → 未支付则取消。
TTL+DLX方案的陷阱
陷阱1:消息阻塞
RabbitMQ判断队列头部消息是否过期。如果队列中有TTL=5分钟的消息A在前面、TTL=1分钟的消息B在后面,B必须等A过期后才能被投递。per-message TTL不同时,这个阻塞问题很严重。
解法:为每种TTL创建独立的延时队列,如delay.queue.5m、delay.queue.30m、delay.queue.1h,各自绑定DLX。
陷阱2:延时精度
RabbitMQ不是精确定时器。消息过期后到被投递到DLX之间有时间差,取决于队列空闲扫描间隔(默认1秒)。对30分钟级别的延时无所谓,对秒级延时需要考虑精度损失。
如果业务要求秒级延时精度,推荐用RabbitMQ的延迟插件(rabbitmq_delayed_message_exchange):
# 安装插件
rabbitmq-plugins enable rabbitmq_delayed_message_exchange
// 使用插件方式
@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);
}
// 发送时设置延迟
rabbitTemplate.convertAndSend("delayed.exchange", "order.delay",
orderId, message -> {
message.getMessageProperties().setDelay(60000); // 60秒
return message;
});
插件方案不会阻塞消息,精度也更高,但牺牲了部分吞吐量。
死信队列监控与告警
死信队列是最后一道防线,本身不能出问题。用RabbitMQ HTTP API监控队列深度:
@Scheduled(fixedRate = 60000)
public void monitorDLQ() {
RestTemplate rest = new RestTemplate();
String url = "http://rabbitmq:15672/api/queues/%2F/order.dlq";
ResponseEntity<Map> resp = rest.getForEntity(url, Map.class);
Map body = resp.getBody();
int messages = (int) body.get("messages");
if (messages > 0) {
alertService.send("死信队列order.dlq有" + messages + "条消息待处理");
}
}
或者用Prometheus + rabbitmq_exporter做更完整的监控:
# Prometheus告警规则
- alert: RabbitMQDeadLetterQueueGrowing
expr: rabbitmq_queue_messages{queue=~".*\\.dlq"} > 10
for: 5m
labels:
severity: warning
annotations:
summary: "死信队列 {{ $labels.queue }} 持续增长"
微服务架构中,消息中间件的可靠性直接决定业务的一致性。死信队列不是”出了问题再看”,而是日常巡检的一部分。每次死信增长都值得追查根因——是下游超时?是消息格式变更?还是消费者有bug?修复后把死信重新投递到业务队列即可恢复。
原创文章,作者:小编,如若转载,请注明出处:https://www.yunthe.com/rabbitmq-si-xin-dui-lie-yu-yan-shi-xiao-xi-fang-an-sheng/