• RabbitMQ


    异步通信-RabbitMQ

    1、面临问题

        之前的请求都是同步调用虽然时效性较强,可以立即得到结果
        但是同步调用存在以下问题:
        1.耦合度高
        2.性能和吞吐能力下降
        3.有额外的资源消耗
        4.有级联失败问题

    2、解决方案-RabbitMQ

    2.1、简介

        MQ全称为Message Queue,即消息队列。“消息队列”是在消息的传输过程中保存消息的容器。它是典型的:生产者、消费者模型。生产者不断向消息队列中生产消息,消费者不断的从队列中获取消息。因为消息的生产和消费都是异步的,而且只关心消息的发送和接收,没有业务逻辑的侵入,这样就实现了生产者和消费者的解耦。
        官网:https://www.rabbitmq.com/
        参考:https://blog.csdn.net/kavito/article/details/91403659

    2.2、主流产品比较

    2.3、优缺点

    【优点】

        1
        3.调用间没有阻塞,不会造成无效的资源占用
        4.耦合度极低,每个服务都可以灵活插拔,可替换
        5.流量削峰:不管发布事件的流量波动多大,都由Broker接收,订阅者可以按照自己的速度去处理事件

    【缺点】

        1.架构复杂了,业务没有明显的流程线,不好管理
        2.需要依赖于Broker(MQ)的可靠、安全、性能

    2.4、运行原理

    2.5、运行流程

    【生产者发送消息流程】

    1、生产者和Broker建立TCP连接。
    ​
    2、生产者和Broker建立通道。
    ​
    3、生产者通过通道消息发送给Broker,由Exchange将消息进行转发。
    ​
    4、Exchange将消息转发到指定的Queue(队列)

    【消费者接收消息流程】

    1、消费者和Broker建立TCP连接
    ​
    2、消费者和Broker建立通道
    ​
    3、消费者监听指定的Queue(队列)
    ​
    4、当有消息到达Queue时Broker默认将消息推送给消费者。
    ​
    5、消费者接收到消息。
    ​
    6、ack回复

    3、使用步骤

    3.1、导入依赖

    1.    <!--AMQP依赖,包含RabbitMQ-->
    2.    <dependency>
    3.        <groupId>org.springframework.boot</groupId>
    4.        <artifactId>spring-boot-starter-amqp</artifactId>
    5.    </dependency>

    3.2、配置

        
    1. spring:
    2.     rabbitmq:
    3.       host: 192.168.248.222 # 主机名
    4.       port: 5672 # 端口
    5.       virtual-host: / # 虚拟主机
    6.       username: itcast # 用户名
    7.       password: 123321 # 密码

    注意:RabbitMQ是可以支持多个交换机绑定一个队列的,也就是说比如test1交换机可以绑定通道aa,然后test2可是可以绑定通道aa,但是如果通道aa在第一次链接被建立的是Topic类型的,那么交换机想要链接也要是Topic类型的,如下图 

    3.3、BasicQueue 简单队列模型

    生产者对应一个队列,这样的话就不需要我们设置交换机,会有一个默认的交换机

        生产者----消费者

    【publisher】

    1. @Autowired
    2.    private RabbitTemplate rabbitTemplate;
    3. @Test
    4.    public void testSimpleQueue() {
    5.        // 队列名称
    6.        String queueName = "simple.queue";
    7.        // 消息
    8.        String message = "hello, spring amqp!";
    9.        // 发送消息
    10.        rabbitTemplate.convertAndSend(queueName, message);
    11.   }

    【consumer】

    1. @RabbitListener(queues = "simple.queue")
    2.    public void listenSimpleQueueMessage(String msg) throws InterruptedException {
    3.        System.out.println("spring 消费者接收到消息:【" + msg + "】");
    4.   }

    3.4、WorkQueue 工作模式

    41431a18fc5f94d1e15dcad100fef0f8.png

        生产者----消费者1|消费者2

    【publisher】

    1.  /**
    2.         * workQueue
    3.         * 向队列中不停发送消息,模拟消息堆积。
    4.         */
    5.    @Test
    6.    public void testWorkQueue() throws InterruptedException {
    7.        // 队列名称
    8.        String queueName = "simple.queue";
    9.        // 消息
    10.        String message = "hello, message_";
    11.        for (int i = 0; i < 50; i++) {
    12.            // 发送消息
    13.            rabbitTemplate.convertAndSend(queueName, message + i);
    14.            Thread.sleep(20);
    15.       }
    16.   }

    【consumer】

    1.  @RabbitListener(queues = "simple.queue")
    2.    public void listenWorkQueue1(String msg) throws InterruptedException {
    3.        System.out.println("消费者1接收到消息:【" + msg + "】" + LocalTime.now());
    4.        Thread.sleep(20);
    5.   }
    6.    @RabbitListener(queues = "simple.queue")
    7.    public void listenWorkQueue2(String msg) throws InterruptedException {
    8.        System.err.println("消费者2........接收到消息:【" + msg + "】" + LocalTime.now());
    9.        Thread.sleep(200);
    10.   }

    这样的话我们只需要设置两个消费者都去指向这个管道就行,其实这样是不需要交换机的,因为就一个管道,有一个默认的交换机

    但是工作模式他有一个机制是“预取”机制,就是我们的消息弄到队列中去之后,消费者消费的时候,他是优先把队列的信息分配好之后,然后消费者再去消费,这样是不好的,因为比如上图的consumer2是效率不太行的,那么他执行的量和consumer1是一样的,这样不如按劳分配,能干的多干,不能干的少干,所以我们可以配置一个东西,让他们每次都只拿一个,拿完一个消费完了再去通道中拿。

    【yml】

    1. spring:
    2.     rabbitmq:
    3.       listener:
    4.         simple:
    5.           prefetch: 1 # 每次只能获取一条消息,处理完成才能获取下一个消息

    3.5、发布/订阅模式

    ① Fanout-扇出/广播模式

    先说fanout模式,也叫广播模式,生产者给交换机信息,交换机会将信息分为两份分别这两个管道,这样两个管道拿到的信息都是一样的:
    首先他是一个交换机对应两个管道,这是就需要我们创建交换机,然后将创建好的管道绑定到交换机中昂,那么我们在创建交换机的时候就可以绑定我们的管道,这样生产者在生产信息给交换机,交换机就会给管道分配信息

    fanout代码实现,有两种方式,一种是在config中配置,一种是在注解中写:

    1. @Configuration
    2. public class FanoutConfig {
    3.  
    4.    // 声明交换机 Fanout类型交换机
    5.    @Bean
    6.    public FanoutExchange fanoutExchange(){
    7.        return new FanoutExchange("itcast.fanout");
    8.   }
    9.    // 第1个队列
    10.    @Bean
    11.    public Queue fanoutQueue1(){
    12.        return new Queue("fanout.queue1");
    13.   }
    14.    // 绑定队列和交换机
    15.    @Bean
    16.    public Binding bindingQueue1(Queue fanoutQueue1, FanoutExchange fanoutExchange){
    17.        return BindingBuilder.bind(fanoutQueue1).to(fanoutExchange);
    18.   }
    19.    // 第1个队列
    20.    @Bean
    21.    public Queue fanoutQueue2(){
    22.        return new Queue("fanout.queue2");
    23.   }
    24.    // 绑定队列和交换机
    25.    @Bean
    26.    public Binding bindingQueue2(Queue fanoutQueue2, FanoutExchange fanoutExchange){
    27.        return BindingBuilder.bind(fanoutQueue2).to(fanoutExchange);
    28.   }
    29. }

    【publisher】

    1. // fanout
    2.    @Test
    3.    public void testFanoutExchange() {
    4.        // 队列名称
    5.        String exchangeName = "itcast.fanout";
    6.        // 消息
    7.        String message = "hello, everyone!";
    8.        rabbitTemplate.convertAndSend(exchangeName, "", message);
    9.   }

    【consumer】

    1. // fanout
    2.    @RabbitListener(queues = "fanout.queue1")
    3.    public void listenFanoutQueue1(String msg) {
    4.        System.out.println("消费者1接收到Fanout消息:【" + msg + "】");
    5.   }
    6.    @RabbitListener(queues = "fanout.queue2")
    7.    public void listenFanoutQueue2(String msg) {
    8.        System.out.println("消费者2接收到Fanout消息:【" + msg + "】");
    9.   }

    ② Direct-路由模式

    743e27a0067256a470d98acb2d1475c7.png

    在Fanout模式中,一条消息,会被所有订阅的队列都消费。但是,在某些场景下,我们希望不同的消息被不同的队列消费。这时就要用到Direct类型的Exchange。

    【publisher】

    1. // direct
    2.    @Test
    3.    public void testSendDirectExchange() {
    4.        // 交换机名称
    5.        String exchangeName = "itcast.direct";
    6.        // 消息
    7.        String message = "红色警报!日本乱排核废水,导致海洋生物变异,惊现哥斯拉!";
    8.        // 发送消息
    9.        rabbitTemplate.convertAndSend(exchangeName, "red", message);
    10.   }

    【consumer】

    1. // direct
    2.    @RabbitListener(bindings = @QueueBinding(
    3.            value = @Queue(name = "direct.queue1"),
    4.            exchange = @Exchange(name = "itcast.direct", type = ExchangeTypes.DIRECT),
    5.            key = {"red", "blue"}
    6.   ))
    7.    public void listenDirectQueue1(String msg) {
    8.        System.out.println("消费者接收到direct.queue1的消息:【" + msg + "】");
    9.   }
    10.    @RabbitListener(bindings = @QueueBinding(
    11.            value = @Queue(name = "direct.queue2"),
    12.            exchange = @Exchange(name = "itcast.direct", type = ExchangeTypes.DIRECT),
    13.            key = {"red", "yellow"}
    14.   ))
    15.    public void listenDirectQueue2(String msg) {
    16.        System.out.println("消费者接收到direct.queue2的消息:【" + msg + "】");
    17.   }

    ③ Topic-主题模式

    Topic类型的ExchangeDirect相比,都是可以根据RoutingKey把消息路由到不同的队列。只不过Topic类型Exchange可以让队列在绑定Routing key 的时候使用通配符!

    Routingkey 一般都是有一个或多个单词组成,多个单词之间以”.”分割,例如: item.insert

    通配符规则:

    #:匹配一个或多个词

    *:匹配不多不少恰好1个词

    【publisher】

    1. // topic
    2.    @Test
    3.    public void testSendTopicExchange() {
    4.        // 交换机名称
    5.        String exchangeName = "itcast.topic";
    6.        // 消息
    7.        String message = "喜报!孙悟空大战哥斯拉,胜!";
    8.        // 发送消息
    9.        rabbitTemplate.convertAndSend(exchangeName, "china.news", message);
    10.   }  

    【consumer】

    1. // topic
    2.    @RabbitListener(bindings = @QueueBinding(
    3.            value = @Queue(name = "topic.queue1"),
    4.            exchange = @Exchange(name = "itcast.topic", type = ExchangeTypes.TOPIC),
    5.            key = "china.#"
    6.   ))
    7.    public void listenTopicQueue1(String msg) {
    8.        System.out.println("消费者接收到topic.queue1的消息:【" + msg + "】");
    9.   }
    10.    @RabbitListener(bindings = @QueueBinding(
    11.            value = @Queue(name = "topic.queue2"),
    12.            exchange = @Exchange(name = "itcast.topic", type = ExchangeTypes.TOPIC),
    13.            key = "#.news"
    14.   ))
    15.    public void listenTopicQueue2(String msg) {
    16.        System.out.println("消费者接收到topic.queue2的消息:【" + msg + "】");
    17.   }

    3.6、消息转换器

    ① 消息是对象

        问题:当消息传递是对象时,会调用jdk的序列化,内存空间占用大

    【publisher】

    1. // 测试jdk序列化
    2.    @Test
    3.    public void testSendMap() throws InterruptedException {
    4.        // 准备消息
    5.        Map<String,Object> msg = new HashMap<>();
    6.        msg.put("name", "Jack");
    7.        msg.put("age", 21);
    8.        // 发送消息
    9.        // messageConverter.toMessage(msg, msg);
    10.        rabbitTemplate.convertAndSend("simple.queue","", msg);
    11.   }

    【consumer】

    1. // 测试jdk序列化
    2.    @RabbitListener(queues = "simple.queue")
    3.    public void listenSimpleQueueMessage(Map msg) throws InterruptedException {
    4.        System.out.println("spring 消费者接收到消息:【" + msg + "】");
    5.   }

    ② 解决

    1. <!--json转换-->
    2.    <dependency>
    3.   <groupId>com.fasterxml.jackson.dataformat</groupId>
    4.   <artifactId>jackson-dataformat-xml</artifactId>
    5. <version>2.9.10</version>
    6.    </dependency>  
    1. @Bean
    2.    public MessageConverter jsonMessageConverter(){
    3.        return new Jackson2JsonMessageConverter();
    4.   }

  • 相关阅读:
    并查集路径压缩
    Mac安装与配置eclipse
    儿童台灯哪个品牌比较好?精选央视消费主张推荐的护眼灯
    uniapp u-tabs表单如何默认选中
    JPA Native Query(本地查询)及查询结果转换
    K线形态识别_跛脚阳线
    SXSSFWorkbook-MinIo-大数据-流式导出
    企业申报“专精特新”,对知识产权有哪些要求?
    uni-app中使用computed解决了tab切换中data()值显示的异常
    用docker搭载redis集群
  • 原文地址:https://blog.csdn.net/hnhroot/article/details/125463948