欢迎光临
我们一直在努力

RabbitMQ 常用模式:本地重试与死信队列

RabbitMQ 常用模式:本地重试与死信队列

推荐方案: AUTO 确认 + Spring 本地有限重试 + RejectAndDontRequeueRecoverer + RabbitMQ DLX。

方案概览

本方案适合秒级、少次数、可幂等的消费失败重试。各组件职责如下:

组件职责
Spring Listener Retry 在当前消费者线程内执行有限次数的本地重试
RejectAndDontRequeueRecoverer 重试耗尽后拒绝消息,并设置 requeue=false
RabbitMQ DLX 将被拒绝的消息路由到死信交换机
DLQ 保存最终处理失败的消息,供告警、排查或人工重放

[!NOTE] Spring AMQP 的 acknowledge-mode: auto 表示由监听容器根据方法是否正常返回来发送 ACK/NACK,不等同于 RabbitMQ 的 autoAck=true;后者在 Spring AMQP 中对应 AcknowledgeMode.NONE。

执行链路如下:

监听方法抛异常

Spring 在当前消费线程内重试

超过最大次数

RejectAndDontRequeueRecoverer
↓ requeue=false
RabbitMQ 将消息投递到 DLX

DLQ 保存失败消息

1. 配置消费者重试

Spring Boot 4.x 使用 max-retries:

spring:
rabbitmq:
listener:
simple:
acknowledge-mode: auto
prefetch: 20
concurrency: 3
max-concurrency: 10
default-requeue-rejected: false
retry:
enabled: true
max-retries: 3
initial-interval: 1s
multiplier: 2
max-interval: 10s

这个例子的执行过程大致是:

首次消费失败
↓ 等待1秒
第1次重试失败
↓ 等待2秒
第2次重试失败
↓ 等待4秒
第3次重试失败

进入恢复逻辑,最终投递 DLQ

如果是 Spring Boot 2.x/3.x,通常使用旧属性:

retry:
enabled: true
max-attempts: 4
initial-interval: 1s
multiplier: 2
max-interval: 10s

max-attempts: 4 包含第一次执行,也就是“首次执行 + 3 次重试”。当前 Spring Boot 配置已经使用 max-retries。参见 Spring Boot RabbitMQ 配置。

注意不要配置错位置:

# 消费者监听重试
spring.rabbitmq.listener.simple.retry

# 生产者发送重试
spring.rabbitmq.template.retry

2. 声明主队列、DLX 和 DLQ

import org.springframework.amqp.core.Binding;
import org.springframework.amqp.core.BindingBuilder;
import org.springframework.amqp.core.DirectExchange;
import org.springframework.amqp.core.Queue;
import org.springframework.amqp.core.QueueBuilder;
import org.springframework.amqp.rabbit.retry.MessageRecoverer;
import org.springframework.amqp.rabbit.retry.RejectAndDontRequeueRecoverer;
import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;

/**
* 功能:
* <p>
* 声明订单消息的主队列和死信拓扑。业务消息正常进入主队列;
* 消费重试耗尽后,消息被拒绝且不重新入队,再由RabbitMQ转发到死信队列。
* </p>
* <p>
* 该配置只负责消息基础设施,不处理订单业务逻辑。
* </p>
*/

@Configuration(proxyBeanMethods = false)
public class OrderRabbitTopologyConfig {

public static final String ORDER_EXCHANGE =
"order.exchange";

public static final String ORDER_ROUTING_KEY =
"order.created";

public static final String ORDER_QUEUE =
"order.created.queue";

public static final String ORDER_DEAD_EXCHANGE =
"order.dead.exchange";

public static final String ORDER_DEAD_ROUTING_KEY =
"order.created.dead";

public static final String ORDER_DEAD_QUEUE =
"order.created.dlq";

@Bean
DirectExchange orderExchange() {
return new DirectExchange(
ORDER_EXCHANGE,
true,
false
);
}

@Bean
Queue orderQueue() {
return QueueBuilder.durable(ORDER_QUEUE)
.deadLetterExchange(ORDER_DEAD_EXCHANGE)
.deadLetterRoutingKey(ORDER_DEAD_ROUTING_KEY)
.build();
}

@Bean
Binding orderBinding(
@Qualifier("orderQueue") Queue queue,
@Qualifier("orderExchange") DirectExchange exchange) {

return BindingBuilder.bind(queue)
.to(exchange)
.with(ORDER_ROUTING_KEY);
}

@Bean
DirectExchange orderDeadExchange() {
return new DirectExchange(
ORDER_DEAD_EXCHANGE,
true,
false
);
}

@Bean
Queue orderDeadQueue() {
return QueueBuilder.durable(ORDER_DEAD_QUEUE)
.build();
}

@Bean
Binding orderDeadBinding(
@Qualifier("orderDeadQueue") Queue queue,
@Qualifier("orderDeadExchange") DirectExchange exchange) {

return BindingBuilder.bind(queue)
.to(exchange)
.with(ORDER_DEAD_ROUTING_KEY);
}

/**
* 功能:
* <p>
* 当监听器的有限重试全部失败后,要求容器拒绝消息且不重新进入原队列。
* 源队列配置DLX后,RabbitMQ会将该消息转发到订单死信队列。
* </p>
*
* @return 消费重试耗尽后的恢复策略
*/

@Bean
MessageRecoverer orderMessageRecoverer() {
return new RejectAndDontRequeueRecoverer(
"订单消息重试耗尽,转入死信队列"
);
}
}

Spring Boot 在开启监听器重试后,默认也会使用 RejectAndDontRequeueRecoverer;这里显式声明,是为了让失败语义更清楚。重试耗尽后消息会被拒绝,配置了 DLX 就进入死信队列,否则会被丢弃。参见 Spring Boot AMQP 文档。

3. 编写消费者

import org.springframework.amqp.rabbit.annotation.RabbitListener;
import org.springframework.stereotype.Component;

/**
* 功能:
* <p>
* 消费订单创建消息,并把具体业务处理交给订单消息服务。
* 监听方法正常返回时由Spring自动ACK;处理失败时必须继续抛出异常,
* 由监听器重试和死信机制统一处理。
* </p>
*/

@Component
public class OrderCreatedMessageListener {

private final OrderMessageService orderMessageService;

public OrderCreatedMessageListener(
OrderMessageService orderMessageService) {
this.orderMessageService = orderMessageService;
}

/**
* 功能:
* <p>
* 消费订单创建事件。业务服务执行完成后方法正常返回,Spring发送ACK;
* 业务服务抛出异常时,本方法不捕获,交给Spring执行有限重试。
* </p>
*
* @param message 订单创建消息,messageId用于消费幂等
*/

@RabbitListener(
queues = OrderRabbitTopologyConfig.ORDER_QUEUE
)
public void consume(OrderCreatedMessage message) {
orderMessageService.consume(message);
}
}

这里最重要的是:不要捕获异常后只打印日志。

错误写法:

@RabbitListener(queues = "order.created.queue")
public void consume(OrderCreatedMessage message) {
try {
orderMessageService.consume(message);
} catch (Exception exception) {
log.error("订单消息处理失败", exception);
}
}

异常被吞掉后,监听方法正常返回:

Spring 认为处理成功

发送 ACK

不会重试,也不会进入 DLQ

如果确实需要记录日志,必须继续抛出:

@RabbitListener(queues = "order.created.queue")
public void consume(OrderCreatedMessage message) {
try {
orderMessageService.consume(message);
} catch (Exception exception) {
log.error(
"订单消息处理失败,messageId={}",
message.getMessageId(),
exception
);
throw exception;
}
}

4. 业务服务必须实现幂等

import org.springframework.stereotype.Service;
import org.springframework.transaction.annotation.Transactional;

/**
* 功能:
* <p>
* 在数据库事务内完成订单消息幂等校验和业务处理。
* 消息重试或ACK丢失时可能发生重复投递,因此通过messageId唯一记录
* 确保相同消息不会重复修改订单。
* </p>
*/

@Service
public class OrderMessageService {

private final ConsumeLogRepository consumeLogRepository;
private final OrderRepository orderRepository;

public OrderMessageService(
ConsumeLogRepository consumeLogRepository,
OrderRepository orderRepository) {
this.consumeLogRepository = consumeLogRepository;
this.orderRepository = orderRepository;
}

/**
* 功能:
* <p>
* 消费订单事件。首次消费时写入幂等记录并更新订单;
* 重复消息直接返回,使监听器可以安全ACK。
* </p>
*
* @param message 待处理的订单创建消息
*/

@Transactional(rollbackFor = Exception.class)
public void consume(OrderCreatedMessage message) {
boolean isFirstConsumption =
consumeLogRepository.tryInsert(
message.getMessageId()
);

if (!isFirstConsumption) {
return;
}

orderRepository.createOrder(
message.getOrderId(),
message.getUserId()
);
}
}

message_id 必须有数据库唯一索引:

CREATE UNIQUE INDEX uk_mq_consume_log_message_id
ON mq_consume_log(message_id);

5. 消息什么时候进入 DLQ

配置成功后:

监听器正常返回
→ Spring ACK
→ 消息删除

监听器抛异常,但重试成功
→ Spring ACK
→ 消息删除

监听器一直失败
→ 重试耗尽
→ basic.reject / basic.nack,requeue=false
→ order.dead.exchange
→ order.created.dlq

RabbitMQ 会在死信消息的 Header 中增加 x-death,记录原队列、死信原因和次数。参见 RabbitMQ DLX 文档。

注意:这里的 Spring 重试是消费者进程内重试:

  • 重试期间消息保持 Unacked;
  • 重试会占用消费者线程;
  • 每次本地重试不会增加 x-death;
  • 只有最终被 RabbitMQ 死信转发时才产生 x-death 记录。

所以这种方案适合秒级、次数较少的重试。分钟级、小时级重试应改成:

主队列
↓ 失败
短延时重试队列
↓ 再失败
长延时重试队列
↓ 超过上限
最终 DLQ

6. 一个容易遇到的部署问题

如果 order.created.queue 已经存在,并且以前没有配置 DLX,再用上面的代码声明,可能报:

PRECONDITION_FAILED – inequivalent arg
'x-dead-letter-exchange'

因为 RabbitMQ 不允许直接修改已有队列的声明参数。

处理方式:

  • 测试环境:删除旧队列后重新声明;
  • 生产环境:优先通过 RabbitMQ Policy 设置 DLX;
  • 不要直接删除仍有消息的生产队列。

RabbitMQ 官方也更推荐使用 Policy 配置 DLX,因为 Policy 可以动态调整,而硬编码的 x-arguments 通常需要重新创建队列。参见 RabbitMQ DLX Policy。

7. 验收测试

至少验证以下场景:

/**
* 业务处理成功时,验证监听器只执行一次且消息不会进入DLQ。
*/

@Test
void shouldAcknowledgeMessageWhenBusinessSucceeds() {
}

/**
* 业务持续失败时,验证达到重试上限后消息进入订单DLQ。
*/

@Test
void shouldMoveMessageToDeadQueueAfterRetriesExhausted() {
}

/**
* 同一个messageId被重复投递时,验证订单只创建一次。
*/

@Test
void shouldKeepBusinessIdempotentWhenMessageIsRedelivered() {
}

集成测试可以从 DLQ 读取结果:

Message deadMessage = rabbitTemplate.receive(
OrderRabbitTopologyConfig.ORDER_DEAD_QUEUE,
15_000
);

assertThat(deadMessage).isNotNull();
assertThat(deadMessage.getMessageProperties().getHeaders())
.containsKey("x-death");

8. 运维与重放建议

  • 为 DLQ 的消息数量、最老消息滞留时间和持续增长趋势配置监控告警。
  • 重放前先定位失败原因并修复消费者,避免消息重新进入“主队列 → 重试 → DLQ”的循环。
  • 重放工具应保留原始 messageId,继续复用消费端幂等校验;同时记录操作人、重放时间、批次和结果。
  • DLQ 是失败消息的隔离区,不等同于自动补偿机制;是否自动重放应根据异常类型、业务风险和重试间隔单独设计。

参考资料

  • Spring Boot RabbitMQ 配置属性
  • Spring Boot AMQP 文档
  • Spring AMQP 监听容器配置
  • RabbitMQ Dead Letter Exchanges
  • RabbitMQ Reliability Guide
赞(0)
未经允许不得转载:171主机测评 » RabbitMQ 常用模式:本地重试与死信队列
分享到: 更多 (0)

评论 抢沙发

  • 昵称 (必填)
  • 邮箱 (必填)
  • 网址