这里主要针对SpringBoot中的RabbitMQTemplate来聊。

SpringBootRabbitMQTemplate消费者确认机制怎么选?

消费者消息确认

RabbitMQ的默认机制,说白了就是“阅后即焚”。一旦确认消息被消费者成功消费,它就会立刻从队列里删除。那它是怎么判断消费者是否成功处理了呢?关键就在于消费者回执——消费者拿到消息后,得向RabbitMQ发一个ACK回执,告诉它:“我处理完了,可以删了。”

来看这样一个场景

这会导致什么问题?消息丢了。所以,什么时候返回ACK,这个时机至关重要。

SpringAMQP提供了三种确认模式

从实际效果来看:

大多数情况下,直接用默认的auto模式就足够了。

配置方式

spring:
  rabbitmq:
    listener:
      simple:
        acknowledge-mode: 确认模式

消费者失败重试机制

当消费者处理消息时抛出异常,消息会被重新放回队列(requeue),然后再次发送给消费者,又异常,又requeue……如此循环往复,MQ的处理压力会飙升,这显然不是什么好事。

本地重试

Spring的retry机制能很好地解决这个问题:消费者出现异常时,在本地进行重试,而不是无休止地把消息丢回MQ队列。

配置方式也很简单,在consumer服务的application.yml中添加以下内容:

spring:
  rabbitmq:
    listener:
      simple:
        retry:
          enabled: true # 开启消费者失败重试
          initial-interval: 1000ms # 初始失败等待时长1秒
          multiplier: 1 # 等待时长倍数,下次等待时长 = multiplier * last-interval
          max-attempts: 3 # 最大重试次数
          stateless: true # true无状态;false有状态。如果业务包含事务,这里改为false

重启consumer服务后重新测试,你会发现:

结论很清晰:

失败策略

从上面的测试可以看出,重试次数耗尽后,消息默认会被丢弃。这是由Spring内部的MessageRecovery接口决定的,它提供了三种不同的实现:

比较优雅的做法是使用RepublishMessageRecoverer,把失败的消息投递到一个专门存放异常消息的队列,后续由人工集中处理,这样既不会丢失消息,也不会影响正常流程。

具体实现分两步:

1)在consumer服务中定义处理失败消息的交换机和队列:

@Bean("error_direct")
public DirectExchange errorMessageExchange(){
    return new DirectExchange("error_direct");
}
@Bean("error_queue")
public Queue errorQueue(){
    return new Queue("error_queue", true);
}
@Bean
public Binding bindingerror(DirectExchange error_direct, Queue error_queue){
    return BindingBuilder.bind(error_queue).to(error_direct).with("error");
}

2)定义一个RepublishMessageRecoverer,关联队列和交换机:

@Bean
public MessageRecoverer messageRecoverer(RabbitTemplate rabbitTemplate){
    return new RepublishMessageRecoverer(rabbitTemplate, "error_direct", "error");
}

完整代码整合如下:

@Bean("error_direct")
public DirectExchange errorMessageExchange(){
    return new DirectExchange("error_direct");
}
@Bean("error_queue")
public Queue errorQueue(){
    return new Queue("error_queue", true);
}
@Bean
public Binding bindingerror(DirectExchange error_direct, Queue error_queue){
    return BindingBuilder.bind(error_queue).to(error_direct).with("error");
}
@Bean
public MessageRecoverer messageRecoverer(RabbitTemplate rabbitTemplate){
    return new RepublishMessageRecoverer(rabbitTemplate, "error_direct", "error");
}

总结

以上就是关于SpringBoot中RabbitMQTemplate消费者确认和重试机制的核心内容。理解这些机制,能帮你在实际项目中更好地保障消息的可靠性,避免数据丢失或重复处理。希望这篇文章能给你一些启发。

本文转载于:https://www.jb51.net/program/368018grh.htm 如有侵犯,请联系zhengruancom@outlook.com删除。
免责声明:正软商城发布此文仅为传递信息,不代表正软商城认同其观点或证实其描述。