当前位置:   article > 正文

消息队列-rabbitmq(生产者.消费者. 消息.可靠性)_rabbitmq 实现可靠任务队列

rabbitmq 实现可靠任务队列

生产者者的可靠性

为了保证我们生产者在发送消息的时候消息不丢失,我们需要保证发送者的可靠性  

1.生产者重试

假如发送消息的时候消息丢失 ,我们可以使用发送者 重试机制,尝试重新发送消息

实现该机制非常简单,只需要在yml文件中 配置 rabbitmq的配置信息就可以了

2.生产者确认机制

在我们 生产者发送消息到交换机的时候,假如 我们发送到交换机 ,但是 队列没有收到消息,会返回ack,发送到交换机,然后发送到队列,消费者没有接收到消息返回ack,但是发送到交换机失败,会返回nack

也就是说 我们 只要消息成功发送到  mq里面去,无论消息是否成功路由被消费,或者没被消费,他都会 返回ack 

下面我们看具体实现代码

生产者代码

  1. @Test
  2. public void publisherConfirm(){
  3. //获得 相关组件
  4. CorrelationData correlationData = new CorrelationData();
  5. correlationData.getFuture().addCallback(new ListenableFutureCallback<CorrelationData.Confirm>() {
  6. @Override
  7. public void onFailure(Throwable ex) {
  8. log.error("消息发送错误{}",ex.getMessage());
  9. }
  10. @Override
  11. public void onSuccess(CorrelationData.Confirm result) {
  12. if(result.isAck()){
  13. log.info("发送成功收到{},但是消息可能未路由",result.getReason());
  14. }else{
  15. log.error("发送失败收到{}",result.getReason());
  16. }
  17. }
  18. });
  19. //我们像一个不的交换机中发送一条消息 验证我们的情况
  20. rabbitTemplate.convertAndSend("myExchange.direct","my","aaa",correlationData);
  21. }

消费者代码

  1. @Component
  2. @Slf4j
  3. @RequiredArgsConstructor
  4. public class Consumer2 {
  5. private final RabbitTemplate rabbitTemplate;
  6. @RabbitListener(bindings = @QueueBinding(
  7. value = @Queue(value = "mySimple.queue",durable = "true",
  8. arguments = @Argument(name="x-queue-mode",value="lazy")),
  9. exchange = @Exchange(value = "myExchange.direct",durable = "true"),
  10. key = "my"
  11. ))
  12. public void consumerMessage(String message){
  13. log.info("收到消息{}",message);
  14. }
  15. }

mq消息可靠性

1.交换机 队列持久化

我们在设置的时候可以设置交换机队列持久化,这样即使我们rabbitmq重启的话,也不用担心交换机队列丢失

2.消息持久化

我们在设置的时候可以手动设置 一个持久化的 message 对象,当作消息传递,然后 持久化之后消息会直接存储到磁盘中去,防止内存中的消息积压同时也可以为消息设置优先级

3.lazyqueue

我们使用lazyqueue 形式发送消息到队列的时候,消息会直接存储到磁盘中去,一直到被需要消费的时候才会被加载

设置lazyqueue也很简单,我们只需要在代码中简单配置即可

消费者可靠性

1.消费者确认机制

消费者确认消息之后会返回ack,并且删除消息,

消费者对消息处理失败返回nack,mq需要重新发送消息

reject 消费者处理消息失败,并且拒绝该消息,并且删除消息

而怎么确认消息也有三种机制

  • none:不处理。即消息投递给消费者后立刻ack,消息会立刻从MQ删除。非常不安全,不建议使用

  • manual:手动模式。需要自己在业务代码中调用api,发送ackreject,存在业务入侵,但更灵活

  • auto:自动模式。SpringAMQP利用AOP对我们的消息处理逻辑做了环绕增强,当业务正常执行时则自动返回ack. 当业务出现异常时,根据异常判断返回不同结果:

    • 如果是业务异常,会自动返回nack

    • 如果是消息处理或校验异常,自动返回reject

我们在yml文件中配置即可

2.消息重试策略

当一条消息发送失败的时候,消费者重新尝试消费消息,达到我们重试的次数之后,消费者返回reject,mq直接删除消息

3.处理失败消息

当消息彻底处理失败,我们消费者设置一个新的交换机队列 重新存储这些失败的消息后续再做处理

  1. package com.itheima.consumer.config;
  2. import org.springframework.amqp.core.Binding;
  3. import org.springframework.amqp.core.BindingBuilder;
  4. import org.springframework.amqp.core.DirectExchange;
  5. import org.springframework.amqp.core.Queue;
  6. import org.springframework.amqp.rabbit.core.RabbitTemplate;
  7. import org.springframework.amqp.rabbit.retry.MessageRecoverer;
  8. import org.springframework.amqp.rabbit.retry.RepublishMessageRecoverer;
  9. import org.springframework.context.annotation.Bean;
  10. @Configuration
  11. @ConditionalOnProperty(name = "spring.rabbitmq.listener.simple.retry.enabled", havingValue = "true")
  12. public class ErrorMessageConfig {
  13. @Bean
  14. public DirectExchange errorMessageExchange(){
  15. return new DirectExchange("error.direct");
  16. }
  17. @Bean
  18. public Queue errorQueue(){
  19. return new Queue("error.queue", true);
  20. }
  21. @Bean
  22. public Binding errorBinding(Queue errorQueue, DirectExchange errorMessageExchange){
  23. return BindingBuilder.bind(errorQueue).to(errorMessageExchange).with("error");
  24. }
  25. @Bean
  26. public MessageRecoverer republishMessageRecoverer(RabbitTemplate rabbitTemplate){
  27. return new RepublishMessageRecoverer(rabbitTemplate, "error.direct", "error");
  28. }
  29. }

死信交换机和延迟消息

死信交换机 ,都是假如一个定时消息过期了,或者发送延迟消息我们直接把该消息传递到我们绑定的死信交换机中,跟上文 消息发送失败了返回rejct之后,消息发送到err交换机是两种不同的策略

我们用死信交换机发送延迟消息的策略可以如下图 

而代码也很简单,只需要再声明队列a的时候绑定死信交换机就可以了

dlx就是死信交换机

声明:本文内容由网友自发贡献,不代表【wpsshop博客】立场,版权归原作者所有,本站不承担相应法律责任。如您发现有侵权的内容,请联系我们。转载请注明出处:https://www.wpsshop.cn/w/酷酷是懒虫/article/detail/911095
推荐阅读
相关标签
  

闽ICP备14008679号