Spring Boot整合RabbitMQ如何配置消费失败消息无限间隔重发
RabbitMQ异常消息固定1秒间隔无限重试实现方案
方案一:基于Spring Retry实现(应用侧重试)
直接自定义Spring Retry的重试策略和退避策略,无需依赖RabbitMQ额外特性,配置简单。
配置步骤
- 定义无限重试拦截器
import org.springframework.retry.policy.AlwaysRetryPolicy; import org.springframework.retry.backoff.FixedBackOffPolicy; import org.springframework.retry.support.RetryTemplate; import org.springframework.amqp.rabbit.config.RetryInterceptorBuilder; import org.springframework.amqp.rabbit.config.RetryOperationsInterceptor; @Bean public RetryOperationsInterceptor infiniteRetryInterceptor() { // 配置固定1秒退避间隔 FixedBackOffPolicy fixedBackOff = new FixedBackOffPolicy(); fixedBackOff.setBackOffPeriod(1000); // 配置永远重试策略 RetryTemplate retryTemplate = new RetryTemplate(); retryTemplate.setBackOffPolicy(fixedBackOff); retryTemplate.setRetryPolicy(new AlwaysRetryPolicy()); return RetryInterceptorBuilder.stateless() .retryOperations(retryTemplate) .build(); }
- 将拦截器绑定到Rabbit监听容器工厂
@Bean public SimpleRabbitListenerContainerFactory rabbitListenerContainerFactory (ConnectionFactory connectionFactory) { SimpleRabbitListenerContainerFactory factory = new SimpleRabbitListenerContainerFactory(); factory.setConnectionFactory(connectionFactory); // 绑定重试拦截器 factory.setAdviceChain(infiniteRetryInterceptor()); return factory; }
注意事项
- 消费逻辑需要保证幂等,避免重复执行业务逻辑产生数据异常
- 如果存在消息格式错误、参数非法等永远无法处理成功的异常,可自定义重试策略过滤这类异常,避免无意义的无限重试:
RetryPolicy infiniteRetryPolicy = new AlwaysRetryPolicy() { @Override public boolean canRetry(RetryContext context) { Throwable e = context.getLastThrowable(); // 仅对业务异常、数据库异常等可恢复异常重试 return e instanceof BusinessException || e instanceof SQLException; } };
方案二:基于RabbitMQ死信队列+延迟队列实现(MQ侧重试)
重试逻辑由RabbitMQ本身实现,不占用应用线程资源,应用重启不影响重试状态。
配置步骤
- 定义队列、死信交换机和绑定关系
public static final String RESEND_DISPOSAL_QUEUE = "RESEND_DISPOSAL"; public static final String RESEND_DISPOSAL_DLX = "RESEND_DISPOSAL_DLX"; public static final String DELAY_QUEUE = "RESEND_DISPOSAL_DELAY_QUEUE"; // 业务队列,配置死信交换机 @Bean public Queue resendDisposalQueue() { Map<String, Object> args = new HashMap<>(); args.put("x-dead-letter-exchange", RESEND_DISPOSAL_DLX); args.put("x-dead-letter-routing-key", RESEND_DISPOSAL_QUEUE); return new Queue(RESEND_DISPOSAL_QUEUE, true, false, false, args); } // 死信交换机 @Bean public DirectExchange dlxExchange() { return new DirectExchange(RESEND_DISPOSAL_DLX); } // 延迟队列,消息停留1秒后自动投递回业务队列 @Bean public Queue delayQueue() { Map<String, Object> args = new HashMap<>(); args.put("x-message-ttl", 1000); args.put("x-dead-letter-exchange", ""); // 使用RabbitMQ默认交换机 args.put("x-dead-letter-routing-key", RESEND_DISPOSAL_QUEUE); return new Queue(DELAY_QUEUE, true, false, false, args); } // 绑定延迟队列到死信交换机 @Bean public Binding delayBinding() { return BindingBuilder.bind(delayQueue()) .to(dlxExchange()) .with(RESEND_DISPOSAL_QUEUE); }
- 修改监听容器配置,关闭消费失败直接重入队
@Bean public SimpleRabbitListenerContainerFactory rabbitListenerContainerFactory (ConnectionFactory connectionFactory) { SimpleRabbitListenerContainerFactory factory = new SimpleRabbitListenerContainerFactory(); factory.setConnectionFactory(connectionFactory); // 消费失败的消息进入死信队列,不直接重新入队 factory.setDefaultRequeueRejected(false); return factory; }
内容的提问来源于stack exchange,提问作者Alex Zhulin
相关产品推荐
相关产品推荐

