springboot整合rabbitmq发布确认高级

这篇具有很好参考价值的文章主要介绍了springboot整合rabbitmq发布确认高级。希望对大家有所帮助。如果存在错误或未考虑完全的地方,请大家不吝赐教,您也可以点击"举报违法"按钮提交疑问。

在生产环境中由于一些不明原因,导致 rabbitmq 重启,在 RabbitMQ 重启期间生产者消息投递失败,导致消息丢失,需要手动处理和恢复。于是,我们如何才能进行 RabbitMQ 的消息可靠投递。

发布确认  

发布确认方案

springboot整合rabbitmq发布确认高级,rabbitmq,java-rabbitmq,spring boot,rabbitmq

 架构

springboot整合rabbitmq发布确认高级,rabbitmq,java-rabbitmq,spring boot,rabbitmq 

配置文件 

在配置文件当中添加 spring.rabbitmq.publisher-confirm-type=correlated

  NONE:禁用发布确认模式,是默认值
  CORREL:ATED发布消息成功到交换器后会触发回调方法
spring.rabbitmq.host=43.139.59.23
spring.rabbitmq.port=5672
spring.rabbitmq.username=admin
spring.rabbitmq.password=123
spring.rabbitmq.publisher-confirm-type=correlated

配置类

import org.springframework.amqp.core.*;
import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;

@Configuration
public class ConfirmConfig {
   public static final String CONFIRM_EXCHANGE_NAME = "confirm.exchange";
   public static final String CONFIRM_QUEUE_NAME = "confirm.queue";
   //声明业务 Exchange
   @Bean("confirmExchange")
   public DirectExchange confirmExchange(){
   return new DirectExchange(CONFIRM_EXCHANGE_NAME);
   }
   // 声明确认队列
   @Bean("confirmQueue")
   public Queue confirmQueue(){
   return QueueBuilder.durable(CONFIRM_QUEUE_NAME).build();
   }
   // 声明确认队列绑定关系
   @Bean
   public Binding queueBinding(@Qualifier("confirmQueue") Queue queue, @Qualifier("confirmExchange") DirectExchange exchange){
    return BindingBuilder.bind(queue).to(exchange).with("key1");
   }
}

 生产者

import com.example.demo.component.MyCallBack;
import lombok.extern.slf4j.Slf4j;
import org.springframework.amqp.rabbit.connection.CorrelationData;
import org.springframework.amqp.rabbit.core.RabbitTemplate;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.PathVariable;
import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.RestController;

import javax.annotation.PostConstruct;
import javax.annotation.Resource;

@RestController
@RequestMapping("/confirm")
@Slf4j
public class Producer {
     public static final String CONFIRM_EXCHANGE_NAME = "confirm.exchange";
     @Resource
     private RabbitTemplate rabbitTemplate;
     @Resource
     private MyCallBack myCallBack;
     //依赖注入 rabbitTemplate 之后再设置它的回调对象
     @PostConstruct
     public void init(){
     rabbitTemplate.setConfirmCallback(myCallBack);
     }
     @GetMapping("sendMessage/{message}")
     public void sendMessage(@PathVariable String message){
        //指定消息 id 为 1
        CorrelationData correlationData1=new CorrelationData("1");
        String routingKey="key1";
        
       rabbitTemplate.convertAndSend(CONFIRM_EXCHANGE_NAME,routingKey,message+routingKey,correlationData1);
        CorrelationData correlationData2=new CorrelationData("2");
        routingKey="key2";
        
       rabbitTemplate.convertAndSend(CONFIRM_EXCHANGE_NAME,routingKey,message+routingKey,correlationData2);
        log.info("发送消息内容:{}",message);
     }
}

 回调接口

import lombok.extern.slf4j.Slf4j;
import org.springframework.amqp.rabbit.connection.CorrelationData;
import org.springframework.amqp.rabbit.core.RabbitTemplate;
import org.springframework.stereotype.Component;

@Component
@Slf4j
public class MyCallBack implements RabbitTemplate.ConfirmCallback {
   /**
   * 交换机不管是否收到消息的一个回调方法
   * CorrelationData
   * 消息相关数据
   * ack
   * 交换机是否收到消息
   */
   @Override
   public void confirm(CorrelationData correlationData, boolean ack, String cause) {
      String id=correlationData!=null?correlationData.getId():"";
      if(ack){
        log.info("交换机已经收到 id 为:{}的消息",id);
      }else{
       log.info("交换机还未收到 id 为:{}消息,由于原因:{}",id,cause);
      }
   }
}

 消费者

import lombok.extern.slf4j.Slf4j;
import org.springframework.amqp.core.Message;
import org.springframework.amqp.rabbit.annotation.RabbitListener;
import org.springframework.stereotype.Component;

@Component
@Slf4j
public class ConfirmConsumer {
   public static final String CONFIRM_QUEUE_NAME = "confirm.queue";
   @RabbitListener(queues =CONFIRM_QUEUE_NAME)
   public void receiveMsg(Message message){
     String msg=new String(message.getBody());
     log.info("接受到队列 confirm.queue 消息:{}",msg);
   }
}

 结果

springboot整合rabbitmq发布确认高级,rabbitmq,java-rabbitmq,spring boot,rabbitmq

可以看到,发送了两条消息,第一条消息的 RoutingKey "key1" ,第二条消息的 RoutingKey
"key2" ,两条消息都成功被交换机接收,也收到了交换机的确认回调,但消费者只收到了一条消息,因为第二条消息的 RoutingKey 与队列的 BindingKey 不一致,也没有其它队列能接收这个消息,所以第二条消息被直接丢弃了。

 回退消息

Mandatory 参数 

在仅开启了生产者确认机制的情况下,交换机接收到消息后,会直接给消息生产者发送确认消息 果发现该消息不可路由,那么消息会被直接丢弃,此时生产者是不知道消息被丢弃这个事件的 。那么如何让无法被路由的消息帮我想办法处理一下?最起码通知我一声,我好自己处理啊。通过设置 mandatory 参 数可以在当消息传递过程中不可达目的地时将消息返回给生产者。

 生产者

import com.example.demo.component.MyCallBack;
import lombok.extern.slf4j.Slf4j;
import org.springframework.amqp.rabbit.connection.CorrelationData;
import org.springframework.amqp.rabbit.core.RabbitTemplate;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.PathVariable;
import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.RestController;

import javax.annotation.PostConstruct;
import javax.annotation.Resource;
import java.util.UUID;

@RestController
@RequestMapping("/confirm")
@Slf4j
public class Producer {
     public static final String CONFIRM_EXCHANGE_NAME = "confirm.exchange";
     @Resource
     private RabbitTemplate rabbitTemplate;
     @Resource
     private MyCallBack myCallBack;
     //依赖注入 rabbitTemplate 之后再设置它的回调对象
     @PostConstruct
     public void init(){
         rabbitTemplate.setConfirmCallback(myCallBack);
         /**
          * true:
          * 交换机无法将消息进行路由时,会将该消息返回给生产者
          * false:
          * 如果发现消息无法进行路由,则直接丢弃
          */
         rabbitTemplate.setMandatory(true);
         //设置回退消息交给谁处理
         rabbitTemplate.setReturnCallback(myCallBack);

     }
    

    @GetMapping("sendMessage")
    public void sendMessage(String message){
        //让消息绑定一个 id 值
        CorrelationData correlationData1 = new CorrelationData(UUID.randomUUID().toString());

        rabbitTemplate.convertAndSend(CONFIRM_EXCHANGE_NAME,"key1",message+"key1",correlationData1)
        ;
        log.info("发送消息 id 为:{}内容为{}",correlationData1.getId(),message+"key1");
        CorrelationData correlationData2 = new CorrelationData(UUID.randomUUID().toString());

        rabbitTemplate.convertAndSend(CONFIRM_EXCHANGE_NAME,"key2",message+"key2",correlationData2)
        ;
        log.info("发送消息 id 为:{}内容为{}",correlationData2.getId(),message+"key2");
    }
}

 回调接口

import lombok.extern.slf4j.Slf4j;
import org.springframework.amqp.core.Message;
import org.springframework.amqp.rabbit.connection.CorrelationData;
import org.springframework.amqp.rabbit.core.RabbitTemplate;
import org.springframework.stereotype.Component;

@Component
@Slf4j
public class MyCallBack implements RabbitTemplate.ConfirmCallback,RabbitTemplate.ReturnCallback {
    /**
     * 交换机不管是否收到消息的一个回调方法
     * CorrelationData
     * 消息相关数据
     * ack
     * 交换机是否收到消息
     */
    @Override
    public void confirm(CorrelationData correlationData, boolean ack, String cause) {
        String id=correlationData!=null?correlationData.getId():"";
        if(ack){
            log.info("交换机已经收到 id 为:{}的消息",id);
        }else{
            log.info("交换机还未收到 id 为:{}消息,由于原因:{}",id,cause);
        }
    }
    //当消息无法路由的时候的回调方法
    @Override
    public void returnedMessage(Message message, int replyCode, String replyText, String
    exchange, String routingKey) {
        log.error(" 消 息 {}, 被交换机 {} 退回,退回原因 :{}, 路 由 key:{}",new
        String(message.getBody()),exchange,replyText,routingKey);
    }
}

 结果

springboot整合rabbitmq发布确认高级,rabbitmq,java-rabbitmq,spring boot,rabbitmq 接收到被退回的消息

备份交换机 

有了 mandatory 参数和回退消息,我们获得了对无法投递消息的感知能力,有机会在生产者的消息无法被投递时发现并处理。但有时候,我们并不知道该如何处理这些无法路由的消息,最多打个日志,然后触发报警,再来手动处理。而通过日志来处理这些无法路由的消息是很不优雅的做法,特别是当生产者所在的服务有多台机器的时候,手动复制日志会更加麻烦而且容易出错。而且设置 mandatory 参数会增加生产者的复杂性,需要添加处理这些被退回的消息的逻辑。如果既不想丢失消息,又不想增加生产者的复杂性,该怎么做呢?前面在设置死信队列的文章中,我们提到,可以为队列设置死信交换机来存储那些处理失败的消息,可是这些不可路由消息根本没有机会进入到队列,因此无法使用死信队列来保存消息。在 RabbitMQ 中,有一种备份交换机的机制存在,可以很好的应对这个问题。什么是备份交换机呢?备份交换机可以理解为 RabbitMQ 中交换机的“备胎”,当我们为某一个交换机声明一个对应的备份交换机时,就是为它创建一个备胎,当交换机接收到一条不可路由消息时,将会把这条消息转发到备份交换机中,由备份交换机来进行转发和处理,通常备份交换机的类型为 Fanout ,这样就能把所有消息都投递到与其绑定的队列中,然后我们在备份交换机下绑定一个队列,这样所有那些原交换机无法被路由的消息,就会都进入这个队列了。

架构

springboot整合rabbitmq发布确认高级,rabbitmq,java-rabbitmq,spring boot,rabbitmq

 配置类

import org.springframework.amqp.core.*;
import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;

@Configuration
public class ConfirmConfig {
     public static final String CONFIRM_EXCHANGE_NAME = "confirm.exchange";
     public static final String CONFIRM_QUEUE_NAME = "confirm.queue";
     public static final String BACKUP_EXCHANGE_NAME = "backup.exchange";
     public static final String BACKUP_QUEUE_NAME = "backup.queue";
     public static final String WARNING_QUEUE_NAME = "warning.queue";
     // 声明确认队列
     @Bean("confirmQueue")
     public Queue confirmQueue(){
       return QueueBuilder.durable(CONFIRM_QUEUE_NAME).build();
     }
     //声明确认队列绑定关系
     @Bean
     public Binding queueBinding(@Qualifier("confirmQueue") Queue queue, @Qualifier("confirmExchange") DirectExchange exchange){
       return BindingBuilder.bind(queue).to(exchange).with("key1");
     }
     //声明备份 Exchange
     @Bean("backupExchange")
     public FanoutExchange backupExchange(){
       return new FanoutExchange(BACKUP_EXCHANGE_NAME);
     }
     //声明确认 Exchange 交换机的备份交换机
     @Bean("confirmExchange")
     public DirectExchange confirmExchange(){
       ExchangeBuilder exchangeBuilder =
       ExchangeBuilder.directExchange(CONFIRM_EXCHANGE_NAME)
       .durable(true)
       //设置该交换机的备份交换机
       .withArgument("alternate-exchange", BACKUP_EXCHANGE_NAME);
       return (DirectExchange)exchangeBuilder.build();
     }
     // 声明警告队列
     @Bean("warningQueue")
     public Queue warningQueue(){
        return QueueBuilder.durable(WARNING_QUEUE_NAME).build();
     }
     // 声明报警队列绑定关系
     @Bean
     public Binding warningBinding(@Qualifier("warningQueue") Queue queue,
     @Qualifier("backupExchange") FanoutExchange backupExchange){
       return BindingBuilder.bind(queue).to(backupExchange);
     }
     // 声明备份队列
     @Bean("backQueue")
     public Queue backQueue(){
       return QueueBuilder.durable(BACKUP_QUEUE_NAME).build();
     }
     // 声明备份队列绑定关系
     @Bean
     public Binding backupBinding(@Qualifier("backQueue") Queue queue, @Qualifier("backupExchange") FanoutExchange backupExchange){
       return BindingBuilder.bind(queue).to(backupExchange);
     }
}

 报警消费者

import lombok.extern.slf4j.Slf4j;
import org.springframework.amqp.core.Message;
import org.springframework.amqp.rabbit.annotation.RabbitListener;
import org.springframework.stereotype.Component;

@Component
@Slf4j
public class ConfirmConsumer {
   public static final String CONFIRM_QUEUE_NAME = "confirm.queue";
   @RabbitListener(queues =CONFIRM_QUEUE_NAME)
   public void receiveMsg(Message message){
     String msg=new String(message.getBody());
     log.info("接受到队列 confirm.queue 消息:{}",msg);
   }
}

 文章来源地址https://www.toymoban.com/news/detail-670579.html

结果 

springboot整合rabbitmq发布确认高级,rabbitmq,java-rabbitmq,spring boot,rabbitmq

mandatory 参数与备份交换机可以一起使用的时候,如果两者同时开启,消息究竟何去何从?谁优先级高,经过上面结果显示答案是备份交换机优先级高

 

到了这里,关于springboot整合rabbitmq发布确认高级的文章就介绍完了。如果您还想了解更多内容,请在右上角搜索TOY模板网以前的文章或继续浏览下面的相关文章,希望大家以后多多支持TOY模板网!

本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若转载,请注明出处: 如若内容造成侵权/违法违规/事实不符,请点击违法举报进行投诉反馈,一经查实,立即删除!

领支付宝红包 赞助服务器费用

相关文章

  • 8. springboot + rabbitmq 消息发布确认机制

    在 RabbitMQ之生产者发布确认原理章节已经介绍了rabbitmq生产者是如何对消息进行发布确认保证消息不丢失的。本章节继续看下springboot整合rabbitmq后是如何保证消息不丢失的。 消息正常是通过生产者生产消息传递到交换机,然后经过交换机路由到消息队列中,最后消费者消费,

    2023年04月25日
    浏览(64)
  • 消息队列-RabbitMQ:发布确认—发布确认逻辑和发布确认的策略

    生产者将信道设置成 confirm 模式,一旦信道进入 confirm 模式,所有在该信道上面发布的消息都将会被指派一个唯一的 ID (从 1 开始),一旦消息被投递到所有匹配的队列之后,broker 就会发送一个确认给生产者 (包含消息的唯一 ID),这就使得生产者知道消息已经正确到达目的队列

    2024年02月21日
    浏览(38)
  • RabbitMQ(一) - 基本结构、SpringBoot整合RabbitMQ、工作队列、发布订阅、直接、主题交换机模式

    Publisher : 生产者 Queue: 存储消息的容器队列; Consumer:消费者 Connection:消费者与消息服务的TCP连接 Channel:信道,是TCP里面的虚拟连接。例如:电缆相当于TCP,信道是一条独立光纤束,一条TCP连接上创建多少条信道是没有限制的。TCP一旦打开,就会出AMQP信道。无论是发布消息

    2024年02月14日
    浏览(55)
  • RabbitMQ(二) - RabbitMQ与消息发布确认与返回、消费确认

    SpringBoot与RabbitMQ整合后,对RabbitClient的“确认”进行了封装、使用方式与RabbitMQ官网不一致; 生产者给交换机发送消息后、若是不管了,则会出现消息丢失; 解决方案1: 交换机接受到消息、给生产者一个答复ack, 若生产者没有收到ack, 可能出现消息丢失,因此重新发送消息;

    2024年02月14日
    浏览(48)
  • RabbitMQ 发布确认机制

    发布确认模式是避免消息由生产者到RabbitMQ消息丢失的一种手段   生产者通过调用channel.confirmSelect方法将信道设置为confirm模式,之后RabbitMQ会返回Confirm.Select-OK命令表示同意生产者将当前信道设置为confirm模式。   confirm模式下的信道所发送的消息都将被应带ack或者nack一次

    2024年02月13日
    浏览(40)
  • rabbitmq的发布确认

    生产者将信道设置成 confirm 模式,一旦信道进入 confirm 模式, 所有在该信道上面发布的 消息都将会被指派一个唯一的 ID (从 1 开始),一旦消息被投递到所有匹配的队列之后,broker 就会发送一个确认给生产者(包含消息的唯一 ID),这就使得生产者知道消息已经正确到达目的队

    2024年02月12日
    浏览(51)
  • 【RabbitMQ教程】第三章 —— RabbitMQ - 发布确认

                                                                       💧 【 R a b b i t M Q 教程】第三章—— R a b b i t M Q − 发布确认 color{#FF1493}{【RabbitMQ教程】第三章 —— RabbitMQ - 发布确认} 【 R abbi tMQ 教程】第三章 —— R abbi tMQ − 发布确认

    2024年02月08日
    浏览(89)
  • RabbitMQ系列(10)--RabbitMQ发布确认模式的概念及实现

    概念:虽然我们可以设置队列和队列中的消息持久化,但任然存在消息在持久化的过程中,即在写入磁盘的过程中,消息未完全写入,然后服务器宕机导致消息丢失的情况,发布确认就是为了解决这种情况的概念,在消息完全写入磁盘后才确认消息完全持久化了 1、发布确认

    2024年02月13日
    浏览(39)
  • Java操作RabbitMq并整合SpringBoot

    秋风阁-北溪入江流 RabbitMq自带有专门的管理界面,可以在其管理界面对RabbitMq进行管理查看等操作。 RabbitMq的管理界面的对外端口为 15672 ,当我们启动RabbitMq后,需要启动管理界面插件后才能访问界面。 通过参数配置连接RabbitMq 通过amqp协议连接RabbitMq queueDeclarePassive: 创建或

    2024年02月16日
    浏览(59)
  • RabbitMQ学习——发布订阅/fanout模式 & topic模式 & rabbitmq回调确认 & 延迟队列(死信)设计

    1.rabbitmq队列方式的梳理,点对点,一对多; 2.发布订阅模式,交换机到消费者,以邮箱和手机验证码为例; 3.topic模式,根据规则决定发送给哪个队列; 4.rabbitmq回调确认,setConfirmCallback和setReturnsCallback; 5.死信队列,延迟队列,创建方法,正常—死信,设置延迟时间; 点对

    2024年02月13日
    浏览(77)

觉得文章有用就打赏一下文章作者

支付宝扫一扫打赏

博客赞助

微信扫一扫打赏

请作者喝杯咖啡吧~博客赞助

支付宝扫一扫领取红包,优惠每天领

二维码1

领取红包

二维码2

领红包