往期开发(已经完成)

今日后端开发

  • 分析目前系统现状

    现状: 目前实现异步的方式是通过本地的从线程池实现的

    1. 无法集中的去限制,只能够单机限制
      假设 AI 服务限制只能有 2-3 个用户同时去调用,目前可以通过设置线程池最大核心线程数为 2-3 来实现。
      假设系统用量极具增大,要求将系统进行分布式改造,多台服务器,每个服务器都要有两个最大核心线程数 2-3
      就会有 2N-3N 个线程数这样总数加起来就超过 AI 服务最大的限制
      解决方案:将所有的服务统一的发送到一个消息管理地方,然后统一的把任务分发给各个分布式系统(集中地去存储并分发要执行的项目)

    2. 通过本地线城池实现的异步化,由于任务是放在内存中执行的,丢失的可能性非常大,假设系统宕机就会导致系统丢失本次任务
      虽然任务会丢失,但是可以通过人工手段,从数据库中取出,通过某些手段,重新执行此任务,但是开发维护成本非常高
      以上就是很经典的重试场景,可以通过定时任务去解决,是不需要我们开发者过于关心任务丢失重试的场景
      解决方案:把任务存储到一个持久化的存储硬盘来进行存储。

    3. 优化
      如果系统功能的越来越多,长耗时任务越来越长,系统越来越复杂,例如要开多个线程池,资源可能会出现资源竞争的情况
      此时我们可以把系统进行拆分,也就是服务拆分,把长耗时任务,消耗很多资源的任务,单独的抽成一个程序,不影响主业务的执行
      解决方案:可以使用一个中间件(中间人),让中间件去帮我们去连接俩个系统(程序),比如核心系统和智能生成业务

  • 引入中间件(RabbitMQ)

    连接多个系统的桥梁,或者帮助多个系统进行通信、协作
    例如 Redis、消息队列、分布式存储 Etcd

    样例图

    • 消息队列

    1. 概念:存储消息的队列
    2. 关键字:存储、消息、队列
    3. 存储:存储数据
    4. 消息:某种数据结构,比如字符串、对象、JSON 字符串、二进制数据等等
    5. 队列:先进先出的数据结构
    6. 见解:消息队列是一种特殊的数据库吗?应该是可以这样理解吧
    7. 应用场景
      在多个不同的系统,或者是应用之间实现消息的传输,也可以进行存储
      不需要考虑传输应用的编程语言、系统、框架等。可以对应用进行解耦
    • 消息队列的模型

    1. 生产者:Producer,可以比作是快递员,发送消息的人(客户端)
    2. 消费者:Consumer,可以比作是取快递的人,接受读取消息的人(客户端)
    3. 消息:Message,可以比作是快递,就是快递员派送的快递也是给消费者传递的消息
    4. 消息队列:Queue,可以比作是快递用快递车来进行派送,存储到车里面进行派送,消息派送的队列

    为什么不直接进行传输,而要用消息队列?
    这是因为生产者不需要消费者什么时候去消息,要不要消费,生产者只需要把东西给消息队列,业务就完成了
    这就将生产者和消费者进行了解耦,俩者互相不影响,自己完成任务即可

    过程图

    • 为什么使用消息队列

      • 异步处理

        1. 生产者发送完消息之后,就可以继续去接受其他任务,不需要等待,
        2. 消费者想什么时候去取出这个消息去消费都可以,不会产生堵塞
      • 削峰填谷

        1. 先把用户请求,让生产者发送到消息队列中,消费者可以按照自己的需求,慢慢的去取消息
        2. 例如:某一时刻来了十万个请求,原本情况下十万个请求要直接到内部进行处理,很快系统就扛不住压力就会宕机了
        3. 解决:请求进来时候,先把消息放到消息队列中,系统按照自己最大的处理能力,去定量的取出消息恒定速率的去处理消息,从而降低了系统的压力,稳定的去处理
      • 消息队列的优势

        1. 数据持久化:可以把消息几种的存储到硬盘里,服务器重启也不会丢失
        2. 可扩展性:可以根据需求,随时的增加或者减少节点,继续保持稳定的服务
        3. 应用解耦:可以连接各个不同语言、框架开发系统,让这些系统能够灵活的传输读取数据等等。
      • 应用解耦的优点

        1. 最早的开发是把所有的功能都放到同一个项目中,调用多个子系统是,一个环节出错或者导致系统挂掉,系统整体就要出错
        2. 使用消息队列进行解耦。一个系统挂了就不会影响另一个系统的运行
        3. 系统挂了在后续恢复之后,仍然可以继续取出消息,重试之前的业务逻辑
        4. 生产者把消息发送到消息队列,就可以立即返回,不用同步调用所有系统,性能也会提高
      • 订阅模式

        1. 如果一个非常大的系统要给其他子系统发送通知,最简单的方式就是一个个的去通知子系统,调用相关的系统去通知。
        2. 这样会产生很大的问题:每次更新发通知就要一个个调用,效率很低,速度还慢,还一直在占用资源。还有可能通知过程中有的调用失败没有通知到
        3. 解决方案:大的核心系统每次通知把消息只发送到一个地方,其他的系统都去订阅这个地方也就是消息队列,去读取这个消息队列的通知信息
      • 消息队列的缺点

        俗话说没有绝对完美的设计,总归这么多好处的消息队列,也有它的缺点
        如果要给系统引入额外的消息中间件,系统就会变得更加复杂,庞大,而且还要花费成本去维护这个中间件,额外的费用去部署
        还要面临,消息丢失,处理的顺序,重复消费,数据的一致性,也就是分布式系统要考虑的问题

    • 主流的消息队列选型

      • 主流技术

        1. activemq
        2. rabbitmq
        3. kafka
        4. rocketmq
        5. zeromq
        6. pulsar
        7. Apache Inlong (Tube)
      • 技术的选型指标

        1. 吞吐量:IO、并发
        2. 时效性:类似延迟、消息的发送、到达时间
        3. 可用性:系统可用的比率 宕机的可能性
        4. 可靠性:消息不丢失,功能正常完成

        选型图

    • RabbitMQ 入门实战

      • 特点

        生态好、好学习、易于理解、时效性强、支持很多不同语言的客户端、可扩展性、可用性都不错。
        学习性价比高的消息队列,适用于绝大多数中小规模的分布式系统
        官网:https://www.rabbitmq.com

      • 基本概念

        AMQP 协议: https://www.rabbitmq.com/tutorials/amqp-concepts.html
        高级消息队列协议 (Advanced Message Queue Protocol)

        1. 生产者:发消息到某个交换机
        2. 消费者: 从某个队列中取消息
        3. 交换机 (Exchange): 负责把消息 转发到对应的队列
        4. 队列 (Queue) : 存储消息的
        5. 路由(Routes): 转发,就是怎么把消息从一个地方转到另一个地方 (比如从生产者转发到某个队列)

        mq

        引入依赖

        1
        2
        3
        4
        5
        6
        <!-- https://mvnrepository.com/artifact/com.rabbitmq/amqp-client -->
        <dependency>
        <groupId>com.rabbitmq</groupId>
        <artifactId>amqp-client</artifactId>
        <version>5.17.0</version>
        </dependency>
      • 单向发送

        一个生产者给一个队列发送消息,一个消费者从这个队列中取消息 1<->1

        1-1

        • 生产者代码
        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.bi.spring.demomq;

        import com.rabbitmq.client.Channel;
        import com.rabbitmq.client.Connection;
        import com.rabbitmq.client.ConnectionFactory;
        import lombok.extern.slf4j.Slf4j;

        import java.nio.charset.StandardCharsets;

        @Slf4j
        public class Producer {

        private static final String QUEUE_NAME = "hello_mq";

        public static void main(String[] args) {
        ConnectionFactory factory = new ConnectionFactory();
        factory.setHost("localhost");
        try {
        // 从工厂中创建连接
        Connection connection = factory.newConnection();
        // 获取通信通道
        Channel channel = connection.createChannel();
        // 声明一个队列
        channel.queueDeclare(QUEUE_NAME, false, false, false, null);
        String message = "Hello,World";
        // 发送消息
        channel.basicPublish("", QUEUE_NAME, null, message.getBytes(StandardCharsets.UTF_8));
        // 打印信息
        System.out.println("Producer send message: " + message);
        } catch (Exception e) {
        log.info("创建链接失败!E:" + e);
        }
        }
        }

        Channel 频道:可以理解为消息队列的 Client(比如 jdbcClient,redisClient),提供了和消息队列简历 server 建立的通信的传输方法为了复用链接,提高传输效率。程序通过操作 Channel 来操作 rabbitmq(收发消息操作)

        创建消息队列的参数

        1. queueName:消息队列名称 同名称和参数的消息队列只能创建一次
        2. durabale:消息队列重启后,消息是否丢失
        3. exclusive:是否只允许当前创建的消息队列的连接操作消息队列
        4. autoDelete:没有任何东西使用消息队列后,是否要删除队列
        • 消费者代码
        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
        package com.bi.spring.demomq;

        import com.rabbitmq.client.*;
        import lombok.extern.slf4j.Slf4j;

        import java.nio.charset.StandardCharsets;

        @Slf4j
        public class Consumer {

        private static final String QUEUE_NAME = "hello_mq";

        public static void main(String[] args) {
        ConnectionFactory factory = new ConnectionFactory();
        factory.setHost("localhost");
        try {
        // 从工厂中创建连接
        Connection connection = factory.newConnection();
        // 获取通信通道
        Channel channel = connection.createChannel();
        // 声明一个队列 没有就会创建 有的话就不再创建
        channel.queueDeclare(QUEUE_NAME, false, false, false, null);
        System.out.println("Consumer waiting for producer to send message");
        // 定义如何处理消息 这个是函数式接口
        DeliverCallback deliverCallback = (consumerTag, delivery) -> {
        String message = new String(delivery.getBody(), StandardCharsets.UTF_8);
        System.out.println("Consumer received message from producer,message:" + message);
        };
        // 消费消息
        channel.basicConsume(QUEUE_NAME, true, deliverCallback, consumer -> {
        });
        } catch (Exception e) {
        log.info("创建链接失败!E:" + e);
        }
        }
        }
      • 多消费者

        场景:多个系统同时去接受并处理任务,尤其是每个系统的处理能力有限
        一个生产者给一个队列发消息,多个消费者从这个队列取消息。1<->n

        1-n

        • 队列持久化
        1
        channel.queueDeclare(TASK_QUEUE_NAME, true, false, false, null);
        • 消息持久化
          指定 MessageProperties.PERSISTENT_TEXT_PLAIN 参数
        1
        2
        channel.basicPublish("", QUEUE_NAME, MessageProperties.PERSISTENT_TEXT_PLAIN,
        message.getBytes(StandardCharsets.UTF_8));
        • 生产者代码

        使用 Scanner 来模拟多次输入

        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.bi.spring.demomq;

        import com.rabbitmq.client.Channel;
        import com.rabbitmq.client.Connection;
        import com.rabbitmq.client.ConnectionFactory;
        import com.rabbitmq.client.MessageProperties;
        import lombok.extern.slf4j.Slf4j;

        import java.nio.charset.StandardCharsets;
        import java.util.Scanner;

        @Slf4j
        public class MultiProducer {

        private static final String QUEUE_NAME = "multi_queue";

        public static void main(String[] args) {
        ConnectionFactory factory = new ConnectionFactory();
        factory.setHost("localhost");
        try {
        // 从工厂中创建连接
        Connection connection = factory.newConnection();
        // 获取通信通道
        Channel channel = connection.createChannel();
        // 声明一个队列
        channel.queueDeclare(QUEUE_NAME, true, false, false, null);
        // 声明一个输入对象
        Scanner scanner = new Scanner(System.in);
        while (true) {
        // 输入信息
        String message = scanner.nextLine();
        // 发送消息
        channel.basicPublish("", QUEUE_NAME, MessageProperties.PERSISTENT_TEXT_PLAIN, message.getBytes(StandardCharsets.UTF_8));
        // 打印信息
        System.out.println("Producer send message: " + message);
        }
        } catch (Exception e) {
        log.info("创建链接失败!E:" + e);
        }
        }
        }
        • 消费者代码
        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
        package com.bi.spring.demomq;

        import com.rabbitmq.client.Channel;
        import com.rabbitmq.client.Connection;
        import com.rabbitmq.client.ConnectionFactory;
        import com.rabbitmq.client.DeliverCallback;
        import lombok.extern.slf4j.Slf4j;

        import java.nio.charset.StandardCharsets;

        @Slf4j
        public class MultiConsumer {

        private static final String QUEUE_NAME = "multi_queue";

        public static void main(String[] args) {
        try {
        ConnectionFactory factory = new ConnectionFactory();
        factory.setHost("localhost");
        // 从工厂中创建连接
        Connection connection = factory.newConnection();
        // 循环创建
        for (int i = 0; i < 2; i++) {
        // 获取通信通道
        Channel channel = connection.createChannel();
        // 声明一个队列 没有就会创建 有的话就不再创建
        channel.queueDeclare(QUEUE_NAME, true, false, false, null);
        System.out.println("Consumer waiting for producer to send message");
        // 定义如何处理消息 这个是函数式接口
        int finalI1 = i;
        DeliverCallback deliverCallback = (consumerTag, delivery) -> {
        try {
        String message = new String(delivery.getBody(), StandardCharsets.UTF_8);
        System.out.println("Consumer received message from producer,message:" + message + "-编号" + finalI1);
        // 确认消息
        channel.basicAck(delivery.getEnvelope().getDeliveryTag(), false);
        // 模拟处理能力
        Thread.sleep(10000);
        } catch (Exception e) {
        e.printStackTrace();
        // 消息失败策略
        channel.basicNack(delivery.getEnvelope().getDeliveryTag(), false, false);
        } finally {
        System.out.println("Consumer is Done!");
        // 如果没有失败最终都要确认消息
        channel.basicAck(delivery.getEnvelope().getDeliveryTag(), false);
        }
        };
        // 消费消息
        channel.basicConsume(QUEUE_NAME, true, deliverCallback, consumer -> {
        });
        }
        } catch (Exception e) {
        log.info("创建链接失败!E:" + e);
        }
        }
        }

        控制单个消费的处理任务的积压数,每个消费者最多同时处理的任务数

        1
        channel.basicQos(1);
      • 交换机

        一个生产者给多个队列发送消息,一个生产者 -> 多个队列
        交换机的作用:提供消息转发功能,类似与网络路由器
        问题:怎么把消息转发到不同的队列上,好让消费者从不同的队列消费

        绑定:交换机可以与队列进行关联,通过路由 key 进行绑定

        1
        channel.queueBind(QUEUE_NAME,EXCHANGE_NAME,"key");

        交换机

        • fanout 扇出、广播

          特点:消息会被转发到所有绑定到该交换机的队列
          场景:很适用于发布订阅的场景。比如写日志,可以多个系统间共享
          注意点:生产者和消费者需要绑定同一个交换机,并且要先创建队列,才能绑定

          fanout

          • 生产者代码
          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
          package com.bi.spring.demomq;

          import com.rabbitmq.client.Channel;
          import com.rabbitmq.client.Connection;
          import com.rabbitmq.client.ConnectionFactory;
          import com.rabbitmq.client.MessageProperties;
          import lombok.extern.slf4j.Slf4j;

          import java.nio.charset.StandardCharsets;
          import java.util.Scanner;

          @Slf4j
          public class FanoutProducer {

          private static final String EXCHANGE_NAME = "fanout-exchange";

          public static void main(String[] args) {
          ConnectionFactory factory = new ConnectionFactory();
          factory.setHost("localhost");
          try {
          // 从工厂中创建连接
          Connection connection = factory.newConnection();
          // 获取通信通道
          Channel channel = connection.createChannel();
          // 声明一个交换机
          channel.exchangeDeclare(EXCHANGE_NAME, "fanout");
          // 声明一个输入对象
          Scanner scanner = new Scanner(System.in);
          while (scanner.hasNext()) {
          // 输入信息
          String message = scanner.nextLine();
          // 发送消息
          channel.basicPublish(EXCHANGE_NAME, "", MessageProperties.PERSISTENT_TEXT_PLAIN, message.getBytes(StandardCharsets.UTF_8));
          // 打印信息
          System.out.println("Producer send message: " + message);
          }
          } catch (Exception e) {
          log.info("创建链接失败!E:" + e);
          }
          }
          }

          • 消费者代码
          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
          package com.bi.spring.demomq;

          import com.rabbitmq.client.Channel;
          import com.rabbitmq.client.Connection;
          import com.rabbitmq.client.ConnectionFactory;
          import com.rabbitmq.client.DeliverCallback;
          import lombok.extern.slf4j.Slf4j;

          import java.nio.charset.StandardCharsets;

          @Slf4j
          public class FanoutConsumer {

          private static final String EXCHANGE_NAME = "fanout-exchange";

          private static final String QUEUE_NAME1 = "fanout-queue-one";

          private static final String QUEUE_NAME2 = "fanout-queue-two";

          public static void main(String[] args) {
          try {
          ConnectionFactory factory = new ConnectionFactory();
          factory.setHost("localhost");
          // 从工厂中创建连接
          Connection connection = factory.newConnection();
          // 创建频道
          Channel channel1 = connection.createChannel();
          Channel channel2 = connection.createChannel();
          // 声明交换机
          channel1.exchangeDeclare(EXCHANGE_NAME, "fanout");
          // 创建队列 1
          channel1.queueDeclare(QUEUE_NAME1, true, false, false, null);
          channel1.queueBind(QUEUE_NAME1, EXCHANGE_NAME, "");
          // 创建队列2
          channel2.queueDeclare(QUEUE_NAME2, true, false, false, null);
          channel2.queueBind(QUEUE_NAME2, EXCHANGE_NAME, "");
          System.out.println("Consumer waiting for producer to send message");
          // 声明消息处理回调函数1
          DeliverCallback deliverCallback1 = (consumerTag, delivery) -> {
          String message = new String(delivery.getBody(), StandardCharsets.UTF_8);
          System.out.println("[1] Received " + message);
          };
          // 声明消息处理回调函数2
          DeliverCallback deliverCallback2 = (consumerTag, delivery) -> {
          String message = new String(delivery.getBody(), StandardCharsets.UTF_8);
          System.out.println("[2] Received " + message);
          };

          // 开启消费监听1
          channel1.basicConsume(QUEUE_NAME1, true, deliverCallback1, consumerTag -> {
          });
          // 开启消费监听2
          channel2.basicConsume(QUEUE_NAME2, true, deliverCallback2, consumerTag -> {
          });
          } catch (Exception e) {
          log.info("创建链接失败!E:" + e);
          }
          }
          }

          效果:所有的消费者都能够接收到消息

        • direct 交给特定的队列

          routingkey:可以让交换机和队列进行关联,可以指定让交换机给那个队列发送消息,路由键,控制消息要给那个队列

          特点:消息会根据路由键转发到指定的队列
          场景:特定的消息只交给特定的系统来处理
          绑定关系:完全匹配字符串

          direct

          多个队列可以绑定同一个路由键

          directT

          • 生产者代码
          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.bi.spring.demomq;

          import com.rabbitmq.client.Channel;
          import com.rabbitmq.client.Connection;
          import com.rabbitmq.client.ConnectionFactory;
          import com.rabbitmq.client.MessageProperties;
          import lombok.extern.slf4j.Slf4j;

          import java.nio.charset.StandardCharsets;
          import java.util.Scanner;

          @Slf4j
          public class DirectProducer {

          private static final String EXCHANGE_NAME = "direct-exchange";

          public static void main(String[] args) {
          ConnectionFactory factory = new ConnectionFactory();
          factory.setHost("localhost");
          try {
          // 从工厂中创建连接
          Connection connection = factory.newConnection();
          // 获取通信通道
          Channel channel = connection.createChannel();
          // 声明一个交换机
          channel.exchangeDeclare(EXCHANGE_NAME, "direct");
          // 声明一个输入对象
          Scanner scanner = new Scanner(System.in);
          while (scanner.hasNext()) {
          // 输入信息
          String userInput = scanner.nextLine();
          String[] s = userInput.split(" ");
          // 消息
          String message = s[0];
          // 路由键
          String routingKey = s[1];
          // 发送消息
          channel.basicPublish(EXCHANGE_NAME, routingKey, MessageProperties.PERSISTENT_TEXT_PLAIN, message.getBytes(StandardCharsets.UTF_8));
          // 打印信息
          System.out.println("Producer send message: " + message + "- routingKey" + routingKey);
          }
          } catch (Exception e) {
          log.info("创建链接失败!E:" + e);
          }
          }
          }
          • 消费者代码
          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
          package com.bi.spring.demomq;

          import com.rabbitmq.client.Channel;
          import com.rabbitmq.client.Connection;
          import com.rabbitmq.client.ConnectionFactory;
          import com.rabbitmq.client.DeliverCallback;
          import lombok.extern.slf4j.Slf4j;

          import java.nio.charset.StandardCharsets;

          @Slf4j
          public class DirectConsumer {

          private static final String EXCHANGE_NAME = "direct-exchange";

          private static final String QUEUE_NAME1 = "direct-queue-one";

          private static final String QUEUE_NAME2 = "direct-queue-two";

          public static void main(String[] args) {
          try {
          ConnectionFactory factory = new ConnectionFactory();
          factory.setHost("localhost");
          // 从工厂中创建连接
          Connection connection = factory.newConnection();
          // 创建频道
          Channel channel1 = connection.createChannel();
          Channel channel2 = connection.createChannel();
          // 声明交换机
          channel1.exchangeDeclare(EXCHANGE_NAME, "fanout");
          // 创建队列 1
          channel1.queueDeclare(QUEUE_NAME1, true, false, false, null);
          // 绑定路由键1
          channel1.queueBind(QUEUE_NAME1, EXCHANGE_NAME, "routing_key_one");
          // 创建队列2
          channel2.queueDeclare(QUEUE_NAME2, true, false, false, null);
          // 绑定路由键2
          channel2.queueBind(QUEUE_NAME2, EXCHANGE_NAME, "routing_key_two");
          System.out.println("Consumer waiting for producer to send message");
          // 声明消息处理回调函数1
          DeliverCallback deliverCallback1 = (consumerTag, delivery) -> {
          String message = new String(delivery.getBody(), StandardCharsets.UTF_8);
          System.out.println("[1] Received " + message);
          };
          // 声明消息处理回调函数2
          DeliverCallback deliverCallback2 = (consumerTag, delivery) -> {
          String message = new String(delivery.getBody(), StandardCharsets.UTF_8);
          System.out.println("[2] Received " + message);
          };

          // 开启消费监听1
          channel1.basicConsume(QUEUE_NAME1, true, deliverCallback1, consumerTag -> {
          });
          // 开启消费监听2
          channel2.basicConsume(QUEUE_NAME2, true, deliverCallback2, consumerTag -> {
          });
          } catch (Exception e) {
          log.info("创建链接失败!E:" + e);
          }
          }
          }

        • topic 会模糊发送特定的队列

          特点:消息会根据一个模糊的路由键转发到指定的队列
          场景:特定的一类的消息可以交给特定的一类系统来处理
          绑定关系:可以模糊匹配多个绑定

          1. *:匹配一个单词,比如 *.orange,可匹配项有 a.orange、b.orange 缺点就是必须要有匹配项不然匹配不到
          2. #: 匹配零个或者多个单词,比如 a.#,可匹配项有 a.a、a.b、a.a.a 也可以不匹配,缺点就是太过灵活

          topic

          示例图

          • 生产者代码
          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
          package com.bi.spring.demomq;

          import com.rabbitmq.client.Channel;
          import com.rabbitmq.client.Connection;
          import com.rabbitmq.client.ConnectionFactory;
          import com.rabbitmq.client.MessageProperties;
          import lombok.extern.slf4j.Slf4j;

          import java.nio.charset.StandardCharsets;
          import java.util.Scanner;

          @Slf4j
          public class TopicProducer {

          private static final String EXCHANGE_NAME = "topic-exchange";

          public static void main(String[] args) {
          ConnectionFactory factory = new ConnectionFactory();
          factory.setHost("localhost");
          try {
          // 从工厂中创建连接
          Connection connection = factory.newConnection();
          // 获取通信通道
          Channel channel = connection.createChannel();
          // 声明一个交换机
          channel.exchangeDeclare(EXCHANGE_NAME, "topic");
          // 声明一个输入对象
          Scanner scanner = new Scanner(System.in);
          while (scanner.hasNext()) {
          // 输入信息
          String userInput = scanner.nextLine();
          String[] s = userInput.split(" ");
          if (s.length < 1) {
          continue;
          }
          // 消息
          String message = s[0];
          // 路由键
          String routingKey = s[1];
          // 发送消息
          channel.basicPublish(EXCHANGE_NAME, routingKey, MessageProperties.PERSISTENT_TEXT_PLAIN, message.getBytes(StandardCharsets.UTF_8));
          // 打印信息
          System.out.println("Producer send message: " + message + "- routingKey" + routingKey);
          }
          } catch (Exception e) {
          log.info("创建链接失败!E:" + e);
          }
          }
          }

          • 消费者代码
          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

          package com.bi.spring.demomq;

          import com.rabbitmq.client.Channel;
          import com.rabbitmq.client.Connection;
          import com.rabbitmq.client.ConnectionFactory;
          import com.rabbitmq.client.DeliverCallback;
          import lombok.extern.slf4j.Slf4j;

          import java.nio.charset.StandardCharsets;

          @Slf4j
          public class TopicConsumer {

          private static final String EXCHANGE_NAME = "topic-exchange";

          private static final String QUEUE_NAME1 = "topic-queue-one";

          private static final String QUEUE_NAME2 = "topic-queue-two";

          private static final String QUEUE_NAME3 = "topic-queue-three";

          public static void main(String[] args) {
          try {
          ConnectionFactory factory = new ConnectionFactory();
          factory.setHost("localhost");
          // 从工厂中创建连接
          Connection connection = factory.newConnection();
          // 创建频道
          Channel channel = connection.createChannel();
          // 声明交换机
          channel.exchangeDeclare(EXCHANGE_NAME, "topic");
          // 创建队列 1
          channel.queueDeclare(QUEUE_NAME1, true, false, false, null);
          // 绑定路由键1
          channel.queueBind(QUEUE_NAME1, EXCHANGE_NAME, "#.队列1.#");
          // 创建队列2
          channel.queueDeclare(QUEUE_NAME2, true, false, false, null);
          // 绑定路由键2
          channel.queueBind(QUEUE_NAME2, EXCHANGE_NAME, "#.队列2.#");
          // 创建队列3
          channel.queueDeclare(QUEUE_NAME3, true, false, false, null);
          // 绑定路由键3
          channel.queueBind(QUEUE_NAME3, EXCHANGE_NAME, "#.队列3.#");

          System.out.println("Consumer waiting for producer to send message");
          // 声明消息处理回调函数1
          DeliverCallback deliverCallback1 = (consumerTag, delivery) -> {
          String message = new String(delivery.getBody(), StandardCharsets.UTF_8);
          System.out.println("[1] Received " + message);
          };
          // 声明消息处理回调函数2
          DeliverCallback deliverCallback2 = (consumerTag, delivery) -> {
          String message = new String(delivery.getBody(), StandardCharsets.UTF_8);
          System.out.println("[2] Received " + message);
          };
          // 声明消息处理回调函数2
          DeliverCallback deliverCallback3 = (consumerTag, delivery) -> {
          String message = new String(delivery.getBody(), StandardCharsets.UTF_8);
          System.out.println("[3] Received " + message);
          };

          // 开启消费监听1
          channel.basicConsume(QUEUE_NAME1, true, deliverCallback1, consumerTag -> {
          });
          // 开启消费监听2
          channel.basicConsume(QUEUE_NAME2, true, deliverCallback2, consumerTag -> {
          });
          // 开启消费监听3
          channel.basicConsume(QUEUE_NAME3, true, deliverCallback3, consumerTag -> {
          });
          } catch (Exception e) {
          log.info("创建链接失败!E:" + e);
          }
          }
          }

        • headers 了解即可

          类似主题和直接交换机,可以根据 headers 中的内容来指定发送到那个队列
          由于性能较差且复杂 不作使用

      • 核心特性
        • 消息过期机制

          可以给每条消息指定一个有效期,一段时间内未被消费者处理,就过期了

          使用场景:清理过期数据、模拟延迟队列的实现,专门让某个程序去处理过期请求

          给队列中的所有消息指定过期时间

          1
          2
          3
          4
          5
          // 创建队列,指定消息过期参数
          Map<String,Object> args = new HasMap<>();
          args.put("x-message-ttl",5000);
          // args指定参数
          channel.queueDeclare(QUEUE_NAME,false,false,false,args);

          如果在过期时间内,还未收到消费者来取消息,消息才会过期
          如果消息已经接收到,但是没有确认,是不会过期的

          如果消息处于待消费状态并且到达过期时间后,消息将会被标记为过期。但是,如果消息已经被消费者消费,并且正在被处理中,即使过期时间到了,消息依旧会被正常处理

          消费者代码:

          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
          package com.bi.spring.demomq;

          import com.rabbitmq.client.Channel;
          import com.rabbitmq.client.Connection;
          import com.rabbitmq.client.ConnectionFactory;
          import com.rabbitmq.client.DeliverCallback;
          import lombok.extern.slf4j.Slf4j;

          import java.nio.charset.StandardCharsets;
          import java.util.HashMap;
          import java.util.Map;

          @Slf4j
          public class TtlConsumer {

          private static final String QUEUE_NAME = "ttl_queue";

          public static void main(String[] args) {
          ConnectionFactory factory = new ConnectionFactory();
          factory.setHost("localhost");
          try {
          // 从工厂中创建连接
          Connection connection = factory.newConnection();
          // 获取通信通道
          Channel channel = connection.createChannel();
          Map<String, Object> canshu = new HashMap<>();
          canshu.put("x-message-ttl", 5000);
          // 声明一个队列 没有就会创建 有的话就不再创建 放入参数
          channel.queueDeclare(QUEUE_NAME, false, false, false, canshu);
          System.out.println("Consumer waiting for producer to send message");
          // 定义如何处理消息 这个是函数式接口
          DeliverCallback deliverCallback = (consumerTag, delivery) -> {
          String message = new String(delivery.getBody(), StandardCharsets.UTF_8);
          System.out.println("Consumer received message from producer,message:" + message);
          };
          // 消费消息
          channel.basicConsume(QUEUE_NAME, true, deliverCallback, consumer -> {
          });
          } catch (Exception e) {
          log.info("创建链接失败!E:" + e);
          }
          }
          }

          生产者代码就正常取出消息消费即可,不作代码演示。

        • 给某条消息指定过期时间

          语法

          1
          2
          AMQP.BasicProperties properties = new AMQP.BasicProperties().builder().expiration("1000").build();
          channel.basicPublish("", QUEUE_NAME, properties, message.getBytes(StandardCharsets.UTF_8));
        • 消息确认机制

          为了保证消息成功被消费,rabbitMQ 提供了消息确认机制,当消费者接收到消息后,要给予反馈:

          ack:消费成功
          nack:消费失败
          reject:拒绝

          只有告诉 rabbitMQ 服务器消费成功,服务器才会放心的移除消息
          支持配置 autoack,会自动执行 ack 命令,接收到消息就立刻成功
          一般情况下 建议 autoack 设置为 false 要根据实际情况,去手动确认

          指定确认某条消息:

          1
          channel.basicAck(delivery.getEnvelope().getDeliveryTag(),false);

          指定拒绝某条消息
          第三个参数表示是否要重试

          1
          channel.basicNack(delivery.getEnvelope().getDeliveryTag(),false,false);
        • 死信队列

          官方文档: https://www.rabbitmq.com/dlx.html

          为了保证消息的可靠性,比如每条消息都成功消费,需要提供一个容错机制,即: 失败的消息怎么处理?

          死信: 过期的消息、拒收的消息、消息队列满了、处理失败的消息的统称死信队列:专门处理死信的队列

          死信队列:专门处理死信的队列(注意,它就是一个普队列,只不过是专门用来处理死信的,你甚至可以理解这个队列的名称叫“死信队列”)

          死信交换机:专门给死信队列转发消息的交换机(注意,它就是一个普通交换机,只不过是专门给死信队列发消息而已,理解为这个交换机的名称就叫 “死信交换机”)。也存在路由绑定

          死信可以通过死信交换机绑定到死信队列。

          • 1.创建死信交换机并绑定关系

          • 2.给失败之后需要容错处理的消息队列绑定死信交换机

          1
          2
          3
          4
          5
          6
          7
          8
          9
          10
          // 指定死信队列参数
          Map<String,Object> args = new HashMap<>();
          // 要绑定到哪个交换机 等于额外绑定一个死信队列 处理特殊情况
          args.put("x-dead-letter-exchange",DEAD EXCHANGE NAME);
          // 指定死信要转发到哪个死信队列
          args.put("x-dead-letter-routing-key""sixin_key");
          // 创建队列,随机分配一个队列名称
          String queueName = "sixi_queue";
          channel.queueDeclare(queueName, truefalse, false, args);
          channel.queueBind(queueName, EXCHANGE NAME,"sixinkey_1");
          • 3.可以给要容错的队列指定死信之后的转发规则,死信应该再换发到那个死信队列
          1
          args.put("x-dead-letter-routing-key","sixin_de_queue")
          • 4.可以通过程序来读取死信队列中的消息,从而进行处理
          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
          Connection connection = factory.newConnection();
          // 获取通信通道
          Channel channel = connection.createChannel();
          // 声明死信交换机
          channel.exchangeDeclare(DEAD_EXCHANGE_NAME, "direct");
          // 创建死信队列
          channel.queueDeclare(DEAD_QUEUE_NAME1, true, false, false, null);
          // 将死信队列绑定到死信交换机并授予routingKey
          channel.queueBind(DEAD_QUEUE_NAME1, DEAD_EXCHANGE_NAME, "routing-key-dead-queue-one");
          // 创建死信队列2
          channel.queueDeclare(DEAD_QUEUE_NAME2, true, false, false, null);
          // 将死信队列绑定到死信交换机并授予routingKey 2
          channel.queueBind(DEAD_QUEUE_NAME2, DEAD_EXCHANGE_NAME, "routing-key-dead-queue-two");
          // 开启死信队列的消费监听 1
          DeliverCallback deliverCallback1 = (consumerTag, delivery) -> {
          String message = new String(delivery.getBody(), StandardCharsets.UTF_8);
          // 拒绝消息
          channel.basicNack(delivery.getEnvelope().getDeliveryTag(), false, false);
          // 打印接收消息
          System.out.println("[1] Dead receive message:" + message + "-key:" + delivery.getEnvelope().getRoutingKey());
          };
          channel.basicConsume(DEAD_QUEUE_NAME1, false, deliverCallback1, consumerTag -> {
          });
          // 开启死信队列的消费监听 2
          DeliverCallback deliverCallback2 = (consumerTag, delivery) -> {
          String message = new String(delivery.getBody(), StandardCharsets.UTF_8);
          // 拒绝消息
          channel.basicNack(delivery.getEnvelope().getDeliveryTag(), false, false);
          // 打印接收消息
          System.out.println("[2] Dead receive message:" + message + "-key:" + delivery.getEnvelope().getRoutingKey());
          };
          channel.basicConsume(DEAD_QUEUE_NAME2, false, deliverCallback2, consumerTag -> {
          });
          • 5.生产者代码
          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
          package com.bi.spring.demomq;

          import com.rabbitmq.client.*;
          import lombok.extern.slf4j.Slf4j;

          import java.nio.charset.StandardCharsets;
          import java.util.Scanner;

          @Slf4j
          public class DlxDirectProducer {

          private static final String DEAD_EXCHANGE_NAME = "dlx-direct-exchange";

          public static final String DEAD_QUEUE_NAME1 = "dlx-one-queue";

          public static final String DEAD_QUEUE_NAME2 = "dlx-two-queue";

          private static final String WORK_EXCHANGE_NAME = "work-direct-exchange";

          public static void main(String[] args) {
          ConnectionFactory factory = new ConnectionFactory();
          factory.setHost("localhost");
          try {
          // 从工厂中创建连接
          Connection connection = factory.newConnection();
          // 获取通信通道
          Channel channel = connection.createChannel();
          // 声明死信交换机
          channel.exchangeDeclare(DEAD_EXCHANGE_NAME, "direct");
          // 创建死信队列
          channel.queueDeclare(DEAD_QUEUE_NAME1, true, false, false, null);
          // 将死信队列绑定到死信交换机并授予routingKey
          channel.queueBind(DEAD_QUEUE_NAME1, DEAD_EXCHANGE_NAME, "routing-key-dead-queue-one");
          // 创建死信队列2
          channel.queueDeclare(DEAD_QUEUE_NAME2, true, false, false, null);
          // 将死信队列绑定到死信交换机并授予routingKey 2
          channel.queueBind(DEAD_QUEUE_NAME2, DEAD_EXCHANGE_NAME, "routing-key-dead-queue-two");
          // 开启死信队列的消费监听 1
          DeliverCallback deliverCallback1 = (consumerTag, delivery) -> {
          String message = new String(delivery.getBody(), StandardCharsets.UTF_8);
          // 拒绝消息
          channel.basicNack(delivery.getEnvelope().getDeliveryTag(), false, false);
          // 打印接收消息
          System.out.println("[1] Dead receive message:" + message + "-key:" + delivery.getEnvelope().getRoutingKey());
          };
          channel.basicConsume(DEAD_QUEUE_NAME1, false, deliverCallback1, consumerTag -> {
          });
          // 开启死信队列的消费监听 2
          DeliverCallback deliverCallback2 = (consumerTag, delivery) -> {
          String message = new String(delivery.getBody(), StandardCharsets.UTF_8);
          // 拒绝消息
          channel.basicNack(delivery.getEnvelope().getDeliveryTag(), false, false);
          // 打印接收消息
          System.out.println("[2] Dead receive message:" + message + "-key:" + delivery.getEnvelope().getRoutingKey());
          };
          channel.basicConsume(DEAD_QUEUE_NAME2, false, deliverCallback2, consumerTag -> {
          });
          // 生产消息
          Scanner scanner = new Scanner(System.in);
          while (scanner.hasNext()) {
          String userInput = scanner.nextLine();
          String[] s = userInput.split(" ");
          if (s.length < 1) {
          continue;
          }
          String message = s[0];
          String routingKey = s[1];
          channel.basicPublish(WORK_EXCHANGE_NAME, routingKey, null, message.getBytes(StandardCharsets.UTF_8));
          System.out.println("[x] Send message:" + message + " with routingKey:" + routingKey);
          }
          } catch (Exception e) {
          log.info("创建链接失败!E:" + e);
          }
          }
          }

          • 6.消费者代码
          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
          package com.bi.spring.demomq;

          import com.rabbitmq.client.Channel;
          import com.rabbitmq.client.Connection;
          import com.rabbitmq.client.ConnectionFactory;
          import com.rabbitmq.client.DeliverCallback;
          import lombok.extern.slf4j.Slf4j;

          import java.nio.charset.StandardCharsets;
          import java.util.HashMap;
          import java.util.Map;

          @Slf4j
          public class DlxDirectConsumer {

          private static final String DEAD_EXCHANGE_NAME = "dlx-direct-exchange";

          public static final String WORK_QUEUE_NAME1 = "work-dlx-one-queue";

          public static final String WORK_QUEUE_NAME2 = "work-dlx-two-queue";

          private static final String WORK_EXCHANGE_NAME = "work-direct-exchange";

          public static void main(String[] args) {
          ConnectionFactory factory = new ConnectionFactory();
          factory.setHost("localhost");
          try {
          // 从工厂中创建连接
          Connection connection = factory.newConnection();
          // 获取通信通道
          Channel channel = connection.createChannel();
          // 声明工作交换机
          channel.exchangeDeclare(WORK_EXCHANGE_NAME, "direct");
          // 指定死信队列参数
          Map<String, Object> parameter = new HashMap<>();
          // 绑定交换机
          parameter.put("x-dead-letter-exchange", DEAD_EXCHANGE_NAME);
          // 指定死信要转发到那个队列
          parameter.put("x-dead-letter-routing-key", "routing-key-dead-queue-one");
          // 创建队列1 以参数形式多绑定一个死信交换机
          channel.queueDeclare(WORK_QUEUE_NAME1, true, false, false, parameter);
          channel.queueBind(WORK_QUEUE_NAME1, WORK_EXCHANGE_NAME, "work-queue-one");
          // 同样也给2进行绑定
          Map<String, Object> parameter2 = new HashMap<>();
          parameter2.put("x-dead-letter-exchange", DEAD_EXCHANGE_NAME);
          parameter2.put("x-dead-letter-routing-key", "routing-key-dead-queue-two");
          // 创建队列2 以参数形式多绑定一个死信交换机
          channel.queueDeclare(WORK_QUEUE_NAME2, true, false, false, parameter);
          channel.queueBind(WORK_QUEUE_NAME2, WORK_EXCHANGE_NAME, "work-queue-two");
          System.out.println("Consumer waiting for producer to send message");
          // 开启工作队列的消费监听 1
          DeliverCallback deliverCallback1 = (consumerTag, delivery) -> {
          String message = new String(delivery.getBody(), StandardCharsets.UTF_8);
          // 拒绝消息
          channel.basicNack(delivery.getEnvelope().getDeliveryTag(), false, false);
          // 打印接收消息
          System.out.println("[1] Dead receive message:" + message + "-key:" + delivery.getEnvelope().getRoutingKey());
          };
          channel.basicConsume(WORK_QUEUE_NAME1, false, deliverCallback1, consumerTag -> {
          });
          // 开启工作队列的消费监听 2
          DeliverCallback deliverCallback2 = (consumerTag, delivery) -> {
          String message = new String(delivery.getBody(), StandardCharsets.UTF_8);
          // 拒绝消息
          channel.basicNack(delivery.getEnvelope().getDeliveryTag(), false, false);
          // 打印接收消息
          System.out.println("[2] Dead receive message:" + message + "-key:" + delivery.getEnvelope().getRoutingKey());
          };
          channel.basicConsume(WORK_QUEUE_NAME2, false, deliverCallback2, consumerTag -> {
          });
          } catch (Exception e) {
          log.info("创建链接失败!E:" + e);
          }
          }
          }

      • rabbit 重点知识
        • 消息队列的概念、模型、应用场景

        • 交换机的类别、路由绑定关系

        • 消息可靠性

          消息确认机制(ack、nack、reject)
          消息持久化(durable)
          消息过期机制
          死信队列

        • 延迟队列 类似于 死信队列

        • 顺序消费、消费幂等性(做了解)

        • 可扩展性

          集群
          故障的回复机制
          镜像

        • 运维监控告警(做了解)

    • RabbitMQ 项目实战

      • 项目中如何使用 RabbitMQ?

        • 使用官方的客户端

        优点:兼容性好,换语言成本低,比较灵活
        缺点:太灵活,要自己去做一些限制或者是一些事情。比如要自己维护管理连接等等。

        • 使用封装好的客户端,比如 Spring Boot RabbitMQ Starter

        优点:简单易用,直接配置直接用,更方便地去管理连接
        缺点:封装的太好了,没基础知识很难去使用或者是配置

      • 基础实战

        • 引入依赖

          引入时候一定要选择与自己当前使用的 Spring Boot 版本一致

          1
          2
          3
          4
          5
          6
          <!-- https://mvnrepository.com/artifact/org.springframework.boot/spring-boot-starter-amqp -->
          <dependency>
          <groupId>org.springframework.boot</groupId>
          <artifactId>spring-boot-starter-amqp</artifactId>
          <version>2.7.2</version>
          </dependency>
        • 在 yml 中进行配置

          1
          2
          3
          4
          5
          6
          spring:
          rabbmit:
          host: localhost
          port: 5672
          username: guest
          passowrd: guest
        • 创建交换机和队列

          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
          package com.bi.spring.bizmq;

          import com.bi.spring.config.RabbitMQConfig;
          import com.rabbitmq.client.Channel;
          import com.rabbitmq.client.Connection;
          import com.rabbitmq.client.ConnectionFactory;
          import lombok.extern.slf4j.Slf4j;
          import org.springframework.beans.factory.annotation.Autowired;
          import org.springframework.beans.factory.annotation.Value;
          import org.springframework.stereotype.Component;
          import org.springframework.stereotype.Service;

          /**
          * 用于创建测试程序用到的交换机和队列(只用在程序启动前执行一次)
          */
          @Slf4j
          public class MqInitMain {

          private static final String EXCHANGE_NAME = "bi-exchange";

          public static void main(String[] args) {
          try {
          // 创建连接
          ConnectionFactory factory = new ConnectionFactory();
          // 设置连接属性
          factory.setHost("localhost");
          Connection connection = factory.newConnection();
          Channel channel = connection.createChannel();
          channel.exchangeDeclare(EXCHANGE_NAME, "direct");
          // 创建队列,随机分配一个队列名称
          String queueName = "test_queue";
          channel.queueDeclare(queueName, true, false, false, null);
          channel.queueBind(queueName, EXCHANGE_NAME, "bi-queue-key");
          } catch (Exception e) {
          log.info("初始化失败:" + e);
          }
          }
          }
          • 生产者代码
          1
          2
          3
          4
          5
          6
          7
          8
          9
          10
          11
          12
          13
          14
          15
          16
          17
          package com.bi.spring.bizmq;

          import org.springframework.amqp.rabbit.core.RabbitTemplate;
          import org.springframework.stereotype.Component;

          import javax.annotation.Resource;

          @Component
          public class MyMessageProducer {

          @Resource
          private RabbitTemplate rabbitTemplate;

          public void sendMessage(String exchange, String routingKey, String message) {
          rabbitTemplate.convertAndSend(exchange, routingKey, 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
          package com.bi.spring.bizmq;

          import com.rabbitmq.client.Channel;
          import lombok.SneakyThrows;
          import lombok.extern.slf4j.Slf4j;
          import org.springframework.amqp.rabbit.annotation.RabbitListener;
          import org.springframework.amqp.support.AmqpHeaders;
          import org.springframework.messaging.handler.annotation.Header;
          import org.springframework.stereotype.Component;

          @Component
          @Slf4j
          public class MyMessageConsumer {

          /**
          * 指定程序监听的消息队列和确认机制
          *
          * @param message
          * @param channel
          * @param deliveryTag
          */
          @SneakyThrows
          @RabbitListener(queues = {"bi_queue"}, ackMode = "MANUAL")
          public void receiveMessage(String message, Channel channel, @Header(AmqpHeaders.DELIVERY_TAG) long deliveryTag) {
          log.info("receiveMessage message = {}", message);
          // 确认收到消息
          channel.basicAck(deliveryTag, false);
          }

          }
          • 单元测试
          1
          2
          3
          4
          5
          6
          7
          8
          9
          10
          11
          12
          13
          14
          15
          16
          17
          18
          19
          20
          package com.bi.spring.bizmq;

          import org.junit.jupiter.api.Test;
          import org.springframework.beans.factory.annotation.Autowired;
          import org.springframework.boot.test.context.SpringBootTest;


          @SpringBootTest
          class MyMessageConsumerTest {

          @Autowired
          private MyMessageProducer myMessageProducer;

          @Test
          void receiveMessage() {
          myMessageProducer.sendMessage("bi-exchange", "bi-queue-key", "消息进来了");
          }

          }

          • 测试结果

      • 智能 BI 改造

        按照上一期的做法,是把用户提交的任务放入本地线程池中,在线程池内排队,但是程序如果中断了,任务就会丢失,就丢了,并且还不会去修改状态。

        • 实现步骤
          • 创建交换机和队列 程序启动前要先启动一次初始化

            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.bi.spring.bizmq;

            import com.rabbitmq.client.Channel;
            import com.rabbitmq.client.Connection;
            import com.rabbitmq.client.ConnectionFactory;

            /**
            * 初始化MQ
            */
            public class BiInitMain {

            public static void main(String[] args) {
            try {
            ConnectionFactory factory = new ConnectionFactory();
            factory.setHost("xxxx");
            factory.setPort(xxxx);
            factory.setUsername("xxxx");
            factory.setPassword("xxxx");
            Connection connection = factory.newConnection();
            Channel channel = connection.createChannel();
            // 声明交换机 类型为直接交换机
            channel.exchangeDeclare(BiMqConstant.BI_EXCHANGE_NAME, "direct");
            // 声明队列
            channel.queueDeclare(BiMqConstant.BI_QUEUE_NAME, true, false, false, null);
            // 绑定队列到交换机
            channel.queueBind(BiMqConstant.BI_QUEUE_NAME, BiMqConstant.BI_EXCHANGE_NAME, BiMqConstant.BI_ROUTING_KEY);
            // 限制任务数
            channel.basicQos(1);
            } catch (Exception e) {
            System.out.println("MQ初始化失败!");
            }

            }
            }

          • 重构 async 接口为新的 API

            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
            @Autowired
            private BiMessageProducer biMessageProducer;

            /**
            * AI异步(消息队列)对话分析
            *
            * @param multipartFile
            * @param genChartByAiRequest
            * @param request
            * @return
            */
            @PostMapping("/gen/async/mq")
            public BaseResponse<String> genChartByAIUseAsyncMQ(@RequestPart("file") MultipartFile multipartFile,
            GenChartByAiRequest genChartByAiRequest, HttpServletRequest request) {
            // 参数获取
            String goal = genChartByAiRequest.getGoal();
            String name = genChartByAiRequest.getName();
            String chartType = genChartByAiRequest.getChartType();
            // 参数校验
            ThrowUtils.throwIf(StringUtils.isEmpty(goal), ErrorCode.PARAMS_ERROR, "目标不能为空");
            ThrowUtils.throwIf(StringUtils.isEmpty(chartType), ErrorCode.PARAMS_ERROR, "类型不能为空");
            ThrowUtils.throwIf(StringUtils.isEmpty(name), ErrorCode.PARAMS_ERROR, "图表名称不能为空");
            // 文件校验
            long size = multipartFile.getSize();
            String originalFilename = multipartFile.getOriginalFilename();
            ThrowUtils.throwIf(size > MAX_FILE_SIZE, ErrorCode.PARAMS_ERROR, "文件过大不得超过1MB");
            ThrowUtils.throwIf(!FILE_NAME_LIST.contains(FileUtil.getSuffix(originalFilename)), ErrorCode.PARAMS_ERROR, "不支持此文件");
            // 判断是否登录
            User loginUser = userService.getLoginUser(request);
            ThrowUtils.throwIf(Objects.isNull(loginUser), ErrorCode.NOT_LOGIN_ERROR);
            // 对用户进行限流
            redisLimiterManager.doRateLimit("genChartByAI_" + loginUser.getId());
            // 用户消息拼接
            String dataStr = ExcelUtils.excelToCsv(multipartFile);
            // 先保存任务到数据库中
            Chart chart = new Chart();
            chart.setName(name);
            chart.setGoal(goal);
            chart.setChartData(dataStr);
            chart.setChartType(chartType);
            chart.setStatus("wait");
            chart.setExecMessage("任务等待中");
            chart.setUserId(loginUser.getId());
            boolean firstChartSave = chartService.save(chart);
            ThrowUtils.throwIf(!firstChartSave, ErrorCode.PARAMS_ERROR, "任务提交失败");
            // 任务如果保存成功 将ID步入消息队列中
            biMessageProducer.sendMessage(String.valueOf(chart.getId()));
            return ResultUtils.success("任务提交成功!");
            }
          • 将线程池中的代码进行重构

            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
            /**
            * 监听消息消费
            *
            * @param message
            * @param channel
            * @param deliverTag
            */
            @RabbitListener(queues = {BiMqConstant.BI_QUEUE_NAME}, ackMode = "MANUAL")
            public void receiveMessage(String message, Channel channel, @Header(AmqpHeaders.DELIVERY_TAG) long deliverTag) {
            try {
            log.info("receiveMessage: {}", message);
            if (StringUtils.isBlank(message)) {
            // 参数 消息的tag 是否要全部拒绝 false为只拒绝当前 是否要重入队列 false为不入
            channel.basicNack(deliverTag, false, false);
            throw new BusinessException(ErrorCode.SYSTEM_ERROR);
            }
            Long id = Long.parseLong(message);
            Chart chart = chartService.getById(id);
            if (Objects.isNull(chart)) {
            channel.basicNack(deliverTag, false, false);
            throw new BusinessException(ErrorCode.SYSTEM_ERROR);
            }
            // 将任务转入任务队列中分配执行
            CompletableFuture.runAsync(() -> {
            Chart updateFirstChart = new Chart();
            updateFirstChart.setId(chart.getId());
            updateFirstChart.setStatus("running");
            updateFirstChart.setExecMessage("任务正在执行中");
            boolean isFirstUpdate = chartService.updateById(updateFirstChart);
            if (!isFirstUpdate) {
            try {
            channel.basicNack(deliverTag, false, false);
            } catch (IOException e) {
            log.error("消息拒绝失败!");
            }
            handleChartUpdateError(chart.getId());
            return;
            }
            // AI接口服务
            String content = aiManager.doChat(buildUserMessage(chart));
            String[] splits = content.split("【【【【【");
            if (splits.length != 3) {
            try {
            channel.basicNack(deliverTag, false, false);
            } catch (IOException e) {
            log.error("消息拒绝失败!");
            }
            throw new RuntimeException("AI 生成错误");
            }
            Chart updateChartData = new Chart();
            updateChartData.setId(chart.getId());
            updateChartData.setStatus("succeed");
            updateFirstChart.setExecMessage("任务已经完成");
            updateChartData.setGenResult(splits[2].trim());
            updateChartData.setGenChart(splits[1].trim());
            boolean isUpdateChartData = chartService.updateById(updateChartData);
            if (!isUpdateChartData) {
            try {
            channel.basicNack(deliverTag, false, false);
            } catch (IOException e) {
            log.error("消息拒绝失败!");
            }
            handleChartUpdateError(chart.getId());
            return;
            }
            }, threadPoolExecutor);
            // 成功执行后确认消息
            channel.basicAck(deliverTag, false);
            } catch (IOException e) {
            log.error("消息消费失败!");
            throw new BusinessException(ErrorCode.SYSTEM_ERROR);
            }
            }

          • 代码中仍然存在缺陷,被拒绝的消息要进行处理,并要修改任务状态

            修改初始化 MQ 代码

            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
            package com.bi.spring.bizmq;

            import com.rabbitmq.client.Channel;
            import com.rabbitmq.client.Connection;
            import com.rabbitmq.client.ConnectionFactory;

            import java.util.HashMap;
            import java.util.Map;

            /**
            * 初始化MQ
            */
            public class BiInitMain {

            public static void main(String[] args) {
            try {
            ConnectionFactory factory = new ConnectionFactory();
            factory.setHost("124.222.153.56");
            factory.setPort(5672);
            factory.setUsername("admin");
            factory.setPassword("admin");
            Connection connection = factory.newConnection();
            Channel channel = connection.createChannel();
            // 声明死信交换机 类型为直接交换机
            channel.exchangeDeclare(BiMqConstant.BI_DEAD_EXCHANGE_NAME, "direct");
            // 声明死信队列
            channel.queueDeclare(BiMqConstant.BI_DEAD_QUEUE_NAME, true, false, false, null);
            // 声明交换机 类型为直接交换机
            channel.exchangeDeclare(BiMqConstant.BI_EXCHANGE_NAME, "direct");
            // 指定死信队列的参数
            Map<String, Object> parameter = new HashMap<>();
            parameter.put("x-dead-letter-exchange", BiMqConstant.BI_DEAD_EXCHANGE_NAME);
            parameter.put("x-dead-letter-routing-key", BiMqConstant.BI_DEAD_ROUTING_KEY);
            // 指定过期时间
            parameter.put("x-message-ttl", 3000L);
            // 声明队列
            channel.queueDeclare(BiMqConstant.BI_QUEUE_NAME, true, false, false, parameter);
            // 绑定队列到交换机
            channel.queueBind(BiMqConstant.BI_QUEUE_NAME, BiMqConstant.BI_EXCHANGE_NAME, BiMqConstant.BI_ROUTING_KEY);
            // 绑定死信交换机和死信队列的routingKey
            channel.queueBind(BiMqConstant.BI_DEAD_QUEUE_NAME, BiMqConstant.BI_DEAD_EXCHANGE_NAME, BiMqConstant.BI_DEAD_ROUTING_KEY);
            // 限制任务数
            channel.basicQos(1);
            } catch (Exception e) {
            System.out.println("MQ初始化失败!");
            }
            }
            }

            添加死信队列消费监听

            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
            /**
            * 处理因为某种情况没修改状态
            *
            * @param message
            * @param channel
            * @param deliverTage
            */
            @RabbitListener(queues = {BiMqConstant.BI_DEAD_QUEUE_NAME}, ackMode = "MANUAL")
            public void receiveDeadMessage(String message, Channel channel, @Header(AmqpHeaders.DELIVERY_TAG) long deliverTage) {
            try {
            log.info("receiveMessage dead message:{}", message);
            if (StringUtils.isBlank(message)) {
            channel.basicAck(deliverTage, false);
            } else {
            Chart chart = chartService.getById(Long.parseLong(message));
            if (!Objects.isNull(chart)) {
            chart.setStatus("failed");
            chart.setExecMessage("更新图表状态失败!");
            boolean save = chartService.save(chart);
            if (!save) {
            // 再次失败 就判定为系统异常
            log.error("更新图表失败状态失败" + chart.getId() + "," + "更新图表状态失败!");
            throw new BusinessException(ErrorCode.SYSTEM_ERROR);
            }
            }
            }
            channel.basicAck(deliverTage, false);
            } catch (IOException e) {
            log.error("死信消息消费失败!");
            throw new BusinessException(ErrorCode.SYSTEM_ERROR);
            }
            }
        • 验证

          验证发现,如果程序中断了,并且消息未被手动确认或者拒绝,也就是没有 ack 或者 nack 无任何反应

          那么这条消息就会重新处于 ready 状态,系统恢复时消费者依然会监听到这条消息并重新消费。

        • 性能优化

          改造之后
          写了一个专门来接受消息的程序,处理任务
          如果程序中断了,消息未被确认,就会重发
          消息集中的发到消息队列,可以部署多个后端,从同一个地方取出消息做处理,实现了分布式的负载均衡

      • 今日后端收获

        后端项目到此就完成了

        1. 完成了 由 同步化->异步化->异步化(消息队列)的优化过程
        2. 学习了如何使用本地线程池来进行异步化接口,从而增强了接口的响应速度,熟悉了线程池参数的意义
        3. 学习了使用消息中间件 RbbitMQ 来作为消息队列,从原本的本地线程池导致数据丢失,生成状态由于程序的中断而导致丢失,卡状态无法反馈的问题。
        4. 此项目中收获最大的就是如何对业务进行解耦拆分,把核心处理业务和其他三方服务或者是耗时长的服务抽离,用户只需要提交任务,后端慢慢的去处理
        5. 初步的使用消息中间件接入项目中,从听理论到结尾、模糊到分析清楚,花费了大量时间。因此此项目对于我收获非常的巨大。
        6. 特别感谢鱼总能够分享技术,带领学习。
  • 今日前端开发

    • 更新后端开发文档 执行 openAPI

    • 创建新的页面 修改请求 API

      1
      2
      3
      4
      5
      6
      7
      8
      9
      10
      11
      12
      13
      14
      15
      16
      17
      18
      19
      20
      21
      22
      const onFinish = async (values: any) => {
      if (submitting) return;
      setSubmitting(true);
      const params = {
      ...values,
      file: undefined,
      };
      try {
      const res = await genChartByAIUseAsyncMQUsingPOST(
      params,
      {},
      values?.file[0]?.originFileObj
      );
      if (res?.data) {
      message.success(res.data);
      form.resetFields();
      }
      } catch (e: any) {
      message.error("分析失败!");
      }
      setSubmitting(false);
      };
    • 提交后 DEBUG 观察消息是否进入后端

    • 到我的图表中观察状态


    • 今日前端收获

      由于前端不是侧重点,最后一期主要是对于后端的开发和优化
      页面架构都是前期搭建好的直接复制过来进行修改即可。