• RabbitMQ


    1.消息队列MQ

    1.MQ定义

    • 1.MQ:消息队列(Message Quene),本质是一个队列,先进先出(FIFO),只是队列中存放的内容是message(消息)
    • 2.MQ:是一种跨进程的通信机制,用于上下游传递消息(上游:发送消息方;下游:接收消息方)
    • 3.MQ:是一种常见的上下游逻辑解耦+物理解耦的消息通信服务,使用MQ发送消息,上游只需要依赖MQ,不需要依赖其他额外服务

    2.MQ三大功能

    1.流量削峰

    • 1.通过MQ并发的请求延缓过滤掉无效请求,以此来削弱流量高峰期对服务器的压力
    • 2.缺点:排队时会增加客户端响应时间,但相对于服务宕机更优

    2.应用解耦

    • 1.通过MQ让系统间的数据传输通过MQ中间传递,应用与应用之间没有直接关联,将应用解耦
      在这里插入图片描述

    3.异步处理

    • 1.生产者消费者的处理不同步,消费者只需要调用MQ不需要响应生产者,可以异步处理生产者的信息
    • 2.生产者调用消费者api,无需等待响应结果,当消费者处理结束后,直接通过MQ告知生产者
      在这里插入图片描述

    3.MQ分类

    1.ActiveMQ

    • 1.优点
      • 1.单机吞吐量万级
      • 2.时效性ms
      • 3.可用性高,可基于主从架构实现高可用性
      • 4.消息可靠,较低的概率丢失数据
    • 2.缺点
      • 1.官方社区对系统(ActiceMQ5.x)维护少
      • 2.高吞吐量场景缺少使用案例

    2.Kafka

    • 1.简介
      • 1.LinkedIn开源的分布式发布-订阅消息系统,目前属于Apache顶级项目
      • 2.一开始目的是用于日志收集和传输0.8版本开始支持复制,不支持事务,对消息的重复,丢失,错误没有严格要求
      • 3.适合产生大量数据的互联网服务的数据收集业务
    • 2.优点
      • 1.单机写入TPS约在百万条/秒,吞吐量高
      • 2.时效性ms级,可用性高
      • 3.基于分布式,一个数据多个副本,不会丢失数据导致系统不可用
      • 4.消费者采用Pull方式获取消息,消息有序,通过控制能够保证所有消息被消费且仅消费一次
      • 5.有优秀的第三方Kafka Web管理界面Kafka-Manager
    • 3.缺点
      • 1.单机超过64队列/分区,加载时CPU会飙高,队列越多占用越高,发送消息响应时间变长
      • 2.使用短轮询方式,实时性取决于轮询间隔时间
      • 3.消费失败不支持重试,可能会丢失
      • 4.支持消息顺序,但是一台代理宕机后,就会产生消息乱序

    3.RocketMQ

    • 1.简介
      • 1.阿里开源的消息中间件,纯Java开发
      • 2.具有高吞吐量,高可用性,适合大规模分布式系统应用的特点
      • 3.RocketMQ思路起源于Kafka,但并不是Kafka的一个Copy,它对消息的可靠传输及事务性做了优化
    • 2.优点
      • 1.单机吞吐量十万级,可用性高
      • 2.分布式架构,消息可以做到0丢失
      • 3.支持10亿级别的消息堆积,不会因为堆积导致性能下降
    • 3.缺点
      • 1.支持的客户端语言不多(Java,C++
      • 2.没有在MQ核心中去实现JMS等接口,某些系统迁移需要修改大量代码

    4.RabbitMQ

    • 1.简介

      • 1.使用Erlang语言开发的开源消息队列系统
      • 2.基于AMQP协议来实现
    • 2.优点

      • 1.由于erlang语言的高并发特性,性能较好,吞吐量万级
      • 2.功能完备,健壮,稳定,跨平台,支持多种语言
    • 3.缺点

      • 1.商业版收费
      • 2.学习成本高

    4.MQ选择

    1.Kafka

    • 1.适合大型互联网公司大量数据的数据收集业务
    • 2.适合有日志采集和传输功能的项目

    2.RocketMQ

    • 1.适合金融互联网领域
    • 2.适合大数据量高并发的业务场景

    3.RabbitMQ

    • 1.适合中小型公司
    • 2.适合数据量适中,时效性微妙级的业务场景

    2.RabbitMQ

    • 1.基于AMQP协议,使用erlang语言开发的开源的消息中间件

    1.Rabbit安装

    • 1.参考Linux软件安装Linux安装RabbitMQ

    2.RabbitMQ项目搭建

    1.创建SpringBoot项目

    • 1.参考Spring Boot文章

    2.导入依赖

    <?xml version="1.0" encoding="UTF-8"?>
    <project xmlns="http://maven.apache.org/POM/4.0.0" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
            xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 https://maven.apache.org/xsd/maven-4.0.0.xsd">
       <modelVersion>4.0.0</modelVersion>
       <groupId>com.demp</groupId>
       <artifactId>RabbitMQ</artifactId>
       <version>0.0.1-SNAPSHOT</version>
       <name>RabbitMQ</name>
       <description>Demo project for ShardingSphere</description>
    
       <properties>
           <java.version>1.8</java.version>
           <project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>
           <project.reporting.outputEncoding>UTF-8</project.reporting.outputEncoding>
           <spring-boot.version>2.3.7.RELEASE</spring-boot.version>
       </properties>
    
       <dependencies>
           <dependency>
               <groupId>org.springframework.boot</groupId>
               <artifactId>spring-boot-starter-web</artifactId>
           </dependency>
           
           <dependency>
               <groupId>org.springframework.boot</groupId>
               <artifactId>spring-boot-starter-amqp</artifactId>
           </dependency>
           <!--操作文件流的一个依赖-->
           <dependency>
               <groupId>commons-io</groupId>
               <artifactId>commons-io</artifactId>
               <version>2.6</version>
           </dependency>
           <dependency>
               <groupId>com.alibaba</groupId>
               <artifactId>fastjson</artifactId>
               <version>1.2.47</version>
           </dependency>
    
           <dependency>
               <groupId>com.alibaba</groupId>
               <artifactId>druid</artifactId>
               <version>1.2.9</version>
           </dependency>
    
           <dependency>
               <groupId>mysql</groupId>
               <artifactId>mysql-connector-java</artifactId>
               <scope>runtime</scope>
           </dependency>
    
           <dependency>
               <groupId>com.baomidou</groupId>
               <artifactId>mybatis-plus-boot-starter</artifactId>
               <version>3.3.1</version>
           </dependency>
    
           <dependency>
               <groupId>org.projectlombok</groupId>
               <artifactId>lombok</artifactId>
               <optional>true</optional>
           </dependency>
    
           <dependency>
               <groupId>org.springframework.boot</groupId>
               <artifactId>spring-boot-starter-test</artifactId>
               <scope>test</scope>
               <exclusions>
                   <exclusion>
                       <groupId>org.junit.vintage</groupId>
                       <artifactId>junit-vintage-engine</artifactId>
                   </exclusion>
               </exclusions>
           </dependency>
           <dependency>
               <groupId>com.alibaba</groupId>
               <artifactId>fastjson</artifactId>
               <version>1.2.47</version>
           </dependency>
       </dependencies>
    
       <dependencyManagement>
           <dependencies>
               <dependency>
                   <groupId>org.springframework.boot</groupId>
                   <artifactId>spring-boot-dependencies</artifactId>
                   <version>${spring-boot.version}</version>
                   <type>pom</type>
                   <scope>import</scope>
               </dependency>
           </dependencies>
       </dependencyManagement>
    
       <build>
           <plugins>
               <plugin>
                   <groupId>org.apache.maven.plugins</groupId>
                   <artifactId>maven-compiler-plugin</artifactId>
                   <version>3.8.1</version>
                   <configuration>
                       <source>1.8</source>
                       <target>1.8</target>
                       <encoding>UTF-8</encoding>
                   </configuration>
               </plugin>
               <plugin>
                   <groupId>org.springframework.boot</groupId>
                   <artifactId>spring-boot-maven-plugin</artifactId>
                   <version>2.3.7.RELEASE</version>
                   <configuration>
                       <mainClass>com.atguigu.shargingjdbcdemo.ShargingJdbcDemoApplication</mainClass>
                   </configuration>
                   <executions>
                       <execution>
                           <id>repackage</id>
                           <goals>
                               <goal>repackage</goal>
                           </goals>
                       </execution>
                   </executions>
               </plugin>
           </plugins>
       </build>
    
    </project>
    
    • 1
    • 2
    • 3
    • 4
    • 5
    • 6
    • 7
    • 8
    • 9
    • 10
    • 11
    • 12
    • 13
    • 14
    • 15
    • 16
    • 17
    • 18
    • 19
    • 20
    • 21
    • 22
    • 23
    • 24
    • 25
    • 26
    • 27
    • 28
    • 29
    • 30
    • 31
    • 32
    • 33
    • 34
    • 35
    • 36
    • 37
    • 38
    • 39
    • 40
    • 41
    • 42
    • 43
    • 44
    • 45
    • 46
    • 47
    • 48
    • 49
    • 50
    • 51
    • 52
    • 53
    • 54
    • 55
    • 56
    • 57
    • 58
    • 59
    • 60
    • 61
    • 62
    • 63
    • 64
    • 65
    • 66
    • 67
    • 68
    • 69
    • 70
    • 71
    • 72
    • 73
    • 74
    • 75
    • 76
    • 77
    • 78
    • 79
    • 80
    • 81
    • 82
    • 83
    • 84
    • 85
    • 86
    • 87
    • 88
    • 89
    • 90
    • 91
    • 92
    • 93
    • 94
    • 95
    • 96
    • 97
    • 98
    • 99
    • 100
    • 101
    • 102
    • 103
    • 104
    • 105
    • 106
    • 107
    • 108
    • 109
    • 110
    • 111
    • 112
    • 113
    • 114
    • 115
    • 116
    • 117
    • 118
    • 119
    • 120
    • 121
    • 122
    • 123
    • 124
    • 125

    3.配置文件

    server:
     port: 8091
    spring:
     profiles:
       active: test
     rabbitmq:
       host: 192.168.73.130
       port: 5672
       username: root
       password: root
       virtual-host: /
       listener:
         simple:
           # SpringBoot自动重试
           retry:
             # 开启消费者重试
             enabled: true
             # 最大重试次数(默认无数次)
             max-attempts: 5
             # 重试间隔次数
             initial-interval: 3000
    
    • 1
    • 2
    • 3
    • 4
    • 5
    • 6
    • 7
    • 8
    • 9
    • 10
    • 11
    • 12
    • 13
    • 14
    • 15
    • 16
    • 17
    • 18
    • 19
    • 20
    • 21

    5.配置类

    package com.rabbit.config;
    
    import org.springframework.amqp.core.*;
    import org.springframework.beans.factory.annotation.Qualifier;
    import org.springframework.context.annotation.Bean;
    import org.springframework.context.annotation.Configuration;
    
    import java.util.HashMap;
    
    @Configuration
    public class TtlQueueConfig {
    
       private final static String EXCHANGE_X = "exchange_x";
       private final static String DEAD_EXCHANGE_Y = "dead_exchange_y";
       private final static String QUEUE_QA = "queue_qa";
       private final static String QUEUE_QB = "queue_qb";
       private final static String DEAD_QUEUE_QD = "dead_queue_qd";
    
       @Bean("xExchange")
       public DirectExchange xExchange(){
           return new DirectExchange(EXCHANGE_X);
       }
    
       @Bean("yDeadExchange")
       public DirectExchange yDeadExchange(){
           return new DirectExchange(DEAD_EXCHANGE_Y);
       }
    
       @Bean("qaQueue")
       public Queue qaQueue(){
           HashMap<String, Object> args = new HashMap<>();
           // 声明当前队列绑定的死信交换机
           args.put("x-dead-letter-exchange",DEAD_EXCHANGE_Y);
           // 声明当前队列的死信路由key
           args.put("x-dead-letter-routing-key","YD");
           // 声明队列的TTL
           args.put("x-message-ttl",10000);
           // return new Queue(QUEUE_QA,true,false,true,args);
           return QueueBuilder.durable(QUEUE_QA).withArguments(args).autoDelete().build();
       }
    
       @Bean("qbQueue")
       public Queue qbQueue(){
           HashMap<String, Object> args = new HashMap<>();
           // 声明当前队列绑定的死信交换机
           args.put("x-dead-letter-exchange",DEAD_EXCHANGE_Y);
           // 声明当前队列的死信路由key
           args.put("x-dead-letter-routing-key","YD");
           // 声明队列的TTL
           args.put("x-message-ttl",40000);
           // return new Queue(QUEUE_QB,true,false,true);
           return QueueBuilder.durable(QUEUE_QB).withArguments(args).autoDelete().build();
       }
    
       @Bean("qdDeadQueue")
       public Queue qdDeadQueue(){
           return new Queue(DEAD_QUEUE_QD,true,false,true);
       }
    
       @Bean
       public Binding xExchangeBindQa(@Qualifier("qaQueue") Queue qaQueue,
                                      @Qualifier("xExchange")  DirectExchange xExchange){
           // return BindingBuilder.bind(qaQueue()).to(xExchange()).with("XA");
           return BindingBuilder.bind(qaQueue).to(xExchange).with("XA");
       }
    
       @Bean
       public Binding xExchangeBindQb(){
           return BindingBuilder.bind(qbQueue()).to(xExchange()).with("XB");
       }
    
       @Bean
       public Binding yDeadExchangeBindQd(){
           return BindingBuilder.bind(qdDeadQueue()).to(yDeadExchange()).with("YD");
       }
    
    }
    
    • 1
    • 2
    • 3
    • 4
    • 5
    • 6
    • 7
    • 8
    • 9
    • 10
    • 11
    • 12
    • 13
    • 14
    • 15
    • 16
    • 17
    • 18
    • 19
    • 20
    • 21
    • 22
    • 23
    • 24
    • 25
    • 26
    • 27
    • 28
    • 29
    • 30
    • 31
    • 32
    • 33
    • 34
    • 35
    • 36
    • 37
    • 38
    • 39
    • 40
    • 41
    • 42
    • 43
    • 44
    • 45
    • 46
    • 47
    • 48
    • 49
    • 50
    • 51
    • 52
    • 53
    • 54
    • 55
    • 56
    • 57
    • 58
    • 59
    • 60
    • 61
    • 62
    • 63
    • 64
    • 65
    • 66
    • 67
    • 68
    • 69
    • 70
    • 71
    • 72
    • 73
    • 74
    • 75
    • 76
    • 77
    • 1.注意
      • 1.搭建SpringBoot项目前,交换机,队列和绑定关系都嵌套在业务代码
      • 2.搭建SpringBoot项目后,交换机,队列和绑定关系使用单独的配置类设置
      • 3.实际开发中使用配置类进行配置,

    3.RabbitMQ四大核心

    1.生产者

    • 1.生产数据发送消息的程序是生产者

    2.消费者

    • 1.消费数据接收消息的程序是消费者

    3.交换机(Exchange)

    • 1.作用:接收生产者的消息,将消息推送到指定队列
    • 2.RabbitMQ消息传递模型的核心思想:生产者的消息不会直接发送到队列,只能将消息发送到交换机(Exchange),并由其推送到指定队列
    • 3.交换机须确切知道如何处理其接收到的消息,是将这些消息推送到特定队列还是推送到多个队列或是把消息丢弃,这个由交换机类型决定
    • 4.如果生产者没有指定交换机则使用默认交换机,默认交换机的类型是direct
    1.交换机类型

    在这里插入图片描述

    • 1.direct:直连交换机
    • 2.fanout:扇形交换机
    • 3.headers:头交换机
    • 4.topic:主题交换机
    1.直连交换机(direct)
    • 1.简单模式工作模式的代码中没有指定交换机仍然能将消息发送到队列,是因为使用的默认交换机,可以通过空字符串进行标识
      channel.basicPublish("","hello",null,message.getBytes())
      
      • 1
      • 说明
        • 1.第一个参数是交换机的名称,空字符串表示默认/无名交换机
        • 2.Exchange通过Routing Key绑定指定队列,消息通过Exchange路由发送到绑定的队列中
        • 3.如果没有指定交换机名称,则后面Routing Key参数是队列名称,可直接发送到该队列
        • 4.如果指定交换机名称,则后面Routing Key参数是路由键
    • 2.默认交换机的类型是直连交换机,可以根据Routing Key绑定投递到不同队列
    • 3.直连交换机绑定类型有两种
      • 1.单个绑定:一个路由键对应一个队列,一个消息通过交换机只会路由到一个队列中
        在这里插入图片描述
      • 2.多个绑定:一个路由键对应多个队列,一个消息通过交换机可以路由到多个队列中,此时类似扇形交换机
        在这里插入图片描述
    2.扇形交换机(fanout)
    • 1.扇形交换机采用广播模式,根据绑定的交换机,路由到与之对应的所有队列
    • 2.发送到扇形交换机的消息都会被转发到与该交换机绑定的所有队列上
    • 3.类似子网广播,每台子网内的主机都获得了一份复制的消息
    • 4.Fanout交换机转发消息是最快的,其路由键可以不用设置,因为广播到所有
      在这里插入图片描述
    3.头交换机(headers)
    • 1.headers交换机不处理路由键,而是根据发送消息内容中的headers属性进行匹配
    • 2.绑定QueueExchange时指定一组键值对,当消息发送到Broker时会取到该消息的headersExchange绑定时指定的键值对进行匹配
    • 3.如果完全匹配则消息会路由到该队列,否则不会路由到该队列
    • 4.headers属性是一个键值对,可以是哈希表,且键值对的值可以是任何类型
    • 5.directfanouttopic 的路由键都需要要字符串形式
    • 6.匹配规则x-match有下列两种类型
      • 1.x-match = all:表示所有的键值对都匹配才能接受到消息
      • 2.x-match = any:表示只要有键值对匹配就能接受到消息
        在这里插入图片描述
    4.主题交换机(topic)
    • 1.topic交换机可以对路由键进行模糊模式匹配后进行投递
      • 1.符号#表示匹配零或一个或多个任意词
      • 2.符号*表示匹配一个任意词
    • 2.topic交换机的路由键不能随意定义,必须是一个单词列表,以.分割
      在这里插入图片描述
    2.绑定(bindings)
    • 1.其是ExchangeQueue之间的桥梁,告知Exchange和哪个或哪些Queue进行了绑定
    • 2.通过绑定关系RoutingKey,决定发送的队列
      在这里插入图片描述

    4.队列

    • 1.队列是RabbitMQ内部使用的一种数据结构,本质上是一个消息缓冲区
    • 2.消息流经 RabbitMQ 和应用程序,但只能存储在队列中
    • 3.队列仅受主机的内存磁盘的限制
    • 4.多个生产者可将消息发送到一个队列,多个消费者也可从一个队列接收消息
    1.临时队列
    • 1.当连接到RabbitMQ时,都需要一个全新的空队列
    • 2.因此可以创建一个具有随机名称的队列,或让服务器选择一个随机队列名称,一旦断开消费者的连接,队列将自动删除
    • 3.创建临时队列的方式如下,其中AD Excl表示临时队列标识
      String queueName = channel.queueDeclare().getQueue()
      
      • 1
    2.持久化队列
    • 1.当需要一个队列下次重启或遇到问题时队列依然存在,则需要创建一个持久化队列
    • 2.持久化队列只需要在声明时采用持久化策略即可,其中D表示持久化队列标识
    3.死信队列
    • 1.死信:无法被消费的信息
    • 2.死信队列:存放无法被消费信息的队列
    • 3.Producer将消息投递到Broker或直接到QueueConsumerQueue取出消息消费,但某时由于特定的原因导致Queue中的某些消息无法被消费,该消息如果没有后续的处理,就变成了死信,存放其的队列是死信队列
      在这里插入图片描述
    1.应用场景
    • 1.保证订单业务的消息数据不丢失,需要使用死信队列机制,当消息发生异常时,将消息投送到死信队列
    • 2.下单成功并点击去支付后在指定时间未支付时自动失效
    2.死信的三大来源
    • 1.消息TTL(存活时间)过期
    • 2.队列达到最大长度(队列满了,无法再添加数据)
    • 3.消息被拒绝(basic.rejectbasic.nack)并且设置不重新入队(requeue=false)
    3.生产者

    在这里插入图片描述

    package com.rabbit.work.dead;
    
    import com.rabbit.util.RabbitMQUtils;
    import com.rabbitmq.client.AMQP;
    import com.rabbitmq.client.Channel;
    
    import java.util.Scanner;
    
    /**
    * 发消息给普通交换机
    */
    public class WorkProducer {
       // 普通交换机名称
       private final static String NORMAL_EXCHANGE_NAME = "normal_exchange";
       // 普通交换机和普通队列绑定键名称
       private final static String NORMAL_EXCHANGE_Binding_QUEUE = "normal_bind";
       public static void main(String[] args) throws Exception {
           Channel channel = RabbitMQUtils.getChannel();
           // 声明交换机 只需声明一次即可,消费者已经声明,生产者可不声明
           // channel.exchangeDeclare(NORMAL_EXCHANGE_NAME, BuiltinExchangeType.DIRECT);
           Scanner scanner = new Scanner(System.in);
           System.out.println("请输入信息");
           while (scanner.hasNext()) {
               String message = scanner.next();
               // 三大来源 1.设置TTL时间 time to live 单位ms
               AMQP.BasicProperties properties = new AMQP.BasicProperties()
                       .builder().expiration("10000").deliveryMode(2).build();
               // 发送消息给普通交换机并将消息持久化 MessageProperties.PERSISTENT_TEXT_PLAIN 本质是设置传递模式为2
               channel.basicPublish(NORMAL_EXCHANGE_NAME,NORMAL_EXCHANGE_Binding_QUEUE,properties,message.getBytes());
               System.out.println("消息发送完毕: " + message);
           }
       }
    }
    
    • 1
    • 2
    • 3
    • 4
    • 5
    • 6
    • 7
    • 8
    • 9
    • 10
    • 11
    • 12
    • 13
    • 14
    • 15
    • 16
    • 17
    • 18
    • 19
    • 20
    • 21
    • 22
    • 23
    • 24
    • 25
    • 26
    • 27
    • 28
    • 29
    • 30
    • 31
    • 32
    • 33
    4.消费者
    package com.rabbit.work.dead;
    
    import com.rabbit.util.RabbitMQUtils;
    import com.rabbitmq.client.BuiltinExchangeType;
    import com.rabbitmq.client.CancelCallback;
    import com.rabbitmq.client.Channel;
    import com.rabbitmq.client.DeliverCallback;
    
    import java.util.HashMap;
    
    /**
    * 接收普通消息
    */
    public class WorkConsumer_01 {
       // 普通交换机名称
       private final static String NORMAL_EXCHANGE_NAME = "normal_exchange";
       // 普通队列名称
       private final static String NORMAL_QUEUE_NAME = "normal_queue";
       // 死信交换机名称
       private final static String DEAD_EXCHANGE_NAME = "dead_exchange";
       // 死信队列名称
       private final static String DEAD_QUEUE_NAME = "dead_queue";
       // 普通交换机和普通队列绑定键名称
       private final static String NORMAL_EXCHANGE_Binding_QUEUE = "normal_bind";
       // 死信交换机和死信队列绑定键名称
       private final static String DEAD_EXCHANGE_Binding_QUEUE = "dead_bind";
    
       public static void main(String[] args) throws Exception {
           Channel channel = RabbitMQUtils.getChannel();
           // 声明普通交换机
           channel.exchangeDeclare(NORMAL_EXCHANGE_NAME, BuiltinExchangeType.DIRECT);
           // 声明死信交换机
           channel.exchangeDeclare(DEAD_EXCHANGE_NAME, BuiltinExchangeType.DIRECT);
    
           // 普通队列绑定死信交换机 声明普通队列时需要设置参数才能连接
    
           HashMap<String, Object> bindArgs = new HashMap<>();
           // 1.普通队列设置关联死信交换机,其中x-dead-letter-exchange是固定键,不能更改
           bindArgs.put("x-dead-letter-exchange",DEAD_EXCHANGE_NAME);
           // 2.普通队列设置死信交换机绑定死信队列的routing key,其中x-dead-letter-routing-key是固定键,不能更改,其值和绑定死信交换机与死信队列的绑定建一致
           bindArgs.put("x-dead-letter-routing-key",DEAD_EXCHANGE_Binding_QUEUE);
           
           //三大来源
           
           // 1.设置过期时间 单位ms 一般在生产者设置消息过期时间更灵活,其中x-message-ttl是固定键,不能更改
           // bindArgs.put("x-message-ttl",10000);
           // 2.设置普通队列的长度限制 一般在声明时设置 同时去掉生产者中过期消息时间的配置 标记是Lim
           // bindArgs.put("x-max-length",6);
    
           // 声明普通持久化队列 并设置参数
           channel.queueDeclare(NORMAL_QUEUE_NAME,true,false,false,bindArgs);
           // 声明死信持久化队列
           channel.queueDeclare(DEAD_QUEUE_NAME,true,false,false,null);
           // 绑定普通交换机与普通队列
           channel.queueBind(NORMAL_QUEUE_NAME,NORMAL_EXCHANGE_NAME,NORMAL_EXCHANGE_Binding_QUEUE);
           // 绑定死信交换机与死信队列
           channel.queueBind(DEAD_QUEUE_NAME,DEAD_EXCHANGE_NAME,DEAD_EXCHANGE_Binding_QUEUE);
    
           System.out.println("WorkConsumer_01等待接收消息....");
    
           DeliverCallback deliverCallback = (consumerTag, message) -> {
    	        // 三大来源 3.拒收消息且不重新入队
               //  if(new String(message.getBody()).equals("测试值")) {
               //      System.out.println("WorkConsumer_01接收到的消息为:" + new String(message.getBody()) + "并拒绝该消息");
               //      channel.basicReject(message.getEnvelope().getDeliveryTag(),false);
               //  }else{
               //      System.out.println("WorkConsumer_01接收到的消息为:" + new String(message.getBody()));
               //      channel.basicAck(message.getEnvelope().getDeliveryTag(),false);
               //  }
               System.out.println("WorkConsumer_01接收到的消息为:" + new String(message.getBody()) + " :: 路由键:" + message.getEnvelope().getRoutingKey());
               /**
                *  肯定确认
                *  1.消息标记,每一个消息都有一个独立的标记
                *  2.是否批量应答未应答消息
                */
               channel.basicAck(message.getEnvelope().getDeliveryTag(),false);
           };
           CancelCallback cancelCallback = (consumerTag) -> System.out.println("消息消费被中断");
           // 设置不公平分发,默认值为0
           channel.basicQos(1);
           // 采用手动应答
           channel.basicConsume(NORMAL_QUEUE_NAME,false,deliverCallback,cancelCallback);
       }
    }
    
    • 1
    • 2
    • 3
    • 4
    • 5
    • 6
    • 7
    • 8
    • 9
    • 10
    • 11
    • 12
    • 13
    • 14
    • 15
    • 16
    • 17
    • 18
    • 19
    • 20
    • 21
    • 22
    • 23
    • 24
    • 25
    • 26
    • 27
    • 28
    • 29
    • 30
    • 31
    • 32
    • 33
    • 34
    • 35
    • 36
    • 37
    • 38
    • 39
    • 40
    • 41
    • 42
    • 43
    • 44
    • 45
    • 46
    • 47
    • 48
    • 49
    • 50
    • 51
    • 52
    • 53
    • 54
    • 55
    • 56
    • 57
    • 58
    • 59
    • 60
    • 61
    • 62
    • 63
    • 64
    • 65
    • 66
    • 67
    • 68
    • 69
    • 70
    • 71
    • 72
    • 73
    • 74
    • 75
    • 76
    • 77
    • 78
    • 79
    • 80
    • 81
    • 82
    • 83
    • 84
    package com.rabbit.work.dead;
    
    import com.rabbit.util.RabbitMQUtils;
    import com.rabbitmq.client.BuiltinExchangeType;
    import com.rabbitmq.client.CancelCallback;
    import com.rabbitmq.client.Channel;
    import com.rabbitmq.client.DeliverCallback;
    
    /**
    * 接收死信队列消息
    */
    public class WorkConsumer_02 {
       // 死信队列名称
       private final static String DEAD_QUEUE_NAME = "dead_queue";
    
       public static void main(String[] args) throws Exception {
           Channel channel = RabbitMQUtils.getChannel();
    
           System.out.println("WorkConsumer_02等待接收消息....");
    
           DeliverCallback deliverCallback = (consumerTag, message) -> {
               System.out.println("WorkConsumer_02接收到的消息为:" + new String(message.getBody()) + " :: 路由键:" + message.getEnvelope().getRoutingKey());
               /**
                *  肯定确认
                *  1.消息标记,每一个消息都有一个独立的标记
                *  2.是否批量应答未应答消息
                */
               channel.basicAck(message.getEnvelope().getDeliveryTag(),false);
           };
           CancelCallback cancelCallback = (consumerTag) -> System.out.println("消息消费被中断");
           // 设置不公平分发,默认值为0
           channel.basicQos(1);
           // 采用手动应答
           channel.basicConsume(DEAD_QUEUE_NAME,false,deliverCallback,cancelCallback);
       }
    }
    
    • 1
    • 2
    • 3
    • 4
    • 5
    • 6
    • 7
    • 8
    • 9
    • 10
    • 11
    • 12
    • 13
    • 14
    • 15
    • 16
    • 17
    • 18
    • 19
    • 20
    • 21
    • 22
    • 23
    • 24
    • 25
    • 26
    • 27
    • 28
    • 29
    • 30
    • 31
    • 32
    • 33
    • 34
    • 35
    • 36
    4.延迟队列
    • 1.延迟队列中的信息是在指定时间取出处理
    • 2.延时队列就是用来存放需要在指定时间被处理信息的队列
    • 3.延迟队列即死信队列三大来源中的消息TTL过期
    • 4.基于死信的延迟队列存在一定问题,因此推荐采用基于插件的死信队列
    1.延迟队列使用场景
    • 1.使用实例
      • 1.订单在十分钟之内未支付则自动取消
      • 2.新创建的店铺等,十天内没有上传过商品则自动发送消息提醒
      • 3.用户注册成功后,如果三天内没有登录则进行短信提醒
      • 4.用户发起退款,如果三天内没有得到处理则通知相关运营人员
      • 5.预定会议后,需要在预定的时间点前十分钟通知各个与会人员参加会议
    • 2.场景总结
      • 1.场景特征:需要在某个事件发生之后或之前的指定时间点完成某一项任务
      • 2.定时任务:每隔一段时间调用一次的任务
      • 1.发送订单生成事件并在十分钟之后检查该订单支付状态,然后将未支付的订单进行关闭
      • 2.如果使用定时任务,一直轮询数据,每隔一段时间处理一次,数据量大时性能很低
      • 3.如果使用延迟队列,可直接发送后设置定期时间,到指定时间后再进行处理
    2.实际应用流程图

    在这里插入图片描述

    • 1.用户抢购车票但未支付
    • 2.将车票订单入库,将状态设置为未支付0,并将订单信息加入到消息队列中设置延迟时间为30分钟
    • 3.如果30分钟内用户支付则更新车票订单状态为已支付1
    • 4.30分钟后延迟队列中对应的监听器监听到消息队列中的信息执行查看该车票订单支付状态,如果未支付则更新订单状态为已失效
    3.实际应用代码

    在这里插入图片描述

    • 1.创建一个direct类型的普通交换机X和一个direct类型的死信交换机Y
    • 2.创建两个普通队列QAQB,创建一个死信队列QD
    • 3.两个普通队列的TTL(time to live)分别设置为10S40S
    • 4.交换机和队列的绑定关系如上图所示,具体项目搭建和配置参考上述RabbitMQ项目搭建
    1.生产者
    package com.rabbit.controller;
    
    import lombok.extern.slf4j.Slf4j;
    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.Resource;
    import java.util.Date;
    
    @Slf4j
    @RestController
    @RequestMapping("/ttl")
    public class SendMessageController {
    
       @Resource
       RabbitTemplate rabbitTemplate;
    
       @GetMapping("/send/{msg}")
       public void sendMessage(@PathVariable String msg){
           log.info("当前时间:{},发送一条信息给两个 TTL 队列:{}", new Date(), msg);
           rabbitTemplate.convertAndSend("exchange_x", "XA", "消息来自 ttl 为 10S 的队列: "+msg);
           rabbitTemplate.convertAndSend("exchange_x", "XB", "消息来自 ttl 为 40S 的队列: "+msg);
       }
    }
    
    • 1
    • 2
    • 3
    • 4
    • 5
    • 6
    • 7
    • 8
    • 9
    • 10
    • 11
    • 12
    • 13
    • 14
    • 15
    • 16
    • 17
    • 18
    • 19
    • 20
    • 21
    • 22
    • 23
    • 24
    • 25
    • 26
    • 27
    2.消费者
    package com.rabbit.listener;
    
    import com.rabbitmq.client.Channel;
    import lombok.extern.slf4j.Slf4j;
    import org.springframework.amqp.core.Message;
    import org.springframework.amqp.rabbit.annotation.RabbitListener;
    import org.springframework.stereotype.Component;
    
    import java.util.Date;
    
    @Slf4j
    @Component
    public class ReceiveMessageListener {
    
       @RabbitListener(queues = "dead_queue_qd")
       public void receiveMessage(Message message, Channel channel){
           String msg = new String(message.getBody());
           log.info("当前时间:{},收到死信队列信息{}", new Date().toString(), msg);
       }
    }
    
    • 1
    • 2
    • 3
    • 4
    • 5
    • 6
    • 7
    • 8
    • 9
    • 10
    • 11
    • 12
    • 13
    • 14
    • 15
    • 16
    • 17
    • 18
    • 19
    • 20
    3.测试结果

    在这里插入图片描述

    4.延迟队列优化

    在这里插入图片描述

    • 1.上述延迟队列的缺点:每增加一个新的时间需求,就要新增一个队列
    • 2.新增一个通用队列QC,其消息过期时间生产者发送时设置,而非队列本身设置
    1.配置类
    • 1.上述的配置类中新增以下设置
    	private final static String QUEUE_QC = "queue_qc";
    
       @Bean("qcQueue")
       public Queue qcQueue(){
           HashMap<String, Object> args = new HashMap<>();
           // 声明当前队列绑定的死信交换机
           args.put("x-dead-letter-exchange",DEAD_EXCHANGE_Y);
           // 声明当前队列的死信路由key
           args.put("x-dead-letter-routing-key","YD");
           // 不声明队列的TTL
           // args.put("x-message-ttl",40000);
           // return new Queue(QUEUE_QB,true,false,true);
           return QueueBuilder.durable(QUEUE_QC).withArguments(args).autoDelete().build();
       }
    
       @Bean
       public Binding xExchangeBindQc(){
           return BindingBuilder.bind(qcQueue()).to(xExchange()).with("XC");
       }
    
    • 1
    • 2
    • 3
    • 4
    • 5
    • 6
    • 7
    • 8
    • 9
    • 10
    • 11
    • 12
    • 13
    • 14
    • 15
    • 16
    • 17
    • 18
    • 19
    2.生产者
    @GetMapping("sendExpire/{message}/{ttlTime}")
       public void sendExpire(@PathVariable String message,@PathVariable String ttlTime) {
           rabbitTemplate.convertAndSend("exchange_x", "XC", message, messageConfig ->{
               messageConfig.getMessageProperties().setExpiration(ttlTime);
               return messageConfig;
           });
           log.info("当前时间:{},发送一条时长{}毫秒 TTL 信息给队列 C:{}", new Date(),ttlTime, message);
       }
    
    • 1
    • 2
    • 3
    • 4
    • 5
    • 6
    • 7
    • 8
    3.消费者
    • 同上述消费者
    4.测试结果

    在这里插入图片描述

    5.基于插件的延迟队列
    1.基于死信的延迟队列问题
    • 1.RabbitMQ 只会检查队列中的第一个消息是否过期
    • 2.如果过期则丢到死信队列,如果第一个消息的延时时长很长,而第二个消息的延时时长很短,则第二个消息并不会优先得到执行
    • 3.通过RabbitMQ插件实现延迟队列可以解决上述问题,本质上通过交换机延迟实现
    2.安装延迟队列插件
    • 1.参考Linux软件安装文章中的安装延迟队列插件
    3.代码架构图
    • 1.基于死信队列
      在这里插入图片描述
    • 2.基于插件
      在这里插入图片描述
    4.配置文件
    	// 延迟队列
       private final static String DELAYED_QUEUE = "delayed_queue";
    
       // 延迟交换机
       private final static String DELAYED_EXCHANGE = "delayed_exchange";
    
       // 绑定键
       private final static String DELAYED_ROUTING_KEY = "delayed.routing_key";
    
       // 声明延迟队列
       @Bean
       public Queue delayedQueue(){
           return new Queue(DELAYED_QUEUE);
       }
    
       // 声明延迟交换机 基于插件自定义交换机
       @Bean
       public CustomExchange delayedExchange(){
           // 参数设置,设置交换机路由类型
           HashMap<String, Object> args = new HashMap<>();
           args.put("x-delayed-type","direct");
           /**
            * 1.交换机的名称
            * 2.交换机的类型
            * 3.是否需要持久化
            * 4.是否需要自动删除
            * 5.其他的参数
            */
           return new CustomExchange(DELAYED_EXCHANGE,"x-delayed-message",true,false,args);
       }
    
       // 绑定
       @Bean
       public Binding delayedQueueBindingExchange(){
           return BindingBuilder.bind(delayedQueue()).to(delayedExchange()).with(DELAYED_ROUTING_KEY).noargs();
       }
    
    • 1
    • 2
    • 3
    • 4
    • 5
    • 6
    • 7
    • 8
    • 9
    • 10
    • 11
    • 12
    • 13
    • 14
    • 15
    • 16
    • 17
    • 18
    • 19
    • 20
    • 21
    • 22
    • 23
    • 24
    • 25
    • 26
    • 27
    • 28
    • 29
    • 30
    • 31
    • 32
    • 33
    • 34
    • 35
    • 36
    5.生产者
       @GetMapping("sendDelayedMsg/{message}/{delayedTime}")
       public void sendDelayedMsg(@PathVariable String message,@PathVariable Integer delayedTime) {
           log.info("当前时间:{},发送一条时长{}毫秒的信息给延迟队列 :{}", new Date(),delayedTime, message);
           rabbitTemplate.convertAndSend("delayed_exchange", "delayed.routing.key", message, messageConfig ->{
           	// 注意此处设置不同于基于死信的延迟队列,延迟时间为Integer类型
               messageConfig.getMessageProperties()setDelay(delayedTime);
               return messageConfig;
           });
       }
    
    • 1
    • 2
    • 3
    • 4
    • 5
    • 6
    • 7
    • 8
    • 9
    6.消费者
    package com.rabbit.listener;
    
    import com.rabbitmq.client.Channel;
    import lombok.extern.slf4j.Slf4j;
    import org.springframework.amqp.core.Message;
    import org.springframework.amqp.rabbit.annotation.RabbitListener;
    import org.springframework.stereotype.Component;
    
    import java.util.Date;
    
    @Slf4j
    @Component
    public class ReceiveMessageListener {
    
       @RabbitListener(queues = "delayed_queue")
       public void receiveMessage(Message message, Channel channel){
           String msg = new String(message.getBody());
           log.info("当前时间:{},收到死信队列信息{}", new Date().toString(), msg);
       }
    }
    
    • 1
    • 2
    • 3
    • 4
    • 5
    • 6
    • 7
    • 8
    • 9
    • 10
    • 11
    • 12
    • 13
    • 14
    • 15
    • 16
    • 17
    • 18
    • 19
    • 20
    7.测试结果

    在这里插入图片描述

    • 1.结果不同于基于死信队列的延迟队列,可以避免该缺陷
    5.优先级队列
    • 1.队列根据消息的优先级进行分发
    • 2.优先级的取值范围为0~255,值越大越优先执行,一般设置为0 ~ 10,过大会耗费CPU性能
    1.使用方式
    • 1.设置队列时添加优先级的范围0~255
    • 2.发送信息时设置消息的优先级,其范围需在队列设置范围之内
    • 3.注意:实现队列优先级需要消费者等待消息全部已发送到队列中再去消费,因为这样才有机会对消息进行排序
    2.配置类
    package com.rabbit.config;
    
    import org.springframework.amqp.core.*;
    import org.springframework.context.annotation.Bean;
    import org.springframework.context.annotation.Configuration;
    
    import java.util.HashMap;
    
    @Configuration
    public class PriConfig {
       public final static String PRIORITY_QUEUE = "priority_queue";
       public final static String LAZY_QUEUE = "lazy_queue";
       public final static String PRIORITY_EXCHANGE = "priority_exchange";
       public final static String LAZY_EXCHANGE = "lazy_exchange";
       public final static String PRIORITY_QUEUE_BINDING_EXCHANGE = "priority";
       public final static String LAZY_QUEUE_BINDING_EXCHANGE = "lazy";
    
       @Bean
       public Queue priorityQueue(){
           HashMap<String, Object> args = new HashMap<>();
           args.put("x-max-priority",10);
           return new Queue(PRIORITY_QUEUE,true,false,false,args);
           // 或 return QueueBuilder.durable(PRIORITY_QUEUE).maxPriority(10).build();
       }
    
       @Bean
       public Queue lazyQueue(){
           HashMap<String, Object> args = new HashMap<>();
           args.put("x-queue-mode","lazy");
           return new Queue(LAZY_QUEUE,true,false,false,args);
           // 或return QueueBuilder.durable(LAZY_QUEUE).lazy().build();
       }
    
       @Bean
       public DirectExchange priorityExchange(){
           return new DirectExchange(PRIORITY_EXCHANGE,true,false,null);
           // 或return ExchangeBuilder.directExchange(PRIORITY_EXCHANGE).build();
       }
    
       @Bean
       public DirectExchange lazyExchange(){
           return new DirectExchange(LAZY_EXCHANGE,true,false,null);
           // 或return ExchangeBuilder.directExchange(LAZY_EXCHANGE).build();
       }
    
       @Bean
       public Binding priorityExchangeBindingQueue(){
           return BindingBuilder.bind(priorityQueue()).to(priorityExchange()).with(PRIORITY_QUEUE_BINDING_EXCHANGE);
       }
    
       @Bean
       public Binding lazyExchangeBindingQueue(){
           return BindingBuilder.bind(lazyQueue()).to(lazyExchange()).with(LAZY_QUEUE_BINDING_EXCHANGE);
           // 或return new Binding(LAZY_QUEUE, Binding.DestinationType.QUEUE,LAZY_EXCHANGE,LAZY_QUEUE_BINDING_EXCHANGE,null);
       }
    }
    
    • 1
    • 2
    • 3
    • 4
    • 5
    • 6
    • 7
    • 8
    • 9
    • 10
    • 11
    • 12
    • 13
    • 14
    • 15
    • 16
    • 17
    • 18
    • 19
    • 20
    • 21
    • 22
    • 23
    • 24
    • 25
    • 26
    • 27
    • 28
    • 29
    • 30
    • 31
    • 32
    • 33
    • 34
    • 35
    • 36
    • 37
    • 38
    • 39
    • 40
    • 41
    • 42
    • 43
    • 44
    • 45
    • 46
    • 47
    • 48
    • 49
    • 50
    • 51
    • 52
    • 53
    • 54
    • 55
    • 56
    3.生产者
    @GetMapping("/sendPri/{message}")
       public void sendPri(@PathVariable String message) {
           for(int i=0; i<10; i++){
               CorrelationData correlationData = new CorrelationData();
               int id = i+1;
               correlationData.setId(String.valueOf(id));
               MessageProperties messageProperties = new MessageProperties();
               Message msg = new Message((message+i).getBytes(),messageProperties);
               correlationData.setReturnedMessage(msg);
               if(i == 5){
                   msg.getMessageProperties().setPriority(5);
                   rabbitTemplate.convertAndSend(PriConfig.PRIORITY_EXCHANGE, PriConfig.PRIORITY_QUEUE_BINDING_EXCHANGE, msg,correlationData);
               }else {
                   rabbitTemplate.convertAndSend(PriConfig.PRIORITY_EXCHANGE, PriConfig.PRIORITY_QUEUE_BINDING_EXCHANGE, msg,correlationData);
               }
           }
           log.info("优先级消息发送完毕");
       }
    
    • 1
    • 2
    • 3
    • 4
    • 5
    • 6
    • 7
    • 8
    • 9
    • 10
    • 11
    • 12
    • 13
    • 14
    • 15
    • 16
    • 17
    • 18
    4.消费者
    @RabbitListener(queues = PriConfig.PRIORITY_QUEUE)
       public void receivePriorityMessage(String Object, Message message, Channel channel){
           long deliveryTag = message.getMessageProperties().getDeliveryTag();
           log.info("收到" + PriConfig.PRIORITY_QUEUE + "队列的 {} 消息: {}", deliveryTag, Object);
           try {
               /**
                * 业务代码
                */
           } catch (Exception e) {
               log.error("签收失败", e);
               /**
                * 记录日志、发送邮件、保存消息到数据库,落库之前判断如果消息已经落库就不保存
                */
               throw new RuntimeException("消息消费失败");
           }
           log.info("消费成功: {},消费内容: {}", deliveryTag, Object);
       }
    
    • 1
    • 2
    • 3
    • 4
    • 5
    • 6
    • 7
    • 8
    • 9
    • 10
    • 11
    • 12
    • 13
    • 14
    • 15
    • 16
    • 17
    5.测试结果

    在这里插入图片描述

    6.惰性队列
    • 1.普通队列将消息存放在内存
    • 2.惰性队列将消息存放在磁盘中,
    • 3.当消费者消费到相应的消息时才会被加载到内存中
    • 4.适用于消息积压的情况
    • 5.相同消息数量,普通队列占用的内存大,响应速度快;而惰性队列占用内存小,但是响应速度慢
    1.配置类
    • 1.参考优先级队列中的配置类

    4.RabbitMQ工作原理图

    在这里插入图片描述

    • 1.Broker:接收和分发消息的应用,RabbitMQ Server就是Broker
    • 2.Virtual host(vhost)
      • 1.出于多租户和安全因素设计,将AMQP(消息队列协议)的基本组件划分到一个虚拟的分组(vhost)中
      • 2.当多个不同用户使用同一个RabbitMQ Server提供的服务时,可以划分出多个vhost,每个用户都在独自的vhost创建exchange/queue
      • 3.多租户Borker包含多个vhost,每个vhost都包含一个ExchangeQueue连接
    • 3.ConnectionProducer/ConsumerBroker之间的TCP连接
    • 4.Channel
      • 1.每次访问RabbitMQ都建立一个TCP Connection,开销大,效率低
      • 2.Channel是在Connection内部建立的逻辑连接,如果应用程序支持多线程,通常每个Thread创建单独的Channel进行通信
      • 3.AMQP包含了Channel id 帮助客户端(消费/生产)和Broker识别Channel,因此Channel之间是完全隔离的
      • 4.Channel作为轻量级的Connection极大地减少了操作系统建立TCP Connection的开销
    • 5.ExchangeMessage达到Broker的第一站,其根据分发规则,匹配查询表中的Rounting key,分发消息到Queue
    • 6.Queue:消息最终被送到这里等待Consumer消费
    • 7.Binding
      • 1.ExchangeQueue之间的虚拟连接
      • 2.Binding中包含Routing keyBinding信息被保存到Exchange查询表中,是Message的分发依据

    5.RabbitMQ六大模式

    1.简单模式(Hello World)

    • 1.简单模式即一个发送一个接收,其使用的是默认直连(direct)交换机
      在这里插入图片描述
    1.生产者
    package com.rabbit;
    
    import com.rabbit.entity.Order;
    import com.rabbit.mapper.OrderMapper;
    import com.rabbitmq.client.*;
    import org.junit.jupiter.api.Test;
    import org.springframework.boot.test.context.SpringBootTest;
    
    import javax.annotation.Resource;
    import java.io.IOException;
    import java.util.List;
    import java.util.concurrent.TimeoutException;
    
    @SpringBootTest
    class RabbitMQStartApplicationTest {
    
       private final static String QUEUE_NAME = "hello";
    
       @Test
       public void producer() throws IOException, TimeoutException {
           // 创建一个连接工厂
           ConnectionFactory factory = new ConnectionFactory();
           // 工厂IP,用于连接RabbitMQ
           factory.setHost("192.168.73.130");
           // 用户名
           factory.setUsername("root");
           // 密码
           factory.setPassword("root");
           // 创建连接
           Connection connection = factory.newConnection();
           // 创建信道 Channel实现了AutoCloseable接口,自动关闭,无需显示关闭
           Channel channel = connection.createChannel();
           /**
            * 声明一个队列
            * 1.队列名称
            * 2.队列中的消息是否持久化(磁盘),默认消息存储在内存中
            * 3.该队列是否私有,true:私有,对当前队列加锁,只能一个消费者消费,其他通道无法访问 false:共享,可以多个消费者访问同一个队列
            * 4.是否自动删除,最后一个消费者断开连接以后,该队列是否自动删除,true:自动删除,false:不自动删除
            * 5.其他参数(延迟消息,死信消息等)
            */
           channel.queueDeclare(QUEUE_NAME,false,false,false,null);
           // 消息体
           String message="hello world";
           /**
            * 发送一个消息
            * 1.exchange 发送到哪个交换机,当前为空,使用默认交换机
            * 2.routingKey 路由到哪个key,当前直接使用队列名称
            * 3.其他的参数信息
            * 4.发送消息的消息体
            */
           channel.basicPublish("",QUEUE_NAME,null,message.getBytes());
           System.out.println("消息发送完毕");
       }
    }
    
    • 1
    • 2
    • 3
    • 4
    • 5
    • 6
    • 7
    • 8
    • 9
    • 10
    • 11
    • 12
    • 13
    • 14
    • 15
    • 16
    • 17
    • 18
    • 19
    • 20
    • 21
    • 22
    • 23
    • 24
    • 25
    • 26
    • 27
    • 28
    • 29
    • 30
    • 31
    • 32
    • 33
    • 34
    • 35
    • 36
    • 37
    • 38
    • 39
    • 40
    • 41
    • 42
    • 43
    • 44
    • 45
    • 46
    • 47
    • 48
    • 49
    • 50
    • 51
    • 52
    • 53
    • 54
    2.消费者
    package com.rabbit;
    
    import com.rabbit.entity.Order;
    import com.rabbit.mapper.OrderMapper;
    import com.rabbitmq.client.*;
    import org.junit.jupiter.api.Test;
    import org.springframework.boot.test.context.SpringBootTest;
    
    import javax.annotation.Resource;
    import java.io.IOException;
    import java.util.List;
    import java.util.concurrent.TimeoutException;
    
    @SpringBootTest
    class RabbitMQStartApplicationTest {
    
       private final static String QUEUE_NAME = "hello";
    
       @Test
       public void consumer() throws IOException, TimeoutException {
           // 创建一个连接工厂
           ConnectionFactory factory = new ConnectionFactory();
           // 工厂IP,用于连接RabbitMQ
           factory.setHost("192.168.73.130");
           // 用户名
           factory.setUsername("root");
           // 密码
           factory.setPassword("root");
           // 创建连接
           Connection connection = factory.newConnection();
           // 创建信道 Channel实现了AutoCloseable接口,自动关闭,无需显示关闭
           Channel channel = connection.createChannel();
    
           System.out.println("等待接收消息....");
    
           // 推送的消息如何进行消费的接口回调
           DeliverCallback deliverCallback = (consumerTag, message) -> {
               // 当前只需要消息的消息体内容,否则输出的消息地址
               String body = new String(message.getBody());
               System.out.println(body);
           };
           // 取消消费的一个回调接口 如在消费的时候队列被删除掉了
           CancelCallback cancelCallback = (consumerTag) -> {
               System.out.println("消息消费被中断");
           };
           /**
            * 消费者消费消息
            * 1.消费队列名称
            * 2.消费成功之后是否要自动应答,true:代表自动应答,false:手动应答
            * 3.消费者未成功消费的回调(即未成功调用的方法,使用Lambda表达式实现函数式接口)
            * 4.消费者取消消费的回调(即消费者取消或中断后调用的方法,使用Lambda表达式实现函数式接口)
            */
           channel.basicConsume(QUEUE_NAME,true,deliverCallback,cancelCallback);
       }
    }
    
    • 1
    • 2
    • 3
    • 4
    • 5
    • 6
    • 7
    • 8
    • 9
    • 10
    • 11
    • 12
    • 13
    • 14
    • 15
    • 16
    • 17
    • 18
    • 19
    • 20
    • 21
    • 22
    • 23
    • 24
    • 25
    • 26
    • 27
    • 28
    • 29
    • 30
    • 31
    • 32
    • 33
    • 34
    • 35
    • 36
    • 37
    • 38
    • 39
    • 40
    • 41
    • 42
    • 43
    • 44
    • 45
    • 46
    • 47
    • 48
    • 49
    • 50
    • 51
    • 52
    • 53
    • 54
    • 55

    2.工作模式(Work quenes)

    • 1.避免大量立即执行的任务,将任务封装为消息发送到队列,后台运行的工作线程依次取出任务并最终执行
    • 2.当后台存在多个工作线程时,这些工作线程共同处理这些任务
    • 3.一个消息只能被执行一次,不可重复执行
    • 4.工作模式中多个消费者采用轮询的方式执行
    • 5.工作模式即一个发送多个争抢接收,采用的是默认直连(direct)交换机
    1.抽取工具类
    package com.rabbit.util;
    
    import com.rabbitmq.client.Channel;
    import com.rabbitmq.client.Connection;
    import com.rabbitmq.client.ConnectionFactory;
    
    public class RabbitMQUtils {
    
       public static Channel getChannel() throws Exception {
           // 创建一个连接工厂
           ConnectionFactory factory = new ConnectionFactory();
           // 工厂IP,用于连接RabbitMQ
           factory.setHost("192.168.73.130");
           // 用户名
           factory.setUsername("root");
           // 密码
           factory.setPassword("root");
           // 创建连接
           Connection connection = factory.newConnection();
           // 创建信道 Channel实现了AutoCloseable接口,自动关闭,无需显示关闭
           return connection.createChannel();
       }
    
    }
    
    • 1
    • 2
    • 3
    • 4
    • 5
    • 6
    • 7
    • 8
    • 9
    • 10
    • 11
    • 12
    • 13
    • 14
    • 15
    • 16
    • 17
    • 18
    • 19
    • 20
    • 21
    • 22
    • 23
    • 24
    2.生产者
    package com.rabbit.work;
    
    import com.rabbit.util.RabbitMQUtils;
    import com.rabbitmq.client.Channel;
    
    import java.util.Scanner;
    
    public class WorkProducer {
       private final static String QUEUE_NAME = "hello";
    
       public static void main(String[] args) throws Exception {
           Channel channel = RabbitMQUtils.getChannel();
           /**
            * 声明一个队列
            * 1.队列名称
            * 2.队列中的消息是否持久化(磁盘),默认消息存储在内存中
            * 3.该队列是否私有,true:私有,对当前队列加锁,只能一个消费者消费,其他通道无法访问 false:共享,可以多个消费者访问同一个队列
            * 4.是否自动删除,最后一个消费者断开连接以后,该队列是否自动删除,true:自动删除,false:不自动删除
            * 5.其他参数(延迟消息,死信消息等)
            */
           channel.queueDeclare(QUEUE_NAME,false,false,false,null);
           Scanner scanner = new Scanner(System.in);
           while (scanner.hasNext()) {
               // 消息体
               String message = scanner.next();
               /**
                * 发送一个消息
                * 1.exchange 发送到哪个交换机,当前为空,使用默认交换机
                * 2.routingKey 路由到哪个key,当前直接使用队列名称
                * 3.其他的参数信息
                * 4.发送消息的消息体
                */
               channel.basicPublish("",QUEUE_NAME,null,message.getBytes());
               System.out.println("消息发送完毕");
           }
       }
    }
    
    • 1
    • 2
    • 3
    • 4
    • 5
    • 6
    • 7
    • 8
    • 9
    • 10
    • 11
    • 12
    • 13
    • 14
    • 15
    • 16
    • 17
    • 18
    • 19
    • 20
    • 21
    • 22
    • 23
    • 24
    • 25
    • 26
    • 27
    • 28
    • 29
    • 30
    • 31
    • 32
    • 33
    • 34
    • 35
    • 36
    • 37

    在这里插入图片描述

    3.消费者
    • 1.需要创建并发的多个消费者,设置允许并发执行
    • 2.每次执行需要修改显示信息,方便查看结果
      package com.rabbit.work;
      
      import com.rabbit.util.RabbitMQUtils;
      import com.rabbitmq.client.CancelCallback;
      import com.rabbitmq.client.Channel;
      import com.rabbitmq.client.DeliverCallback;
      
      public class WorkConsumer {
          private final static String QUEUE_NAME = "hello";
      
          public static void main(String[] args) throws Exception {
              Channel channel = RabbitMQUtils.getChannel();
      
              System.out.println("Work1等待接收消息....");
      
              // 推送的消息如何进行消费的接口回调
              DeliverCallback deliverCallback = (consumerTag, message) -> {
                  // 当前只需要消息的消息体内容,且需要封装为String类,否则输出的消息地址
                  String body = new String(message.getBody());
                  System.out.println("接收到的消息为:" + body);
              };
              // 取消消费的一个回调接口 如在消费的时候队列被删除掉了
              CancelCallback cancelCallback = (consumerTag) -> {
                  System.out.println("消息消费被中断");
              };
              /**
               * 消费者消费消息
               * 1.消费队列名称
               * 2.消费成功之后是否要自动应答,true:代表自动应答,false:手动应答
               * 3.消费者未成功消费的回调(即未成功调用的方法,使用Lambda表达式实现函数式接口)
               * 4.消费者取消消费的回调(即消费者取消或中断后调用的方法,使用Lambda表达式实现函数式接口)
               */
              channel.basicConsume(QUEUE_NAME,true,deliverCallback,cancelCallback);
          }
      }
      
      • 1
      • 2
      • 3
      • 4
      • 5
      • 6
      • 7
      • 8
      • 9
      • 10
      • 11
      • 12
      • 13
      • 14
      • 15
      • 16
      • 17
      • 18
      • 19
      • 20
      • 21
      • 22
      • 23
      • 24
      • 25
      • 26
      • 27
      • 28
      • 29
      • 30
      • 31
      • 32
      • 33
      • 34
      • 35
      在这里插入图片描述
      在这里插入图片描述
      在这里插入图片描述

    3.发布/订阅模式(Publish/Subscribe)

    • 1.发布订阅模式类似广播模式,将Exchange收到的消息广播给所有其绑定的队列
    • 2.发布订阅模式即一个发送绑定都接收,使用的是fanout类型的交换机
      在这里插入图片描述
    1.生产者
    package com.rabbit.work.fanout;
    
    import com.rabbit.util.RabbitMQUtils;
    import com.rabbitmq.client.Channel;
    import com.rabbitmq.client.MessageProperties;
    
    import java.util.Scanner;
    
    /**
     * 发消息给交换机
     */
    public class WorkProducer {
       // 交换机名称
       private final static String EXCHANGE_NAME = "fanout_message";
       public static void main(String[] args) throws Exception {
           Channel channel = RabbitMQUtils.getChannel();
           // 声明交换机
           channel.exchangeDeclare(EXCHANGE_NAME,"fanout");
    
           Scanner scanner = new Scanner(System.in);
           System.out.println("请输入信息:");
           while (scanner.hasNext()) {
               String message = scanner.next();
               // 发送消息给交换机并将消息持久化
               channel.basicPublish(EXCHANGE_NAME,"", MessageProperties.PERSISTENT_TEXT_PLAIN,message.getBytes());
               System.out.println("消息发送完毕: " + message);
           }
       }
    }
    
    • 1
    • 2
    • 3
    • 4
    • 5
    • 6
    • 7
    • 8
    • 9
    • 10
    • 11
    • 12
    • 13
    • 14
    • 15
    • 16
    • 17
    • 18
    • 19
    • 20
    • 21
    • 22
    • 23
    • 24
    • 25
    • 26
    • 27
    • 28
    • 29
    2.消费者
    package com.rabbit.work.fanout;
    
    import com.rabbit.util.RabbitMQUtils;
    import com.rabbit.util.SleepUtils;
    import com.rabbitmq.client.CancelCallback;
    import com.rabbitmq.client.Channel;
    import com.rabbitmq.client.DeliverCallback;
    /**
    * 接收消息1
    */
    public class WorkConsumer_01 {
       // 交换机名称
       private final static String EXCHANGE_NAME = "fanout_message";
    
       public static void main(String[] args) throws Exception {
           Channel channel = RabbitMQUtils.getChannel();
           // 声明交换机
           channel.exchangeDeclare(EXCHANGE_NAME,"fanout");
           // 声明一个临时队列,队列名称随机,消费者断开连接后,队列自动删除
           String queueName = channel.queueDeclare().getQueue();
           // 绑定交换机与队列
           channel.queueBind(queueName,EXCHANGE_NAME,"");
    
           System.out.println("WorkConsumer_01等待接收消息....");
    
           DeliverCallback deliverCallback = (consumerTag, message) -> {
               System.out.println("WorkConsumer_01接收到的消息为:" + new String(message.getBody()));
               /**
                *  肯定确认
                *  1.消息标记,每一个消息都有一个独立的标记
                *  2.是否批量应答未应答消息
                */
               channel.basicAck(message.getEnvelope().getDeliveryTag(),false);
           };
           CancelCallback cancelCallback = (consumerTag) -> System.out.println("消息消费被中断");
           // 设置不公平分发,默认值为0
           channel.basicQos(1);
           // 采用手动应答
           channel.basicConsume(queueName,false,deliverCallback,cancelCallback);
       }
    }
    
    • 1
    • 2
    • 3
    • 4
    • 5
    • 6
    • 7
    • 8
    • 9
    • 10
    • 11
    • 12
    • 13
    • 14
    • 15
    • 16
    • 17
    • 18
    • 19
    • 20
    • 21
    • 22
    • 23
    • 24
    • 25
    • 26
    • 27
    • 28
    • 29
    • 30
    • 31
    • 32
    • 33
    • 34
    • 35
    • 36
    • 37
    • 38
    • 39
    • 40
    • 41
    package com.rabbit.work.fanout;
    
    import com.rabbit.util.RabbitMQUtils;
    import com.rabbit.util.SleepUtils;
    import com.rabbitmq.client.CancelCallback;
    import com.rabbitmq.client.Channel;
    import com.rabbitmq.client.DeliverCallback;
    /**
    * 接收消息2
    */
    public class WorkConsumer_02 {
       // 交换机名称
       private final static String EXCHANGE_NAME = "fanout_message";
    
       public static void main(String[] args) throws Exception {
           Channel channel = RabbitMQUtils.getChannel();
           // 声明交换机
           channel.exchangeDeclare(EXCHANGE_NAME,"fanout");
           // 声明一个临时队列,队列名称随机,消费者断开连接后,队列自动删除
           String queueName = channel.queueDeclare().getQueue();
           // 绑定交换机与队列
           channel.queueBind(queueName,EXCHANGE_NAME,"");
    
           System.out.println("WorkConsumer_02等待接收消息....");
    
           DeliverCallback deliverCallback = (consumerTag, message) -> {
               System.out.println("WorkConsumer_02接收到的消息为:" + new String(message.getBody()));
               /**
                *  肯定确认
                *  1.消息标记,每一个消息都有一个独立的标记
                *  2.是否批量应答未应答消息
                */
               channel.basicAck(message.getEnvelope().getDeliveryTag(),false);
           };
           CancelCallback cancelCallback = (consumerTag) -> System.out.println("消息消费被中断");
           // 设置不公平分发,默认值为0
           channel.basicQos(1);
           // 采用手动应答
           channel.basicConsume(queueName,false,deliverCallback,cancelCallback);
       }
    }
    
    • 1
    • 2
    • 3
    • 4
    • 5
    • 6
    • 7
    • 8
    • 9
    • 10
    • 11
    • 12
    • 13
    • 14
    • 15
    • 16
    • 17
    • 18
    • 19
    • 20
    • 21
    • 22
    • 23
    • 24
    • 25
    • 26
    • 27
    • 28
    • 29
    • 30
    • 31
    • 32
    • 33
    • 34
    • 35
    • 36
    • 37
    • 38
    • 39
    • 40
    • 41

    在这里插入图片描述
    在这里插入图片描述
    在这里插入图片描述

    4.路由模式(Routing)

    • 1.路由模式即一个发送指定的接收,使用的是direct交换机
    • 2.路由模式和简单以及工作模式的区别
      • 1.简单模式使用默认direct交换机:生产者发送消息到默认交换机,然后再发送给指定名称的队列,最后再由队列发送给一个消费者
      • 2.工作模式使用默认direct交换机:生产者发送消息到默认交换机,然后再发送给指定名称的队列,最后再由队列发送给多个消费者争抢
      • 3.路由模式使用手动创建direct交换机:生产者发送消息到指定交换机,然后交换机根据Routing Key发送给指定的队列,最后再由队列发送给消费者
    • 3.如果路由模式的Routing Key都一样则类似发布订阅模式
    1.生产者
    package com.rabbit.work.routing;
    
    import com.rabbit.util.RabbitMQUtils;
    import com.rabbitmq.client.BuiltinExchangeType;
    import com.rabbitmq.client.Channel;
    import com.rabbitmq.client.MessageProperties;
    
    import java.util.Scanner;
    
    /**
    * 发消息给交换机
    */
    public class WorkProducer {
       // 交换机名称
       private final static String EXCHANGE_NAME = "direct_message";
       public static void main(String[] args) throws Exception {
           Channel channel = RabbitMQUtils.getChannel();
           // 声明交换机
           channel.exchangeDeclare(EXCHANGE_NAME, BuiltinExchangeType.DIRECT);
    
           Scanner scanner = new Scanner(System.in);
           System.out.println("请输入信息+路由键:");
           while (scanner.hasNext()) {
               String whole = scanner.next();
               String[] messageKey = whole.split(",");
               // 发送消息给交换机并将消息持久化
               channel.basicPublish(EXCHANGE_NAME,messageKey[1], MessageProperties.PERSISTENT_TEXT_PLAIN,messageKey[0].getBytes());
               System.out.println("消息发送完毕: " + messageKey[0]);
           }
       }
    }
    
    • 1
    • 2
    • 3
    • 4
    • 5
    • 6
    • 7
    • 8
    • 9
    • 10
    • 11
    • 12
    • 13
    • 14
    • 15
    • 16
    • 17
    • 18
    • 19
    • 20
    • 21
    • 22
    • 23
    • 24
    • 25
    • 26
    • 27
    • 28
    • 29
    • 30
    • 31
    2.消费者
    package com.rabbit.work.routing;
    
    import com.rabbit.util.RabbitMQUtils;
    import com.rabbitmq.client.BuiltinExchangeType;
    import com.rabbitmq.client.CancelCallback;
    import com.rabbitmq.client.Channel;
    import com.rabbitmq.client.DeliverCallback;
    
    /**
    * 接收状态消息
    */
    public class WorkConsumer_01 {
       // 交换机名称
       private final static String EXCHANGE_NAME = "direct_message";
       // 队列名称
       private final static String STATUE_QUEUE_NAME = "statue_message";
    
       public static void main(String[] args) throws Exception {
           Channel channel = RabbitMQUtils.getChannel();
           // 声明交换机
           channel.exchangeDeclare(EXCHANGE_NAME, BuiltinExchangeType.DIRECT);
           // 声明一个持久化队列
           channel.queueDeclare(STATUE_QUEUE_NAME,true,false,false,null);
           // 绑定交换机与队列
           channel.queueBind(STATUE_QUEUE_NAME,EXCHANGE_NAME,"success");
           channel.queueBind(STATUE_QUEUE_NAME,EXCHANGE_NAME,"error");
    
           System.out.println("WorkConsumer_01等待接收消息....");
    
           DeliverCallback deliverCallback = (consumerTag, message) -> {
    
               System.out.println("WorkConsumer_01接收到的消息为:" + new String(message.getBody()) + " :: 路由键:" + message.getEnvelope().getRoutingKey());
               /**
                *  肯定确认
                *  1.消息标记,每一个消息都有一个独立的标记
                *  2.是否批量应答未应答消息
                */
               channel.basicAck(message.getEnvelope().getDeliveryTag(),false);
           };
           CancelCallback cancelCallback = (consumerTag) -> System.out.println("消息消费被中断");
           // 设置不公平分发,默认值为0
           channel.basicQos(1);
           // 采用手动应答
           channel.basicConsume(STATUE_QUEUE_NAME,false,deliverCallback,cancelCallback);
       }
    }
    
    • 1
    • 2
    • 3
    • 4
    • 5
    • 6
    • 7
    • 8
    • 9
    • 10
    • 11
    • 12
    • 13
    • 14
    • 15
    • 16
    • 17
    • 18
    • 19
    • 20
    • 21
    • 22
    • 23
    • 24
    • 25
    • 26
    • 27
    • 28
    • 29
    • 30
    • 31
    • 32
    • 33
    • 34
    • 35
    • 36
    • 37
    • 38
    • 39
    • 40
    • 41
    • 42
    • 43
    • 44
    • 45
    • 46
    package com.rabbit.work.routing;
    
    import com.rabbit.util.RabbitMQUtils;
    import com.rabbitmq.client.BuiltinExchangeType;
    import com.rabbitmq.client.CancelCallback;
    import com.rabbitmq.client.Channel;
    import com.rabbitmq.client.DeliverCallback;
    
    /**
    * 接收内容消息
    */
    public class WorkConsumer_02 {
       // 交换机名称
       private final static String EXCHANGE_NAME = "direct_message";
       // 队列名称
       private final static String CONTEXT_QUEUE_NAME = "context_message";
    
       public static void main(String[] args) throws Exception {
           Channel channel = RabbitMQUtils.getChannel();
           // 声明交换机
           channel.exchangeDeclare(EXCHANGE_NAME, BuiltinExchangeType.DIRECT);
           // 声明一个持久化队列
           channel.queueDeclare(CONTEXT_QUEUE_NAME,true,false,false,null);
           // 绑定交换机与队列
           channel.queueBind(CONTEXT_QUEUE_NAME,EXCHANGE_NAME,"context");
    
           System.out.println("WorkConsumer_02等待接收消息....");
    
           DeliverCallback deliverCallback = (consumerTag, message) -> {
               System.out.println("WorkConsumer_02接收到的消息为:" + new String(message.getBody()) + " :: 路由键:" + message.getEnvelope().getRoutingKey());
               /**
                *  肯定确认
                *  1.消息标记,每一个消息都有一个独立的标记
                *  2.是否批量应答未应答消息
                */
               channel.basicAck(message.getEnvelope().getDeliveryTag(),false);
           };
           CancelCallback cancelCallback = (consumerTag) -> System.out.println("消息消费被中断");
           // 设置不公平分发,默认值为0
           channel.basicQos(1);
           // 采用手动应答
           channel.basicConsume(CONTEXT_QUEUE_NAME,false,deliverCallback,cancelCallback);
       }
    }
    
    • 1
    • 2
    • 3
    • 4
    • 5
    • 6
    • 7
    • 8
    • 9
    • 10
    • 11
    • 12
    • 13
    • 14
    • 15
    • 16
    • 17
    • 18
    • 19
    • 20
    • 21
    • 22
    • 23
    • 24
    • 25
    • 26
    • 27
    • 28
    • 29
    • 30
    • 31
    • 32
    • 33
    • 34
    • 35
    • 36
    • 37
    • 38
    • 39
    • 40
    • 41
    • 42
    • 43
    • 44

    在这里插入图片描述
    在这里插入图片描述
    在这里插入图片描述

    5.主题模式(Topics)

    • 1.主题模式及一个发送匹配的接收,其使用的是topics交换机
    • 2.主题模式可以对路由键进行模糊匹配,将交换机接收到的消息路由到匹配成功的队列上
    1.模糊匹配
    • 1.案列
      • 1.*.orange.*:绑定的是中间是orange且3个单词的单词列表
      • 2.*.*.rabbit:绑定的是最后是rabbit且3个单词的单词列表
      • 3.lazy.#:绑定的是开始是lazy且长度不限的单词列表
    • 2.注意
      • 1.当绑定键是#,该绑定队列将接收所有数据,类似fanout
      • 2.当绑定键中没有#*出现,该队列绑定类型类似direct
    2.生产者
    package com.rabbit.work.topics;
    
    import com.rabbit.util.RabbitMQUtils;
    import com.rabbitmq.client.BuiltinExchangeType;
    import com.rabbitmq.client.Channel;
    import com.rabbitmq.client.MessageProperties;
    
    import java.util.Scanner;
    
    /**
    * 发消息给交换机
    */
    public class WorkProducer {
       // 交换机名称
       private final static String EXCHANGE_NAME = "topics_message";
       public static void main(String[] args) throws Exception {
           Channel channel = RabbitMQUtils.getChannel();
           // 声明交换机
           channel.exchangeDeclare(EXCHANGE_NAME, BuiltinExchangeType.TOPIC);
    
           Scanner scanner = new Scanner(System.in);
           System.out.println("请输入信息+路由键:");
           while (scanner.hasNext()) {
               String whole = scanner.next();
               String[] messageKey = whole.split(",");
               // 发送消息给交换机并将消息持久化
               channel.basicPublish(EXCHANGE_NAME,messageKey[1], MessageProperties.PERSISTENT_TEXT_PLAIN,messageKey[0].getBytes());
               System.out.println("消息发送完毕: " + messageKey[0]);
           }
       }
    }
    
    • 1
    • 2
    • 3
    • 4
    • 5
    • 6
    • 7
    • 8
    • 9
    • 10
    • 11
    • 12
    • 13
    • 14
    • 15
    • 16
    • 17
    • 18
    • 19
    • 20
    • 21
    • 22
    • 23
    • 24
    • 25
    • 26
    • 27
    • 28
    • 29
    • 30
    • 31
    3.消费者
    package com.rabbit.work.topics;
    
    import com.rabbit.util.RabbitMQUtils;
    import com.rabbitmq.client.BuiltinExchangeType;
    import com.rabbitmq.client.CancelCallback;
    import com.rabbitmq.client.Channel;
    import com.rabbitmq.client.DeliverCallback;
    
    /**
    * 接收状态消息
    */
    public class WorkConsumer_01 {
       // 交换机名称
       private final static String EXCHANGE_NAME = "topics_message";
       // 队列名称
       private final static String STATUE_QUEUE_NAME = "queue_statue_message";
    
       public static void main(String[] args) throws Exception {
           Channel channel = RabbitMQUtils.getChannel();
           // 声明交换机
           channel.exchangeDeclare(EXCHANGE_NAME, BuiltinExchangeType.TOPIC);
           // 声明一个持久化队列
           channel.queueDeclare(STATUE_QUEUE_NAME,true,false,false,null);
           // 绑定交换机与队列
           channel.queueBind(STATUE_QUEUE_NAME,EXCHANGE_NAME,"success.#");
           channel.queueBind(STATUE_QUEUE_NAME,EXCHANGE_NAME,"error.*");
    
    
           System.out.println("WorkConsumer_01等待接收消息....");
    
           DeliverCallback deliverCallback = (consumerTag, message) -> {
    
               System.out.println("WorkConsumer_01接收到的消息为:" + new String(message.getBody()) + " :: 路由键:" + message.getEnvelope().getRoutingKey());
               /**
                *  肯定确认
                *  1.消息标记,每一个消息都有一个独立的标记
                *  2.是否批量应答未应答消息
                */
               channel.basicAck(message.getEnvelope().getDeliveryTag(),false);
           };
           CancelCallback cancelCallback = (consumerTag) -> System.out.println("消息消费被中断");
           // 设置不公平分发,默认值为0
           channel.basicQos(1);
           // 采用手动应答
           channel.basicConsume(STATUE_QUEUE_NAME,false,deliverCallback,cancelCallback);
       }
    }
    
    • 1
    • 2
    • 3
    • 4
    • 5
    • 6
    • 7
    • 8
    • 9
    • 10
    • 11
    • 12
    • 13
    • 14
    • 15
    • 16
    • 17
    • 18
    • 19
    • 20
    • 21
    • 22
    • 23
    • 24
    • 25
    • 26
    • 27
    • 28
    • 29
    • 30
    • 31
    • 32
    • 33
    • 34
    • 35
    • 36
    • 37
    • 38
    • 39
    • 40
    • 41
    • 42
    • 43
    • 44
    • 45
    • 46
    • 47
    package com.rabbit.work.topics;
    
    import com.rabbit.util.RabbitMQUtils;
    import com.rabbitmq.client.BuiltinExchangeType;
    import com.rabbitmq.client.CancelCallback;
    import com.rabbitmq.client.Channel;
    import com.rabbitmq.client.DeliverCallback;
    
    /**
    * 接收内容消息
    */
    public class WorkConsumer_02 {
       // 交换机名称
       private final static String EXCHANGE_NAME = "topics_message";
       // 队列名称
       private final static String CONTEXT_QUEUE_NAME = "queue_context_message";
    
       public static void main(String[] args) throws Exception {
           Channel channel = RabbitMQUtils.getChannel();
           // 声明交换机
           channel.exchangeDeclare(EXCHANGE_NAME, BuiltinExchangeType.TOPIC);
           // 声明一个持久化队列
           channel.queueDeclare(CONTEXT_QUEUE_NAME,true,false,false,null);
           // 绑定交换机与队列
           channel.queueBind(CONTEXT_QUEUE_NAME,EXCHANGE_NAME,"context.*");
    
           System.out.println("WorkConsumer_02等待接收消息....");
    
           DeliverCallback deliverCallback = (consumerTag, message) -> {
               System.out.println("WorkConsumer_02接收到的消息为:" + new String(message.getBody()) + " :: 路由键:" + message.getEnvelope().getRoutingKey());
               /**
                *  肯定确认
                *  1.消息标记,每一个消息都有一个独立的标记
                *  2.是否批量应答未应答消息
                */
               channel.basicAck(message.getEnvelope().getDeliveryTag(),false);
           };
           CancelCallback cancelCallback = (consumerTag) -> System.out.println("消息消费被中断");
           // 设置不公平分发,默认值为0
           channel.basicQos(1);
           // 采用手动应答
           channel.basicConsume(CONTEXT_QUEUE_NAME,false,deliverCallback,cancelCallback);
       }
    }
    
    • 1
    • 2
    • 3
    • 4
    • 5
    • 6
    • 7
    • 8
    • 9
    • 10
    • 11
    • 12
    • 13
    • 14
    • 15
    • 16
    • 17
    • 18
    • 19
    • 20
    • 21
    • 22
    • 23
    • 24
    • 25
    • 26
    • 27
    • 28
    • 29
    • 30
    • 31
    • 32
    • 33
    • 34
    • 35
    • 36
    • 37
    • 38
    • 39
    • 40
    • 41
    • 42
    • 43
    • 44

    在这里插入图片描述
    在这里插入图片描述
    在这里插入图片描述

    6.发布确认模式(Publisher Confirms)

    • 1.由于某些原因,导致RabbitMQ重启或不可用,重启或不可用期间生产者消息投递失败,导致消息丢失,需手动处理和恢复
    • 2.极端情况(RabbitMQ集群不可用)可通过发布确认模式保证RabbitMQ的可靠性投递
    • 3.交换机队列都有可能发生问题导致生产者消息丢失,首先需要对发送的消息进行缓存,如果发送失败则可以通过缓存对未成功发送的消息重新投递
      在这里插入图片描述
    1.开启发布确认模式
    • 1.发布确认模式默认是关闭的,需要手动开启
    • 2.下述发布确认机制中可以看到测试中开启发布模式的方式,但开发中一般通过配置文件开启发布确认模式
    • 3.配置文件设置开启发布确认模式,如果不开启则不会生效
      spring.rabbitmq.publisher-confirm-type=correlated
      // 参数选项
      NONE # 禁用发布确认模式,是默认值
      CORRELATED	# 发布消息成功到交换器后会触发回调方法 一般采用该方法
      SIMPLE # 经测试有两种效果
      # 其一和 CORRELATED 值一样会触发回调方法
      # 其二发布消息成功后使用 rabbitTemplate 调用 waitForConfirms 或 waitForConfirmsOrDie 方法等待 broker 节点返回发送结果,根据返回结果来判定下一步的逻辑;注意:waitForConfirmsOrDie 方法如果返回 false 则会关闭 channel,则接下来无法发送消息到 broker
      
      • 1
      • 2
      • 3
      • 4
      • 5
      • 6
      • 7
      在这里插入图片描述
    2.配置类

    在这里插入图片描述

    package com.rabbit.config;
    
    import org.springframework.amqp.core.*;
    import org.springframework.context.annotation.Bean;
    import org.springframework.context.annotation.Configuration;
    
    /**
    * @description: 发布确认模式
    */
    @Configuration
    public class ConfirmConfig {
    
       private final static String CONFIRM_QUEUE = "confirm_queue";
       private final static String CONFIRM_EXCHANGE = "confirm_exchange";
       private final static String CONFIRM_ROUTING_KEY = "confirm";
    
       @Bean
       public Queue confirmQueue(){
           return new Queue(CONFIRM_QUEUE);
       }
    
       @Bean
       public DirectExchange confirmExchange(){
           return new ExchangeBuilder(CONFIRM_EXCHANGE, ExchangeTypes.DIRECT).durable(true).build();
       }
    
       @Bean
       public Binding confirmExchangeBindingQueue(){
           return BindingBuilder.bind(confirmQueue()).to(confirmExchange()).with(CONFIRM_ROUTING_KEY);
       }
    }
    
    • 1
    • 2
    • 3
    • 4
    • 5
    • 6
    • 7
    • 8
    • 9
    • 10
    • 11
    • 12
    • 13
    • 14
    • 15
    • 16
    • 17
    • 18
    • 19
    • 20
    • 21
    • 22
    • 23
    • 24
    • 25
    • 26
    • 27
    • 28
    • 29
    • 30
    • 31
    3.生产者
    @GetMapping("sendConfirm/{message}")
       public void sendConfirm(@PathVariable String message) {
           System.out.println("---------------------------------------------------------------------------------------------------------------------------------------------------------");
           CorrelationData correlationData = new CorrelationData("1");
           correlationData.setReturnedMessage(new Message(message.getBytes(),null));
           log.info("发送一条信息给confirm_queue队列: {}", message);
           rabbitTemplate.convertAndSend(ConfirmConfig.CONFIRM_EXCHANGE, ConfirmConfig.CONFIRM_ROUTING_KEY", message,correlationData);
       }
    
    • 1
    • 2
    • 3
    • 4
    • 5
    • 6
    • 7
    • 8
    4.消费者
    package com.rabbit.listener;
    
    import com.rabbit.config.ConfirmConfig;
    import com.rabbitmq.client.Channel;
    import lombok.extern.slf4j.Slf4j;
    import org.springframework.amqp.core.Message;
    import org.springframework.amqp.rabbit.annotation.RabbitListener;
    import org.springframework.stereotype.Component;
    
    import java.util.Date;
    
    @Slf4j
    @Component
    public class ReceiveMessageListener {
    
       @RabbitListener(queues = ConfirmConfig.CONFIRM_QUEUE)
       public void receiveConfirmMessage(Message message, Channel channel){
           String msg = new String(message.getBody());
           log.info("收到confirm_queue队列信息: {}", msg);
       }
    }
    
    • 1
    • 2
    • 3
    • 4
    • 5
    • 6
    • 7
    • 8
    • 9
    • 10
    • 11
    • 12
    • 13
    • 14
    • 15
    • 16
    • 17
    • 18
    • 19
    • 20
    • 21
    5.回调处理
    package com.rabbit.service.impl;
    
    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.Service;
    
    import javax.annotation.PostConstruct;
    import javax.annotation.Resource;
    
    /**
    * 回调接口
    * 注意:该实现类实现的是内部接口,所以需要将该类注入到内部类的指定接口中
    */
    @Slf4j
    @Service // 或@Component
    public class MyCallBack implements RabbitTemplate.ConfirmCallback{
    
       @Resource
       private RabbitTemplate rabbitTemplate;
    
       @PostConstruct
       public void init(){
           // 1.注入到指定类的指定接口上
           // 2.顺序不能更改
           // 3.首先加载MyCallBack,然后加载RabbitTemplate,最后加载该方法将该实现类注入到RabbitTemplate的指定接口上
           rabbitTemplate.setConfirmCallback(this);
       }
    
       /**
        * 交换机确认回调方法
        * 1.发消息 交换机收到 回调
        *  1.1 correlationData 保存回调消息的ID及相关信息
        *  1.2 交换机收到消息 ack = true
        *  1.3 cause null
        * 2.发消息 交换机接收失败 回调
        *  2.1 correlationData 保存回调消息的ID及相关信息
        *  2.2 交换机未收到消息 ack = false
        *  2.3 cause 失败的原因
        */
       @Override
       public void confirm(CorrelationData correlationData, boolean ack, String cause) {
           String id = correlationData.getId() != null ? correlationData.getId() : "";
           Message returnedMessage = correlationData.getReturnedMessage() != null ? correlationData.getReturnedMessage() : null;
           if(ack) {
               log.info("交换机已经收到id为: {} 的消息: {}",id, new String(returnedMessage.getBody()));
           }else {
               log.info("交换机未收到id为: {} 的消息: {},由于: {}",id,new String(returnedMessage.getBody()),cause);
           }
       }
    }
    
    • 1
    • 2
    • 3
    • 4
    • 5
    • 6
    • 7
    • 8
    • 9
    • 10
    • 11
    • 12
    • 13
    • 14
    • 15
    • 16
    • 17
    • 18
    • 19
    • 20
    • 21
    • 22
    • 23
    • 24
    • 25
    • 26
    • 27
    • 28
    • 29
    • 30
    • 31
    • 32
    • 33
    • 34
    • 35
    • 36
    • 37
    • 38
    • 39
    • 40
    • 41
    • 42
    • 43
    • 44
    • 45
    • 46
    • 47
    • 48
    • 49
    • 50
    • 51
    • 52
    • 1.回调处理针对的是交换机未接收消息导致消息丢失的处理方案
    6.回退处理
    package com.rabbit.service.impl;
    
    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.Service;
    
    import javax.annotation.PostConstruct;
    import javax.annotation.Resource;
    
    /**
    * 回调接口
    * 注意:该实现类实现的是内部接口,所以需要将该类注入到内部类的指定接口中
    */
    @Slf4j
    @Service // 或@Component
    public class MyCallBack implements >RabbitTemplate.ConfirmCallback,RabbitTemplate.ReturnCallback{
    
       @Resource
       private RabbitTemplate rabbitTemplate;
    
       @PostConstruct
       public void init(){
           // 1.注入回调具体处理类到指定类的指定接口上
           // 2.顺序不能更改
           // 3.首先加载MyCallBack,然后加载RabbitTemplate,最后加载该方法将该实现类注入到RabbitTemplate的指定接口上
           rabbitTemplate.setConfirmCallback(this);
           // 1.注入回退具体处理类到指定类的指定接口上
           rabbitTemplate.setReturnCallback(this);
       }
    
       /**
        * 交换机确认回调方法
        * 1.发消息 交换机收到 回调
        *  1.1 correlationData 保存回调消息的ID及相关信息
        *  1.2 交换机收到消息 ack = true
        *  1.3 cause null
        * 2.发消息 交换机接收失败 回调
        *  2.1 correlationData 保存回调消息的ID及相关信息
        *  2.2 交换机未收到消息 ack = false
        *  2.3 cause 失败的原因
        */
       @Override
       public void confirm(CorrelationData correlationData, boolean ack, String cause) {
        //记录日志、发送邮件通知、落库定时任务扫描重发
           String id = correlationData.getId() != null ? correlationData.getId() : "";
           Message returnedMessage = correlationData.getReturnedMessage() != null ? correlationData.getReturnedMessage() : null;
           if(ack) {
               log.info("交换机已经收到id为: {} 的消息: {}",id, new String(returnedMessage.getBody()));
           }else {
               log.info("交换机未收到id为: {} 的消息: {},由于: {}",id,new String(returnedMessage.getBody()),cause);
           }
       }
    
       @Override
       public void returnedMessage(Message message, int replyCode, String replyText, String exchange, String routingKey) {
        //记录日志、发送邮件通知、落库定时任务扫描重发
           log.info("消息:{}被服务器退回,退回原因:{}, 交换机是:{}, 路由 key:{}",
                   new String(message.getBody()),replyText, exchange, routingKey);
       }
    }
    
    • 1
    • 2
    • 3
    • 4
    • 5
    • 6
    • 7
    • 8
    • 9
    • 10
    • 11
    • 12
    • 13
    • 14
    • 15
    • 16
    • 17
    • 18
    • 19
    • 20
    • 21
    • 22
    • 23
    • 24
    • 25
    • 26
    • 27
    • 28
    • 29
    • 30
    • 31
    • 32
    • 33
    • 34
    • 35
    • 36
    • 37
    • 38
    • 39
    • 40
    • 41
    • 42
    • 43
    • 44
    • 45
    • 46
    • 47
    • 48
    • 49
    • 50
    • 51
    • 52
    • 53
    • 54
    • 55
    • 56
    • 57
    • 58
    • 59
    • 60
    • 61
    • 62
    • 1.回退处理针对的交换机无法路由到队列导致消息丢失的处理方案
    • 2.回退处理也需要在配置文件中配置
      在这里插入图片描述
    7.测试结果
    • 1.回调处理
      在这里插入图片描述
    • 2.回退处理
      在这里插入图片描述
    8.备份交换机
    • 1.当交换机接收到不可路由消息时会把该消息转发到备份交换机中,由备份交换机来进行转发和处理
    • 2.通常备份交换机的类型为 Fanout ,能把所有消息都投递到与其绑定的队列中
      在这里插入图片描述
    1.配置类
    package com.rabbit.config;
    
    import org.springframework.amqp.core.*;
    import org.springframework.context.annotation.Bean;
    import org.springframework.context.annotation.Configuration;
    
    /**
    * @description: 发布确认模式
    */
    @Configuration
    public class ConfirmConfig {
    
       public final static String CONFIRM_QUEUE = "confirm_queue";
       public final static String CONFIRM_EXCHANGE = "confirm_exchange";
       public final static String CONFIRM_ROUTING_KEY = "confirm";
       public final static String BACKUP_EXCHANGE = "backup_exchange";
       public final static String BACKUP_QUEUE = "backup_queue";
       public final static String WARNING_QUEUE = "warning_queue";
    
       @Bean
       public Queue confirmQueue(){
           return new Queue(CONFIRM_QUEUE);
       }
    
       @Bean
       public Queue backupQueue(){
           return QueueBuilder.durable(BACKUP_QUEUE).build();
       }
    
       @Bean
       public Queue warningQueue(){
           return new Queue(WARNING_QUEUE);
       }
    
       @Bean
       public DirectExchange confirmExchange(){
           return new ExchangeBuilder(CONFIRM_EXCHANGE, ExchangeTypes.DIRECT).durable(true)
                   // 设置该交换机的备份交换机
                   .withArgument("alternate-exchange",BACKUP_EXCHANGE).build();
       }
    
       @Bean
       public FanoutExchange backupExchange(){
           return ExchangeBuilder.fanoutExchange(BACKUP_EXCHANGE).build();
       }
    
       @Bean
       public Binding confirmExchangeBindingQueue(){
           return BindingBuilder.bind(confirmQueue()).to(confirmExchange()).with(CONFIRM_ROUTING_KEY);
       }
    
       @Bean
       public Binding backupExchangeBindingQueue(){
           return BindingBuilder.bind(backupQueue()).to(backupExchange());
       }
    
       @Bean
       public Binding warningExchangeBindingQueue(){
           return BindingBuilder.bind(warningQueue()).to(backupExchange());
       }
    }
    
    • 1
    • 2
    • 3
    • 4
    • 5
    • 6
    • 7
    • 8
    • 9
    • 10
    • 11
    • 12
    • 13
    • 14
    • 15
    • 16
    • 17
    • 18
    • 19
    • 20
    • 21
    • 22
    • 23
    • 24
    • 25
    • 26
    • 27
    • 28
    • 29
    • 30
    • 31
    • 32
    • 33
    • 34
    • 35
    • 36
    • 37
    • 38
    • 39
    • 40
    • 41
    • 42
    • 43
    • 44
    • 45
    • 46
    • 47
    • 48
    • 49
    • 50
    • 51
    • 52
    • 53
    • 54
    • 55
    • 56
    • 57
    • 58
    • 59
    • 60
    • 61
    • 1.注意:重新启动项目的时候需要把原来的 confirm_exchange 删除,因为修改了其绑定属性,否则报以下错
      在这里插入图片描述
    2.生产者
       @GetMapping("sendConfirm/{message}")
       public void sendConfirm(@PathVariable String message) {
           System.out.println("---------------------------------------------------------------------------------------------------------------------------------------------------------");
           CorrelationData correlationData = new CorrelationData("1");
           correlationData.setReturnedMessage(new Message(message.getBytes(),null));
           log.info("发送一条信息给confirm_queue队列: {}", message);
           rabbitTemplate.convertAndSend(ConfirmConfig.CONFIRM_EXCHANGE, ConfirmConfig.CONFIRM_ROUTING_KEY+"123", message,correlationData);
       }
    
    • 1
    • 2
    • 3
    • 4
    • 5
    • 6
    • 7
    • 8
    3.消费者
    package com.rabbit.listener;
    
    import com.rabbit.config.ConfirmConfig;
    import com.rabbitmq.client.Channel;
    import lombok.extern.slf4j.Slf4j;
    import org.springframework.amqp.core.Message;
    import org.springframework.amqp.rabbit.annotation.RabbitListener;
    import org.springframework.stereotype.Component;
    
    import java.util.Date;
    
    @Slf4j
    @Component
    public class ReceiveMessageListener {
    
       @RabbitListener(queues = ConfirmConfig.CONFIRM_QUEUE)
       public void receiveConfirmMessage(Message message, Channel channel){
           String msg = new String(message.getBody());
           log.info("收到confirm_queue队列信息: {}", msg);
       }
    
       @RabbitListener(queues = ConfirmConfig.BACKUP_QUEUE)
       public void receiveBackupMessage(Message message, Channel channel){
           String msg = new String(message.getBody());
           log.info("备份队列收到不可路由信息: {}", msg);
       }
    
       @RabbitListener(queues = ConfirmConfig.WARNING_QUEUE)
       public void receiveWarningMessage(Message message, Channel channel){
           String msg = new String(message.getBody());
           log.info("报警队列收到不可路由信息: {}", msg);
       }
    }
    
    • 1
    • 2
    • 3
    • 4
    • 5
    • 6
    • 7
    • 8
    • 9
    • 10
    • 11
    • 12
    • 13
    • 14
    • 15
    • 16
    • 17
    • 18
    • 19
    • 20
    • 21
    • 22
    • 23
    • 24
    • 25
    • 26
    • 27
    • 28
    • 29
    • 30
    • 31
    • 32
    • 33
    4.回退处理
    • 1.回退处理同上
    4.测试结果

    在这里插入图片描述

    • 1.回退处理备份交换机一起使用时,经上面结果对比是备份交换机优先级高

    6.消息应答机制

    • 1.消费者完成任务可能需要一段时间,如果消费者处理任务并仅完成部分突然挂掉则任务丢失
    • 2.RabbitMQ一旦向消费者传递了一条消息,便立即将该消息标记为删除
    • 3.该情况下如果突然有消费者挂掉,将丢失正在处理的消息以及后续发送给该消费者的消息,因为其无法接收
    • 4.为了保证消息在发送过程中不丢失,RabbitMQ引入消息应答机制
    • 5.消息应答机制:消费者在接收到消息并处理该消息之后,需告诉RabbitMQ其已经处理完成,RabbitMQ才可以把该消息删除

    1.自动应答

    • 1.消息发送后立即被认为已传送成功,并把该消息删除
    • 2.该模式仅适合满足一定条件的环境,不适合极端环境
    • 3.该模式如果消息在接收到之前,消费者出现连接或Channel关闭,会消息丢失
    • 4.该模式消费者如果接收过载的消息,即没有对传递的消息数量进行限制,导致消息的积压,最终可能使得内存耗尽,该模式没有批量应答

    2.手动应答

    1.三种方式
    • 1.Channel.basicAck()肯定确认RabbitMQ确认成功处理该消息,可将其从队列删除
      /**
        * 确认一个或多个接收到的消息
        * @param deliveryTag 标签
        * @param multiple True表示确认所有的消息;False表示只确认当前消息
        */
      void basicAck(long deliveryTag, boolean multiple) throws IOException;
      
      • 1
      • 2
      • 3
      • 4
      • 5
      • 6
    • 2.Channel.basicNack()否定确认RabbitMQ确认未成功处理该消息,不可将其从队列删除
       /**
        * 拒绝一个或多个接收到的消息
        * @param deliveryTag 标签
        * @param multiple True表示拒绝所有的消息;False表示只拒绝当前消息
        * the supplied delivery tag; false to reject just the supplied
        * delivery tag.
        * @param requeue true表示被拒绝的消息应该被重新排序而不是被丢弃为死信
        */
       void basicNack(long deliveryTag, boolean multiple, boolean requeue)
               throws IOException;
      
      • 1
      • 2
      • 3
      • 4
      • 5
      • 6
      • 7
      • 8
      • 9
      • 10
    • 3.Channel.basicReject()否定确认,不处理该消息直接拒绝,并可以将其丢弃
       /**
        * 拒绝消息
        * @param deliveryTag 标签
        * @param requeue true表示被拒绝的消息应该被重新排序而不是被丢弃为死信
        */
       void basicReject(long deliveryTag, boolean requeue) throws IOException;
      
      • 1
      • 2
      • 3
      • 4
      • 5
      • 6
    2.批量应答
    • 1.参数:boolean multiple
    • 2.如果为true:会批量应答该Channel上其他未应答的消息
    • 3.如果为false:只会应答当前消息
      在这里插入图片描述
    • 4.批量应答可减少网络拥堵,但一般不建议使用,因为可能会丢失消息
      channel.basicAck(deiveryTag,true)
      
      • 1
    3.消息自动重新入队
    • 1.消费者由于某些原因失去连接(其通道已关闭,连接已关闭或TCP连接丢失),导致消费者未发送ACK确认,RabbitMQ将认为该消息未完全处理,并将对其重新排队
    • 2.此时其他消费者可以继续处理该重新入队的消息,即使某个消费者偶尔死亡,也可确保不会丢失任何消息
      在这里插入图片描述
    4.消息手动应答代码
    • 1.默认消息采用的是自动应答,所以要想实现消息消费过程不丢失,需把自动应答改为手动应答
    • 2.手动应答代码处于消费者端
    • 3.消息消费未确认时应将消息放回队列等待消费者重新消费

    生产者

    package com.rabbit.work.ack_work;
    
    import com.rabbit.util.RabbitMQUtils;
    import com.rabbitmq.client.Channel;
    
    import java.util.Scanner;
    
    public class WorkProducer {
       private final static String QUEUE_NAME = "ack_queue";
       public static void main(String[] args) throws Exception {
           Channel channel = RabbitMQUtils.getChannel();
           channel.queueDeclare(QUEUE_NAME,false,false,false,null);
           Scanner scanner = new Scanner(System.in);
           System.out.println("请输入信息:");
           while (scanner.hasNext()) {
               String message = scanner.next();
               channel.basicPublish("",QUEUE_NAME,null,message.getBytes());
               System.out.println("消息发送完毕: " + message);
           }
       }
    }
    
    • 1
    • 2
    • 3
    • 4
    • 5
    • 6
    • 7
    • 8
    • 9
    • 10
    • 11
    • 12
    • 13
    • 14
    • 15
    • 16
    • 17
    • 18
    • 19
    • 20
    • 21

    消费者

    package com.rabbit.work.ack_work;
    
    import com.rabbit.util.RabbitMQUtils;
    import com.rabbit.util.SleepUtils;
    import com.rabbitmq.client.CancelCallback;
    import com.rabbitmq.client.Channel;
    import com.rabbitmq.client.DeliverCallback;
    
    public class WorkConsumer_01 {
       private final static String QUEUE_NAME = "ack_queue";
    
       public static void main(String[] args) throws Exception {
           Channel channel = RabbitMQUtils.getChannel();
    
           System.out.println("WorkConsumer_01等待接收消息处理消息时间短");
    
           DeliverCallback deliverCallback = (consumerTag, message) -> {
               String body = new String(message.getBody());
               SleepUtils.sleep(1);
               System.out.println("接收到的消息为:" + body);
               /**
                *  肯定确认
                *  1.消息标记,每一个消息都有一个独立的标记
                *  2.是否批量应答未应答消息
                */
               channel.basicAck(message.getEnvelope().getDeliveryTag(),false);
           };
           CancelCallback cancelCallback = (consumerTag) -> {
               System.out.println("消息消费被中断");
           };
           //采用手动应答
           boolean autoAck = false;
           >channel.basicConsume(QUEUE_NAME,autoAck,deliverCallback,cancelCallback);
       }
    }
    
    • 1
    • 2
    • 3
    • 4
    • 5
    • 6
    • 7
    • 8
    • 9
    • 10
    • 11
    • 12
    • 13
    • 14
    • 15
    • 16
    • 17
    • 18
    • 19
    • 20
    • 21
    • 22
    • 23
    • 24
    • 25
    • 26
    • 27
    • 28
    • 29
    • 30
    • 31
    • 32
    • 33
    • 34
    • 35
    package com.rabbit.work.ack_work;
    
    import com.rabbit.util.RabbitMQUtils;
    import com.rabbit.util.SleepUtils;
    import com.rabbitmq.client.CancelCallback;
    import com.rabbitmq.client.Channel;
    import com.rabbitmq.client.DeliverCallback;
    
    public class WorkConsumer_02 {
       private final static String QUEUE_NAME = "ack_queue";
    
       public static void main(String[] args) throws Exception {
           Channel channel = RabbitMQUtils.getChannel();
    
           System.out.println("WorkConsumer_02等待接收消息处理消息时间长");
    
           DeliverCallback deliverCallback = (consumerTag, message) -> {
               String body = new String(message.getBody());
               SleepUtils.sleep(30);
               System.out.println("接收到的消息为:" + body);
               /**
                *  肯定确认
                *  1.消息标记,每一个消息都有一个独立的标记
                *  2.是否批量应答未应答消息
                */
               channel.basicAck(message.getEnvelope().getDeliveryTag(),false);
           };
           CancelCallback cancelCallback = (consumerTag) -> {
               System.out.println("消息消费被中断");
           };
           //采用手动应答
           boolean autoAck = false;
           >channel.basicConsume(QUEUE_NAME,autoAck,deliverCallback,cancelCallback);
       }
    }
    
    • 1
    • 2
    • 3
    • 4
    • 5
    • 6
    • 7
    • 8
    • 9
    • 10
    • 11
    • 12
    • 13
    • 14
    • 15
    • 16
    • 17
    • 18
    • 19
    • 20
    • 21
    • 22
    • 23
    • 24
    • 25
    • 26
    • 27
    • 28
    • 29
    • 30
    • 31
    • 32
    • 33
    • 34
    • 35

    工具类

    package com.rabbit.util;
    
    public class SleepUtils {
    
       public static void sleep(int second){
           try {
               Thread.sleep(1000 * second);
           } catch (InterruptedException e) {
               e.printStackTrace();
           }
       }
    }
    
    • 1
    • 2
    • 3
    • 4
    • 5
    • 6
    • 7
    • 8
    • 9
    • 10
    • 11
    • 12

    在这里插入图片描述
    在这里插入图片描述
    在这里插入图片描述

    3.使用方式

    1.配置文件中开启手动应答

    在这里插入图片描述

    2.生产者
    @GetMapping("sendConfirm/{message}")
       public void sendConfirm(@PathVariable String message) {
           System.out.println("---------------------------------------------------------------------------------------------------------------------------------------------------------");
           CorrelationData correlationData = new CorrelationData("1");
           correlationData.setReturnedMessage(new Message(message.getBytes(),null));
           log.info("发送一条信息给confirm_queue队列: {}", message);
           rabbitTemplate.convertAndSend(ConfirmConfig.CONFIRM_EXCHANGE, ConfirmConfig.CONFIRM_ROUTING_KEY, message,correlationData);
       }
    
    • 1
    • 2
    • 3
    • 4
    • 5
    • 6
    • 7
    • 8
    3.消费者
    @RabbitListener(queues = ConfirmConfig.CONFIRM_QUEUE)
       public void receiveAckMessage(String Object, Message message, Channel channel){
           long deliveryTag = message.getMessageProperties().getDeliveryTag();
           log.info("收到" + ConfirmConfig.CONFIRM_QUEUE + "队列的 {} 消息: {}", deliveryTag, Object);
           try {
               /**
                * 执行业务代码
                */
               channel.basicAck(deliveryTag,false);
               log.info("消费成功: {},消费内容: {}", deliveryTag, Object);
           } catch (Exception e) {
               log.error("手动签收失败",e);
               try {
                   channel.basicNack(deliveryTag,false,true);
               } catch (Exception exception) {
                   log.error("手动拒签失败",exception);
               }
           }
       }
    
    • 1
    • 2
    • 3
    • 4
    • 5
    • 6
    • 7
    • 8
    • 9
    • 10
    • 11
    • 12
    • 13
    • 14
    • 15
    • 16
    • 17
    • 18
    • 19
    4.测试结果

    在这里插入图片描述

    • 1.如果使用上述方案的代码,一旦发生一次消费报错就会死循环
    • 2.因为 basicNack 方法的第三个参数代表是否重回队列
      • 1.如果填 false 消息就直接丢弃,相当于没有保障消息可靠
      • 2.如果填 true ,当发生消费报错之后,该消息会重回消息队列顶端,继续推送到消费端,继续消费该消息
    • 3.通常代码报错并不会因为重试就能解决,所以该消息将会出现死循环情况
      • 1.继续被消费,继续报错,重回队列,继续被消费…

    4.解决方式

    • 1.当消费失败后将此消息存到 Redis,记录消费次数,如果消费了三次还是失败,就丢弃掉消息,记录日志落库保存
    • 2.直接填 false ,不重回队列,记录日志、发送邮件等待开发手动处理
    • 3.不启用手动 ack ,使用 SpringBoot 提供的消息重试
    1.SpringBoot消息重试

    在这里插入图片描述

    2.消费者
    @RabbitListener(queues = ConfirmConfig.CONFIRM_QUEUE)
       public void receiveAckMessage(String Object, Message message, Channel channel){
           long deliveryTag = message.getMessageProperties().getDeliveryTag();
           log.info("收到" + ConfirmConfig.CONFIRM_QUEUE + "队列的 {} 消息: {}", deliveryTag, Object);
           try {
               /**
                * 执行业务代码
                */
           } catch (Exception e) {
               log.error("签收失败", e);
               /**
                * 记录日志、发送邮件、保存消息到数据库,落库之前判断如果消息已经落库就不保存
                */
               throw new RuntimeException("消息消费失败");
           }
           log.info("消费成功: {},消费内容: {}", deliveryTag, Object);
       }
    
    • 1
    • 2
    • 3
    • 4
    • 5
    • 6
    • 7
    • 8
    • 9
    • 10
    • 11
    • 12
    • 13
    • 14
    • 15
    • 16
    • 17
    • 1.注意:一定要手动 throw 一个异常,因为 SpringBoot 触发重试是根据方法中发生未捕捉的异常来决定的
    • 2.该重试是 SpringBoot 提供的,重新执行消费者方法,而不是让 RabbitMQ 重新推送消息

    7.持久化

    • 1.保障RabbitMQ服务停掉后消息生产者发送到队列中的消息不丢失
    • 2.默认RabbitMQ重启或由于某种原因崩溃时,队列和消息将会丢失
    • 3.确保消息不会丢失需要做两件事:队列持久化消息持久化
      在这里插入图片描述
    • 4.实际开发中,队列和消息都做持久化且不自动删除

    1.队列持久化

    • 1.队列持久化是在生产者端,只需声明队列时修改持久化策略为true
      //让消息队列持久化
      channel.queueDeclare(ACK_QUEUE_NAME,true,false,false,null)
      
      • 1
      • 2
    • 2.如果修改的队列已经存在且为非持久化,需先把原先队列删除或重新创建一个持久化的队列,不然就会出现错误
      在这里插入图片描述
    • 3.队列修改为持久化后,即使重启RabbitMQ,队列依然存在
    1.问题
    • 1.队列持久化并不能保证消息不丢失
    • 2.因为队列持久化后虽然存在,但是消息并未持久化,意外情况下依然会丢失
    • 3.因此需要将消息持久化,保存在磁盘

    2.消息持久化

    • 1.消息实现持久化是在生产者端,只需发送消息时修改参数为MessageProperties.PERSISTENT_TEXT_PLAIN
      //未修改前
      channel.basicPublish("",ACK_QUEUE_NAME,null,message.getBytes("UTF-8"));
      //修改后
      channel.basicPublish("",ACK_QUEUE_NAME,MessageProperties.PERSISTEMT_TEXT_PLAIN,message.getBytes("UTF-8"))
      
      • 1
      • 2
      • 3
      • 4
    1.问题
    • 1.将消息标记为持久化并不能完全保证不会丢失消息
    • 2.尽管RabbitMQ将消息保存到磁盘,但存在当消息刚存储到磁盘时,还没存储完,消息还在缓存的一个间隔点
    • 3.此时并没有真正写入磁盘,持久性保证并不强,遇到意外情况消息依然会丢失

    8.消费者分发策略

    • 1.分发策略主要取决于消费者所连接Channel中的预取值
    • 2.默认值为0,即Channel中没有预取值,有消息就分发(轮询),消息过多则会阻塞
    • 3.如果设置了预取值,则该Channel会预取指定值个数的消息,如果接收到超过预取值个数的消息则分发给其他消费者
       /**
        * Request a specific prefetchCount "quality of service" settings for this channel.
        * 为这个Channel设置请求一个预取值
        * Note the prefetch count must be between 0 and 65535 (unsigned short in AMQP 0-9-1).
        * 取值范围为:0~65535
        * @param prefetchCount 服务器将传递的最大消息数,默认值为0
        */
       void basicQos(int prefetchCount) throws IOException;
      
      • 1
      • 2
      • 3
      • 4
      • 5
      • 6
      • 7
      • 8

    1.不公平分发

    • 1.默认RabbitMQ分发消息采用的轮询分发,即预期值默认是0
    • 2.但是该策略并不适合所有场景下(例:一个消费者处理任务的速度非常快,而另一个消费者处理速度却很慢)
    • 3.为了避免这种情况,可设置参数channel.basicQos(1),表示采用不公平分发
      // 设置不公平分发,默认值为0
      channel.basicQos(1);
      
      • 1
      • 2
    • 4.设置为1是不公平分发,设置为0为轮询分发
    • 5.设置不公平分发策略在消费者端,只需在所有的消费者端的消费信息之前加上该设置

    2.指定值分发

    • 1.消息的发送是异步送的,Channel中不止一个消息
    • 2.消费者的手动应答本质上是异步的,因此存在一个未确认消息缓冲区
    • 3.使用basicQos方法设置预取值来限制此缓冲区的大小,避免缓冲区里无限制的未确认消息
    • 4.预取值定义通道上允许的未确认消息最大数量,一旦数量达到配置的数量,RabbitMQ将停止在该通道上传递更多消息,除非至少有一个未处理的消息被确认
    • 5.增加预取值将提高向消费者传递消息的速度,可指定每一个信道里面有多少个消息可消费,轮询方式分发,如果满了就不能再插入信息

    3.对比自动应答

    • 1.虽然自动应答传输消息速率是最佳的,但在该情况下已传递但尚未处理的消息数量也会增加,从而增加了消费者RAM消耗(随即存取存储器)
    • 2.避免使用具有无限预处理自动确认模式或手动确认模式,因为消费者消费了大量的信息如果没有确认的话,会导致消费者连接节点的内存消耗变大
    • 3.找到合适的预取值是一个反复实验的过程,不同的负载该值取值也不同
      • 1.100300范围内的值通常可提供最佳的吞吐量,并且不会给消费者带来太大的风险
      • 2.1是最保守的,但这将使吞吐量变得很低,特别不利于消费者延迟很严重或消费者连接等待时间较长的情况

    4.实际项目配置

    在这里插入图片描述

    9.发布确认机制

    1.消息持久化问题

    • 1.消息持久化到磁盘时需要一定时间,如果这时遇到意外终止,则消息会丢失
    • 2.因此生产者需要开启发布确认机制确保消息完全不丢失

    2.发布确认原理

    • 1.生产者将信道(Channel)设置成 Confirm 模式,一旦信道进入该模式,所有在该信道上面发布的消息都将会被指派一个唯一 ID(从 1 开始)
    • 2.一旦消息被投递到所有匹配的队列之后,Broker就会发送一个确认给生产者(包含消息的唯一 ID),使得生产者知道消息已经正确到达目的队列
    • 3.如果消息和队列是可持久化的,那么确认消息会在将消息写入磁盘后发出,保证了消息存储的可靠性
    • 4.Broker 回传给生产者的确认消息中 delivery-tag 包含了确认消息的序列号
    • 5.发布确认机制有三种模式
      • 1.单个确认模式:生产者每发一个信息都需要等待Broker发送确认后才可以发送下一条消息
      • 2.批量确认模式:生产者指定发一定数量的信息才需要等待Broker发送确认后才可以发送下一批消息
      • 3.异步确认模式:生产者可以发布任意条消息,然后通过异步回调等待Broker发送的确认和错误消息,并可以通过回调对该信息进行处理

    3.确保消息完全不丢失步骤

    • 1.生产者
      • 1.开启发布确认模式(消息保存到磁盘后,Broker需要向生产者发送确认信息)
    • 2.Broker
      • 1.开启队列持久化
      • 2.开启消息持久化
    • 3.消费者
      • 1.开启手动应答模式

    4.开启发布确认模式

    • 1.发布确认模式默认未开启,如果开启需在生产者发送消息前调用方法confirmSelect
    • 2.每当生产者需要开启发布确认模式,都需要在Channel上调用该方法
      Channel channel = connection.createChannel();
      channel.confirmSelect();
      
      • 1
      • 2

    5.单个确认发布

    • 1.一种同步确认发布方式,发布一个消息之后只有其被确认发布,后续的消息才能继续发布
    • 2.缺点:发布速度慢,因为没有确认发布的消息就会阻塞所有后续消息的发布
    • 3.waitForConfirms这个方法只有在消息被确认
      的时候才返回,如果在指定时间范围内这个消息没有被确认那么它将抛出异常
      /**
        * 等待直到发布的所有消息都已被Broker确认或未确认
        * 注意当在non-Confirm channel上调用时,waitForConfirms会抛出一个IllegalStateException异常
        */
       boolean waitForConfirms() throws InterruptedException;
      
      • 1
      • 2
      • 3
      • 4
      • 5
      package com.rabbit.work.confirm;
      
      import com.rabbit.util.RabbitMQUtils;
      import com.rabbitmq.client.Channel;
      import com.rabbitmq.client.MessageProperties;
      
      import java.util.UUID;
      
      /**
       * 单个消息等待发布状态 消息持久化共耗时3099毫秒 消息不持久化共耗时393毫秒
       */
      public class WorkProducer {
          // 发送消息的总个数
          public static final int MESSAGE_COUNT = 1000;
      
          public static void main(String[] args) throws Exception {
              Channel channel = RabbitMQUtils.getChannel();
              String QUEUE_NAME = UUID.randomUUID().toString();
              channel.queueDeclare(QUEUE_NAME,true,false,false,null);
              // 开启发布确认模式
              channel.confirmSelect();
              // 开始时间
              long begin = System.currentTimeMillis();
              // 批量发送消息
              for (int i=0; i<MESSAGE_COUNT; i++) {
                  String message = i + "";
                  //消息不持久化
                  channel.basicPublish("",QUEUE_NAME,null,message.getBytes());
                  //消息持久化
                  //channel.basicPublish("",QUEUE_NAME,MessageProperties.PERSISTENT_TEXT_PLAIN,message.getBytes());
                  // 单个消息等待发布状态
                  boolean flag = channel.waitForConfirms();
                  if(flag){
                      System.out.println("信息发送成功");
                  }else {
                      System.out.println("信息发送失败");
                  }
              }
              // 结束时间
              long end = System.currentTimeMillis();
              System.out.println("共耗时"+ (end-begin) +"毫秒");
          }
      }
      
      • 1
      • 2
      • 3
      • 4
      • 5
      • 6
      • 7
      • 8
      • 9
      • 10
      • 11
      • 12
      • 13
      • 14
      • 15
      • 16
      • 17
      • 18
      • 19
      • 20
      • 21
      • 22
      • 23
      • 24
      • 25
      • 26
      • 27
      • 28
      • 29
      • 30
      • 31
      • 32
      • 33
      • 34
      • 35
      • 36
      • 37
      • 38
      • 39
      • 40
      • 41
      • 42
      • 43

    6.批量确认发布

    • 1.一批消息一批确认可极大提高吞吐量
    • 2.缺点:当发生故障导致发布(发送)出现问题时,不知道是哪个消息出现问题
      package com.rabbit.work.confirm;
      
      import com.rabbit.util.RabbitMQUtils;
      import com.rabbitmq.client.Channel;
      import com.rabbitmq.client.MessageProperties;
      
      import java.util.UUID;
      
      /**
       * 批量消息等待发布状态 消息持久化共耗时660毫秒 消息不持久化共耗时121毫秒
       */
      public class BatchWorkProducer {
          // 发送消息的总个数
          public static final int MESSAGE_COUNT = 1000;
      
          public static void main(String[] args) throws Exception {
              Channel channel = RabbitMQUtils.getChannel();
              String QUEUE_NAME = UUID.randomUUID().toString();
              channel.queueDeclare(QUEUE_NAME,true,false,false,null);
              // 开启发布确认模式
              channel.confirmSelect();
              // 开始时间
              long begin = System.currentTimeMillis();
              // 批量步长
              int batch = 100;
              // 批量发送消息
              for (int i=0; i<MESSAGE_COUNT; i++) {
                  String message = i + "";
                  //消息不持久化
                  channel.basicPublish("",QUEUE_NAME,null,message.getBytes());
                  //消息持久化
      			//channel.basicPublish("",QUEUE_NAME, MessageProperties.PERSISTENT_TEXT_PLAIN,message.getBytes());
                  // 单个消息等待发布状态
                  if(i%batch==0){
                      //如果该判断放在循环后面表示全部发送完再确认,现在是每满100条确认一次
                      boolean flag = channel.waitForConfirms();
                      if(flag){
                          System.out.println("信息发送成功");
                      }else {
                          System.out.println("信息发送失败");
                      }
                  }
              }
              // 结束时间
              long end = System.currentTimeMillis();
              System.out.println("共耗时"+ (end-begin) +"毫秒");
              //关闭,否则会内存泄露造成程序崩溃
              channel.close();
          }
      }
      
      • 1
      • 2
      • 3
      • 4
      • 5
      • 6
      • 7
      • 8
      • 9
      • 10
      • 11
      • 12
      • 13
      • 14
      • 15
      • 16
      • 17
      • 18
      • 19
      • 20
      • 21
      • 22
      • 23
      • 24
      • 25
      • 26
      • 27
      • 28
      • 29
      • 30
      • 31
      • 32
      • 33
      • 34
      • 35
      • 36
      • 37
      • 38
      • 39
      • 40
      • 41
      • 42
      • 43
      • 44
      • 45
      • 46
      • 47
      • 48
      • 49
      • 50

    7.异步确认发布

    在这里插入图片描述

    • 1.异步确认编程逻辑复杂,但可靠性和效率高
    • 2.生产者利用监听器监听Broker回调函数达到消息可靠性传递,Broker也是通过监听函数回调来保证是否向消费者投递成功
       /**
        * Add a lambda-based {@link ConfirmListener}.
        * @see ConfirmListener
        * @see ConfirmCallback
        * @param ackCallback callback on ack
        * @param nackCallback call on nack (negative ack)
        * @return the listener that wraps the callbacks
        */
       ConfirmListener addConfirmListener(ConfirmCallback ackCallback, ConfirmCallback nackCallback);
      
      • 1
      • 2
      • 3
      • 4
      • 5
      • 6
      • 7
      • 8
      • 9
      package com.rabbit.work.confirm;
      
      import com.rabbit.util.RabbitMQUtils;
      import com.rabbitmq.client.Channel;
      import com.rabbitmq.client.ConfirmCallback;
      import com.rabbitmq.client.MessageProperties;
      
      import java.util.UUID;
      
      /**
       * 异步发布等待发布状态 消息持久化共耗时74毫秒 消息不持久化共耗时34毫秒
       */
      public class AsynWorkProducer {
          // 发送消息的总个数
          public static final int MESSAGE_COUNT = 1000;
      
          public static void main(String[] args) throws Exception {
              Channel channel = RabbitMQUtils.getChannel();
              String QUEUE_NAME = UUID.randomUUID().toString();
              channel.queueDeclare(QUEUE_NAME,true,false,false,null);
              // 开启发布确认模式
              channel.confirmSelect();
              // 开始时间
              long begin = System.currentTimeMillis();
              // 消息确认成功 回调函数
              ConfirmCallback ackCallback = (deliveryTag, multiple) -> {
                  System.out.println("确认消息的标识: " + deliveryTag + " :: " + (multiple?"批量处理":"不批量处理"));
              };
              //  消息未确认成功 回调函数
              ConfirmCallback nackCallback = (deliveryTag, multiple) -> {
                  System.out.println("未确认消息的标识: " + deliveryTag);
              };
              //发送消息前准备监听器,监听Broker的回调函数,异步通知
              /**
               * 1.监听确认成功消息
               * 2.监听确认失败消息
               */
              channel.addConfirmListener(ackCallback,nackCallback);
              // 批量发送消息
              for (int i=0; i<MESSAGE_COUNT; i++) {
                  String message = i + "";
                  //消息不持久化
                  channel.basicPublish("",QUEUE_NAME,null,message.getBytes());
                  // 消息持久化
                  // channel.basicPublish("",QUEUE_NAME, MessageProperties.PERSISTENT_TEXT_PLAIN,message.getBytes());
              }
              // 结束时间
              long end = System.currentTimeMillis();
              System.out.println("共耗时"+ (end-begin) +"毫秒");
              //异步不能关闭通道,否则无法监听回调
      		//channel.close();
          }
      }
      
      • 1
      • 2
      • 3
      • 4
      • 5
      • 6
      • 7
      • 8
      • 9
      • 10
      • 11
      • 12
      • 13
      • 14
      • 15
      • 16
      • 17
      • 18
      • 19
      • 20
      • 21
      • 22
      • 23
      • 24
      • 25
      • 26
      • 27
      • 28
      • 29
      • 30
      • 31
      • 32
      • 33
      • 34
      • 35
      • 36
      • 37
      • 38
      • 39
      • 40
      • 41
      • 42
      • 43
      • 44
      • 45
      • 46
      • 47
      • 48
      • 49
      • 50
      • 51
      • 52
      • 53

    8.处理异步未确认消息

    • 1.首先把所有消息都保存起来,然后过滤出未确认的消息
    • 2.将未确认的消息放到一个基于内存的能被发布线程访问的集合
    • 3.ConcurrentSkipListMap集合可以在ConfirmCallback发布线程之间进行消息的传递
      package com.rabbit.work.confirm;
      
      import com.rabbit.util.RabbitMQUtils;
      import com.rabbitmq.client.Channel;
      import com.rabbitmq.client.ConfirmCallback;
      
      import java.util.UUID;
      import java.util.concurrent.ConcurrentNavigableMap;
      import java.util.concurrent.ConcurrentSkipListMap;
      
      /**
       * 异步发布未确认消息处理
       */
      public class AsyWorkNoDailProducer {
          // 发送消息的总个数
          public static final int MESSAGE_COUNT = 1000;
      
          public static void main(String[] args) throws Exception {
              Channel channel = RabbitMQUtils.getChannel();
              String QUEUE_NAME = UUID.randomUUID().toString();
              channel.queueDeclare(QUEUE_NAME,true,false,false,null);
              // 开启发布确认模式
              channel.confirmSelect();
              /**
               * 线程安全有序的哈希表,适用于高并发情况,该程序存在监听和发送两个线程
               * 1.序号
               * 2.内容
              */
              ConcurrentSkipListMap<Long, String> messageContainer = new ConcurrentSkipListMap<>();
      
              // 消息确认成功 回调函数
              ConfirmCallback ackCallback = (deliveryTag, multiple) -> {
                  System.out.println("确认消息的标识: " + deliveryTag + " :: " + (multiple?"批量处理":"不批量处理"));
                  // 如果批量则批量清除集合中确认的信息,其余则为未确认信息
                  if(multiple){
                      // 确认消息的标识
                      ConcurrentNavigableMap<Long, String> ackMessages = messageContainer.headMap(deliveryTag);
                      ackMessages.clear();
                  }else {
                      messageContainer.remove(deliveryTag);
                  }
      
              };
              //  消息未确认成功 回调函数
              ConfirmCallback nackCallback = (deliveryTag, multiple) -> {
                  String message = messageContainer.get(deliveryTag);
                  System.out.println("未确认消息的标识: " + deliveryTag + " :: 内容: " + message);
              };
              //发送消息前准备监听器,监听Broker的回调函数,异步通知
              /**
               * 1.监听确认成功消息
               * 2.监听确认失败消息
               */
              channel.addConfirmListener(ackCallback,nackCallback);
              // 开始时间
              long begin = System.currentTimeMillis();
              // 批量发送消息
              for (int i=0; i<MESSAGE_COUNT; i++) {
                  String message = i + "";
                  //消息不持久化
                  channel.basicPublish("",QUEUE_NAME,null,message.getBytes());
                  // 消息持久化
                  // channel.basicPublish("",QUEUE_NAME, MessageProperties.PERSISTENT_TEXT_PLAIN,message.getBytes());
                  // 保存发送的信息
                  messageContainer.put(channel.getNextPublishSeqNo(),message);
              }
              // 结束时间
              long end = System.currentTimeMillis();
              System.out.println("共耗时"+ (end-begin) +"毫秒");
              //异步不能关闭通道,否则无法监听回调
      //        channel.close();
          }
      }
      
      • 1
      • 2
      • 3
      • 4
      • 5
      • 6
      • 7
      • 8
      • 9
      • 10
      • 11
      • 12
      • 13
      • 14
      • 15
      • 16
      • 17
      • 18
      • 19
      • 20
      • 21
      • 22
      • 23
      • 24
      • 25
      • 26
      • 27
      • 28
      • 29
      • 30
      • 31
      • 32
      • 33
      • 34
      • 35
      • 36
      • 37
      • 38
      • 39
      • 40
      • 41
      • 42
      • 43
      • 44
      • 45
      • 46
      • 47
      • 48
      • 49
      • 50
      • 51
      • 52
      • 53
      • 54
      • 55
      • 56
      • 57
      • 58
      • 59
      • 60
      • 61
      • 62
      • 63
      • 64
      • 65
      • 66
      • 67
      • 68
      • 69
      • 70
      • 71
      • 72
      • 73

    9.实际开发处理

    • 1.实际开发处理参考发布确认模式

    10. 消息重复消费(幂等性)

    • 1.同一操作发起的一次请求或者多次请求的结果是一致的,不会因为多次请求而产生了副作用(例:保证支付消息只被消费一次)

    1.消息重复消费

    • 1.MQ把消息发送给消费者,消费者消费MQ中的消息
    • 2.消费完成后,消费者在给MQ返回ack时网络中断
    • 3.MQ未收到确认消息,该条消息会重新发送给其他消费者或在网络重连后再次发送给该消费者
    • 4.但实际上该消费者已经成功消费了该信息,造成消费者消费了重复的消息
    • 5.一般海量订单生成的业务高峰期,消费者有可能会重复消费消息
    • 6.解决方式一般有两种
      • 1.第一种:确保消费端只执行一次
      • 2.第二种:允许消费端执行多次,保证数据不受影响

    2.解决方式

    • 1.确保消费端只执行一次
      • 1.全局唯一标识(例:UUID)
      • 2.利用redis的原子性
    • 2.允许消费端执行多次,保证数据不受影响
      • 1.数据库唯一键约束
        • 1.如果消费端业务是新增操作,可以利用数据库的唯一键约束,如果重复消费将会插入两条相同的记录,数据库会报错
      • 2.数据库乐观锁思想
        • 1.如果消费端业务是更新操作,可以给业务表加一个 version 字段,每次更新把 version 作为条件,更新之后 version + 1
        • 2.由于 MySQLinnoDB 是行锁,当其中一个请求成功更新之后,另一个请求才能进来,由于版本号 version 已经变成 2,必定更新的 SQL 语句影响行数为 0,不会影响数据库数据
    1.全局唯一标识
    • 1.全局唯一标识:自定义规则或时间戳生成的全局唯一标识,基本都是由业务规则拼接
    • 2.需要保证唯一性,利用查询语句进行判断这个id是否存在数据库中
    • 3.优点:实现简单,拼接id然后查询判断该id是否重复
    • 4.缺点:高并发时,单个数据库会有写入性能瓶颈
    • 5.不推荐使用该方式
    2.Redis原子性
    • 1.同上述一样生成全局唯一标识
    • 2.利用redis执行setnx命令,天然具有幂等性,从而实现不重复消费
    3.代码实现
    @RabbitListener(queues = ConfirmConfig.CONFIRM_QUEUE)
       public void receiveRepeatMessage(String Object, Message message, Channel channel){
           long deliveryTag = message.getMessageProperties().getDeliveryTag();
           log.info("收到" + ConfirmConfig.CONFIRM_QUEUE + "队列的 {} 消息: {}", deliveryTag, Object);
           try {
               String uuid = UUID.randomUUID().toString().replace("", "\\-");
               Boolean flag = redisTemplate.hasKey(uuid);
               if(Boolean.TRUE.equals(flag)){
                   return;
               }
               /**
                * 执行业务代码
                */
               redisTemplate.opsForValue().set(uuid,"1",10,TimeUnit.SECONDS);
           } catch (Exception e) {
               log.error("签收失败", e);
               /**
                * 记录日志、发送邮件、保存消息到数据库,落库之前判断如果消息已经落库就不保存
                */
               throw new RuntimeException("消息消费失败");
           }
           log.info("消费成功: {},消费内容: {}", deliveryTag, Object);
       }
    
    • 1
    • 2
    • 3
    • 4
    • 5
    • 6
    • 7
    • 8
    • 9
    • 10
    • 11
    • 12
    • 13
    • 14
    • 15
    • 16
    • 17
    • 18
    • 19
    • 20
    • 21
    • 22
    • 23

    11.消息积压

    1.产生原因

    • 1.一般是由于消费端消费速度远小于生产者发消息速度,导致大量消息在 RabbitMQ 的队列中无法消费

    2.解决方式

    • 1.生产者发消息适当限流(不推荐,影响用户体验)
    • 2.多部署几台消费者实例(推荐)
    • 3.适当增加 prefetch 的数量,使消费端一次多接受一些消息(推荐,可以和第二种方案一起用)

    12.消息顺序性

    • 1.某些业务场景需要让消息顺序消费
    • 2.例:使用 canal 订阅 MySQLbinary 日志来更新 Redis,通常会将 canal 订阅到的数据变化发送到消息队列
    • 3.如果不保证 RabbitMQ 的顺序消费, Redis 中就有可能会出现脏数据
      在这里插入图片描述

    1.单个消费者实例

    • 1.一般情况下单个消费者本身是有顺序的,不过也可能因为网络问题RabbitMQ集群导致无序
    • 2.可以在消费者方使用redis中的zset对消息进行排序后再进行消费
    • 3.也可以设置 RabbitMQ 每次只推送一个消息,再开启手动 ack 即可
      在这里插入图片描述

    2.多个消费者实例

    • 1.当消费者是多个实例时,队列中的消息会分发到所有实例进行消费(同一个消息只能发给一个消费者实例)
    • 2.这样不能保证消息顺序的消费,因为不能确保哪台机器执行消费端业务代码的速度快
    • 3.可以在多个消费者前面设置一个统一的redis,消息先进入redis中通过zset排序,然后再分发给多个消费者处理

    12.RabiitMQ集群

    • 1.一般单机版RabbitMQ无法满足目实际项目要求
      • 1.例:RabbitMQ服务器遇到了内存崩溃机器宕机或者重启等情况
      • 2.例:应用需RabbitMQ服务满足每秒10万条消息的吞吐量,但单台RabbitMQ服务器只能满足每秒1000条消息的吞吐量
    • 2.此时为了满足业务需求,需要搭建RabbitMQ集群
    • 3.RabbitMQ集群不是主从模式的集群,而是对等集群,即节点之间是对等,不存在主从关系(镜像除外)
    • 4.不同节点存放的是不同的队列,不同的消息,集群内所有节点的所有队列的消息共同组成该集群的所有消息
    1.高可用
    • 1.普通RabbitMQ集群无法实现高可用性,可用使用镜像队列实现高可用性
    • 1.镜像队列中的某些节点作为该队列的主master节点的mirror节点,存在主从关系,该主从主要用于实现队列消息的冗余存储实现队列和队列内部消息的高可用,而不是用于读写分离实现高并发
    2.高吞吐
    • 1.RabbitMQ通过集群可以解决单节点处理海量消息时的性能瓶颈,从而实现高吞吐量
    • 2.RabbitMQ集群处理海量消息是通过在集群的多个节点建立多个不同的队列来分散消息到多个不同节点的,所以在业务层需要对消息进行细分类别来放到不同队列中
    3.高并发
    • 1.大部分中间件的集群(RedisZookeeperMySQL等),采用读写分离的方式来应对高并发数据请求
    • 2.RabbitMQ集群是一个对等集群,集群内部不同的节点存放不同队列(所以也就存放着不同的数据)
    • 3.从客户端的角度看,RabbitMQ集群其实就是一个逻辑上的单节点,即客户端可以连接RabbitMQ集群的任何一个节点,而不需要约束为只能连接该客户端所要请求的消息对应的队列所在的节点
    • 4.任何一个节点接收到客户端的请求后,如果该请求对应的消息属于自身节点的队列,则由该节点处理,否则将该请求转发给该消息所属队列所在的节点处理,这个过程对客户端来说是透明的
    • 5.所以通过这种方式RabbitMQ集群可以通过增加节点个数来应对高并发请求,因为客户端可以连接任意一个节点
    4.搭建步骤
    • 1.参考Linux软件安装文章中的RabbitMQ集群搭建
    5.优缺点

    在这里插入图片描述

    • 1.优点
      • 1.搭建普通集群模式,可以从多个节点消费队列数据
    • 2.缺点
      • 1.该模式虽然可以从多个节点消费队列数据,但本质上只是从原队列节点上拉取数据
      • 2.如果真正存放数据的节点宕机,则其他节点不可复用,无法实现高可用性,且该模式无法解决单节点性能瓶颈
      • 3.因此需要使用镜像队列解决该问题

    1.镜像队列

    在这里插入图片描述

    • 1.上述搭建的普通集群中的队列不可复用,即node1上创建的队列node2node3上并不存在,一旦node1宕机,其上面的队列就会消失
    • 2.引入镜像队列(Mirror Queue)机制,将队列镜像到集群中的其他Broker节点之上
    • 3.如果集群中的一个节点失效,队列能自动地切换/手动指定到镜像中的另一个节点上以保证服务的可用性
    • 4.镜像队列的实质:将指定队列备份到节点上
    • 5.该模式可以解决队列不可复用问题,实现高可用性
    • 6.一般备份一个节点即可,即一个主机一个备机,让服务器自动选择备份,这种方式只要集群中还有一台节点可用,都可以保证信息不会丢失
    1普通队列问题
    • 1.如图所示,普通集群的队列,只存在当前节点上
      在这里插入图片描述
      在这里插入图片描述
    • 2.如果当前连接节点宕机则会报错,且队列消失,所以需要负载均衡器镜像队列实现高可用
      在这里插入图片描述
      在这里插入图片描述
      在这里插入图片描述
    2.镜像队列搭建
    • 1.确保节点都是启动状态
      在这里插入图片描述
    • 2.任意节点上设置镜像策略
      在这里插入图片描述
      在这里插入图片描述
    • 3.设置成功
      在这里插入图片描述
      在这里插入图片描述
    • 4.测试成功,即使整个集群只剩下一个节点,依然能消费队列里面的消息
      在这里插入图片描述

    2.Haproxy+Keepalive实现高可用负载均衡

    1.整体架构图

    在这里插入图片描述

    2.Haproxy实现负载均衡
    • 1.HAProxy 提供高可用性、负载均衡及基于 TCPHTTP 应用的代理
    • 2.类似Nginx,这里作为负载均衡器
    3.搭建步骤(存在问题,待完善)
    • 1.下载 haproxy(选择两个节点分别下载,一个主一个备)
      # node1和node2节点分别下载
      yum -y install haproxy
      
      • 1
      • 2
      在这里插入图片描述
    • 2.修改 node1node2haproxy.cfg
      vim /etc/haproxy/haproxy.cfg
      
      • 1
      在这里插入图片描述
      在这里插入图片描述
    • 注意:设置如果报错按这个配置,代理端口修改为5672
      在这里插入图片描述
    • 3.两个节点分别启动 haproxy(注意:启动前需要先安装httpd
      yum -y install httpd
      haproxy -f /etc/haproxy/haproxy.cfg
      systemctl status haproxy
      ps -ef | grep haproxy
      
      • 1
      • 2
      • 3
      • 4
      在这里插入图片描述
      在这里插入图片描述
      在这里插入图片描述
    • 4.访问地址(尽管haproxy状态显示failed,任然可以转发)
      http://192.168.73.130:8000
      
      • 1
      在这里插入图片描述
    4.Keepalived实现双机(主备)热备
    • 1.如果配置的 HAProxy 主机突然宕机或者网卡失效,那么虽然 RbbitMQ 集群没有任何故障,但是对于外界的客户端来说所有的连接都会被断开,结果将是灾难性
    • 2.为了确保负载均衡服务的可靠性同样十分重要,引入 Keepalived 通过自身健康检查、资源接管功能做高可用(双机热备),实现故障转移
    5.搭建步骤(存在问题,待完善)
    • 1.下载 keepalived

       yum -y install keepalived
      
      • 1
    • 2.修改节点 node1 配置文件

      # 把资料里面的 keepalived.conf 修改之后替换
      vim /etc/keepalived/keepalived.conf
      
      • 1
      • 2
    • 3.修改节点 node2 配置文件

      需要修改 global_defs 的 router_id,如:nodeB
      其次要修改 vrrp_instance_VI 中 state 为"BACKUP";
      最后要将 priority 设置为小于 100 的值
      
      • 1
      • 2
      • 3
    • 4.添加 haproxy_chk.sh

      为了防止 HAProxy 服务挂掉之后 Keepalived 还在正常工作而没有切换到Backup 上
      所以这里需要编写一个脚本来检测 HAProxy 务的状态
      当 HAProxy 服务挂掉之后该脚本会自动重启HAProxy 的服务
      如果不成功则关闭 Keepalived 服务,这样便可以切换到Backup 继续工作
      vim /etc/keepalived/haproxy_chk.sh(可以直接上传文件)
      修改权限 chmod 777 /etc/keepalived/haproxy_chk.sh
      
      • 1
      • 2
      • 3
      • 4
      • 5
      • 6
    • 5.启动 keepalive 命令(node1node2 启动)

      systemctl start keepalived
      
      • 1
    • 6.观察 Keepalived 的日志

      tail -f /var/log/messages -n 200
      
      • 1
    • 7.观察最新添加的 vip

      ip add show
      
      • 1
    • 8.node1 模拟 keepalived 关闭状态

      systemctl stop keepalived 
      
      • 1
    • 9.使用 vip 地址来访问 rabbitmq 集群

    3.Federation Exchange

    4.Federation Queue

    5.Shovel

  • 相关阅读:
    springbean的生命周期
    Android 引入库报错 Null extracted folder for artifact 解决方案
    软考140-上午题-【软件工程】-软件工具
    【深入理解Kotlin协程】Google的工程师们是这样理解Flow的?
    【初学者入门C语言】之运算符及表达式(二)
    【思科】MPLS VPN 实验配置
    Ubuntu20.04下载opencv3.4--未完善
    按照规则来,为什么还是提图片必须是一条字符串啊!
    集成学习思想
    ubuntu 18.04 双系统安装
  • 原文地址:https://blog.csdn.net/wu246051/article/details/126178462