精华内容
下载资源
问答
  • 首先介绍一下rabbitmq三种模式 Direct–路由模式 任何发送到Direct Exchange的消息都会被转发到RouteKey指定的Queue。 这种模式下不需要将Exchange进行任何绑定(binding)操作。 消息传递时需要一个“RouteKey”,...
  • 基于SpringBoot整合RabbitMQ发送邮件通知---构建springcloud微服务资源搭建。
  • SpringBoot整合Rabbitmq发送接收消息实战 另外,博主发起了SpringBoot整合Rabbitmq这一系列的gitchat交流会。刚兴趣的童鞋可以进入交流:https://gitbook.cn/gitchat/activity/5b90f9214fb1bd5c9acd4338 交流QQ:...
  • springboot整合rabbitmq,开启手工确认,保证消息100%投递。springboot整合rabbitmq,开启手工确认,保证消息100%投递。
  • springboot整合RabbitMQ实现延时队列的两种方式 教程及源码。参考博客:https://blog.csdn.net/qq_29914837/article/details/94070677
  • Springboot整合RabbitMQ最简单demo,适用于springcloud项目,作为消息总线适用,需要安装RabbitMQ,Mac linux可以使用命令行一键安装,在项目配置文件配置好端口即可(已默认配置),启动项目访问8080端口,参数见controller.
  • springboot整合RabbitMQ实现死信/死信队列及实现源码及教程,参考博客:https://blog.csdn.net/qq_29914837/article/details/93334313
  • SpringBoot整合RabbitMQ 实现消息发送确认与消息接收确认机制 源码及教材 可以参考博客: https://blog.csdn.net/qq_29914837/article/details/93376741
  • 1 SpringBoot整合RabbitMQ实战系列教程-整合配置篇-源码数据库
  • springboot-rabbitmq
  • Springboot 整合RabbitMq ,用心看完这一篇就够了

    万次阅读 多人点赞 2019-09-03 23:21:11
    该篇文章内容较多,包括有rabbitMq相关的一些简单理论介绍,provider消息推送实例,consumer消息消费实例,Direct、Topic、Fanout的使用,消息回调、手动确认等。 (但是关于rabbitMq的安装,就不介绍了) 在安装...

    该篇文章内容较多,包括有rabbitMq相关的一些简单理论介绍,provider消息推送实例,consumer消息消费实例,Direct、Topic、Fanout的使用,消息回调、手动确认等。 (但是关于rabbitMq的安装,就不介绍了)
     

    在安装完rabbitMq后,输入http://ip:15672/ ,是可以看到一个简单后台管理界面的。

    在这个界面里面我们可以做些什么?
    可以手动创建虚拟host,创建用户,分配权限,创建交换机,创建队列等等,还有查看队列消息,消费效率,推送效率等等。

    以上这些管理界面的操作在这篇暂时不做扩展描述,我想着重介绍后面实例里会使用到的。

    首先先介绍一个简单的一个消息推送到接收的流程,提供一个简单的图:
      

    JCccc-RabbitMq
    RabbitMq -JCccc

    黄色的圈圈就是我们的消息推送服务,将消息推送到 中间方框里面也就是 rabbitMq的服务器,然后经过服务器里面的交换机、队列等各种关系(后面会详细讲)将数据处理入列后,最终右边的蓝色圈圈消费者获取对应监听的消息。

    常用的交换机有以下三种,因为消费者是从队列获取信息的,队列是绑定交换机的(一般),所以对应的消息推送/接收模式也会有以下几种:

    Direct Exchange 

    直连型交换机,根据消息携带的路由键将消息投递给对应队列。

    大致流程,有一个队列绑定到一个直连交换机上,同时赋予一个路由键 routing key 。
    然后当一个消息携带着路由值为X,这个消息通过生产者发送给交换机时,交换机就会根据这个路由值X去寻找绑定值也是X的队列。

    Fanout Exchange

    扇型交换机,这个交换机没有路由键概念,就算你绑了路由键也是无视的。 这个交换机在接收到消息后,会直接转发到绑定到它上面的所有队列。

    Topic Exchange

    主题交换机,这个交换机其实跟直连交换机流程差不多,但是它的特点就是在它的路由键和绑定键之间是有规则的。
    简单地介绍下规则:

    *  (星号) 用来表示一个单词 (必须出现的)
    #  (井号) 用来表示任意数量(零个或多个)单词
    通配的绑定键是跟队列进行绑定的,举个小例子
    队列Q1 绑定键为 *.TT.*          队列Q2绑定键为  TT.#
    如果一条消息携带的路由键为 A.TT.B,那么队列Q1将会收到;
    如果一条消息携带的路由键为TT.AA.BB,那么队列Q2将会收到;

    主题交换机是非常强大的,为啥这么膨胀?
    当一个队列的绑定键为 "#"(井号) 的时候,这个队列将会无视消息的路由键,接收所有的消息。
    当 * (星号) 和 # (井号) 这两个特殊字符都未在绑定键中出现的时候,此时主题交换机就拥有的直连交换机的行为。
    所以主题交换机也就实现了扇形交换机的功能,和直连交换机的功能。

    另外还有 Header Exchange 头交换机 ,Default Exchange 默认交换机,Dead Letter Exchange 死信交换机,这几个该篇暂不做讲述。

    好了,一些简单的介绍到这里为止,  接下来我们来一起编码。

    本次实例教程需要创建2个springboot项目,一个 rabbitmq-provider (生产者),一个rabbitmq-consumer(消费者)。

    首先创建 rabbitmq-provider,

    pom.xml里用到的jar依赖:

            <!--rabbitmq-->
            <dependency>
                <groupId>org.springframework.boot</groupId>
                <artifactId>spring-boot-starter-amqp</artifactId>
            </dependency>
            <dependency>
                <groupId>org.springframework.boot</groupId>
                <artifactId>spring-boot-starter-web</artifactId>
            </dependency>

    然后application.yml:

    ps:里面的虚拟host配置项不是必须的,我自己在rabbitmq服务上创建了自己的虚拟host,所以我配置了;你们不创建,就不用加这个配置项。

    server:
      port: 8021
    spring:
      #给项目来个名字
      application:
        name: rabbitmq-provider
      #配置rabbitMq 服务器
      rabbitmq:
        host: 127.0.0.1
        port: 5672
        username: root
        password: root
        #虚拟host 可以不设置,使用server默认host
        virtual-host: JCcccHost

    接着我们先使用下direct exchange(直连型交换机),创建DirectRabbitConfig.java(对于队列和交换机持久化以及连接使用设置,在注释里有说明,后面的不同交换机的配置就不做同样说明了):

    import org.springframework.amqp.core.Binding;
    import org.springframework.amqp.core.BindingBuilder;
    import org.springframework.amqp.core.DirectExchange;
    import org.springframework.amqp.core.Queue;
    import org.springframework.context.annotation.Bean;
    import org.springframework.context.annotation.Configuration;
    
    /**
     * @Author : JCccc
     * @CreateTime : 2019/9/3
     * @Description :
     **/
    @Configuration
    public class DirectRabbitConfig {
    
        //队列 起名:TestDirectQueue
        @Bean
        public Queue TestDirectQueue() {
            // durable:是否持久化,默认是false,持久化队列:会被存储在磁盘上,当消息代理重启时仍然存在,暂存队列:当前连接有效
            // exclusive:默认也是false,只能被当前创建的连接使用,而且当连接关闭后队列即被删除。此参考优先级高于durable
            // autoDelete:是否自动删除,当没有生产者或者消费者使用此队列,该队列会自动删除。
            //   return new Queue("TestDirectQueue",true,true,false);
    
            //一般设置一下队列的持久化就好,其余两个就是默认false
            return new Queue("TestDirectQueue",true);
        }
    
        //Direct交换机 起名:TestDirectExchange
        @Bean
        DirectExchange TestDirectExchange() {
          //  return new DirectExchange("TestDirectExchange",true,true);
            return new DirectExchange("TestDirectExchange",true,false);
        }
    
        //绑定  将队列和交换机绑定, 并设置用于匹配键:TestDirectRouting
        @Bean
        Binding bindingDirect() {
            return BindingBuilder.bind(TestDirectQueue()).to(TestDirectExchange()).with("TestDirectRouting");
        }
    
    
    
        @Bean
        DirectExchange lonelyDirectExchange() {
            return new DirectExchange("lonelyDirectExchange");
        }
    
    
    
    }

    然后写个简单的接口进行消息推送(根据需求也可以改为定时任务等等,具体看需求),SendMessageController.java:

    import org.springframework.amqp.rabbit.core.RabbitTemplate;
    import org.springframework.beans.factory.annotation.Autowired;
    import org.springframework.web.bind.annotation.GetMapping;
    import org.springframework.web.bind.annotation.RestController;
    import java.time.LocalDateTime;
    import java.time.format.DateTimeFormatter;
    import java.util.HashMap;
    import java.util.Map;
    import java.util.UUID;
    
    /**
     * @Author : JCccc
     * @CreateTime : 2019/9/3
     * @Description :
     **/
    @RestController
    public class SendMessageController {
    
        @Autowired
        RabbitTemplate rabbitTemplate;  //使用RabbitTemplate,这提供了接收/发送等等方法
    
        @GetMapping("/sendDirectMessage")
        public String sendDirectMessage() {
            String messageId = String.valueOf(UUID.randomUUID());
            String messageData = "test message, hello!";
            String createTime = LocalDateTime.now().format(DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss"));
            Map<String,Object> map=new HashMap<>();
            map.put("messageId",messageId);
            map.put("messageData",messageData);
            map.put("createTime",createTime);
            //将消息携带绑定键值:TestDirectRouting 发送到交换机TestDirectExchange
            rabbitTemplate.convertAndSend("TestDirectExchange", "TestDirectRouting", map);
            return "ok";
        }
    
    
    }

    把rabbitmq-provider项目运行,调用下接口:

    因为我们目前还没弄消费者 rabbitmq-consumer,消息没有被消费的,我们去rabbitMq管理页面看看,是否推送成功:


    再看看队列(界面上的各个英文项代表什么意思,可以自己查查哈,对理解还是有帮助的):

    很好,消息已经推送到rabbitMq服务器上面了。

     


    接下来,创建rabbitmq-consumer项目:

    pom.xml里的jar依赖:

            <!--rabbitmq-->
            <dependency>
                <groupId>org.springframework.boot</groupId>
                <artifactId>spring-boot-starter-amqp</artifactId>
            </dependency>
            <dependency>
                <groupId>org.springframework.boot</groupId>
                <artifactId>spring-boot-starter</artifactId>
            </dependency>

    然后是 application.yml:

    
    server:
      port: 8022
    spring:
      #给项目来个名字
      application:
        name: rabbitmq-consumer
      #配置rabbitMq 服务器
      rabbitmq:
        host: 127.0.0.1
        port: 5672
        username: root
        password: root
        #虚拟host 可以不设置,使用server默认host
        virtual-host: JCcccHost

    然后一样,创建DirectRabbitConfig.java(消费者单纯的使用,其实可以不用添加这个配置,直接建后面的监听就好,使用注解来让监听器监听对应的队列即可。配置上了的话,其实消费者也是生成者的身份,也能推送该消息。):

    import org.springframework.amqp.core.Binding;
    import org.springframework.amqp.core.BindingBuilder;
    import org.springframework.amqp.core.DirectExchange;
    import org.springframework.amqp.core.Queue;
    import org.springframework.context.annotation.Bean;
    import org.springframework.context.annotation.Configuration;
    
    /**
     * @Author : JCccc
     * @CreateTime : 2019/9/3
     * @Description :
     **/
    @Configuration
    public class DirectRabbitConfig {
    
        //队列 起名:TestDirectQueue
        @Bean
        public Queue TestDirectQueue() {
            return new Queue("TestDirectQueue",true);
        }
    
        //Direct交换机 起名:TestDirectExchange
        @Bean
        DirectExchange TestDirectExchange() {
            return new DirectExchange("TestDirectExchange");
        }
    
        //绑定  将队列和交换机绑定, 并设置用于匹配键:TestDirectRouting
        @Bean
        Binding bindingDirect() {
            return BindingBuilder.bind(TestDirectQueue()).to(TestDirectExchange()).with("TestDirectRouting");
        }
    }
    

    然后是创建消息接收监听类,DirectReceiver.java:

    @Component
    @RabbitListener(queues = "TestDirectQueue")//监听的队列名称 TestDirectQueue
    public class DirectReceiver {
    
        @RabbitHandler
        public void process(Map testMessage) {
            System.out.println("DirectReceiver消费者收到消息  : " + testMessage.toString());
        }
    
    }

    然后将rabbitmq-consumer项目运行起来,可以看到把之前推送的那条消息消费下来了:

    然后可以再继续调用rabbitmq-provider项目的推送消息接口,可以看到消费者即时消费消息:

     

    那么直连交换机既然是一对一,那如果咱们配置多台监听绑定到同一个直连交互的同一个队列,会怎么样?

    可以看到是实现了轮询的方式对消息进行消费,而且不存在重复消费。

     

    接着,我们使用Topic Exchange 主题交换机。

    在rabbitmq-provider项目里面创建TopicRabbitConfig.java:


    import org.springframework.amqp.core.Binding;
    import org.springframework.amqp.core.BindingBuilder;
    import org.springframework.amqp.core.Queue;
    import org.springframework.amqp.core.TopicExchange;
    import org.springframework.context.annotation.Bean;
    import org.springframework.context.annotation.Configuration;
    
    /**
     * @Author : JCccc
     * @CreateTime : 2019/9/3
     * @Description :
     **/
    
    @Configuration
    public class TopicRabbitConfig {
        //绑定键
        public final static String man = "topic.man";
        public final static String woman = "topic.woman";
    
        @Bean
        public Queue firstQueue() {
            return new Queue(TopicRabbitConfig.man);
        }
    
        @Bean
        public Queue secondQueue() {
            return new Queue(TopicRabbitConfig.woman);
        }
    
        @Bean
        TopicExchange exchange() {
            return new TopicExchange("topicExchange");
        }
    
    
        //将firstQueue和topicExchange绑定,而且绑定的键值为topic.man
        //这样只要是消息携带的路由键是topic.man,才会分发到该队列
        @Bean
        Binding bindingExchangeMessage() {
            return BindingBuilder.bind(firstQueue()).to(exchange()).with(man);
        }
    
        //将secondQueue和topicExchange绑定,而且绑定的键值为用上通配路由键规则topic.#
        // 这样只要是消息携带的路由键是以topic.开头,都会分发到该队列
        @Bean
        Binding bindingExchangeMessage2() {
            return BindingBuilder.bind(secondQueue()).to(exchange()).with("topic.#");
        }
    
    }

    然后添加多2个接口,用于推送消息到主题交换机:

        @GetMapping("/sendTopicMessage1")
        public String sendTopicMessage1() {
            String messageId = String.valueOf(UUID.randomUUID());
            String messageData = "message: M A N ";
            String createTime = LocalDateTime.now().format(DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss"));
            Map<String, Object> manMap = new HashMap<>();
            manMap.put("messageId", messageId);
            manMap.put("messageData", messageData);
            manMap.put("createTime", createTime);
            rabbitTemplate.convertAndSend("topicExchange", "topic.man", manMap);
            return "ok";
        }
    
        @GetMapping("/sendTopicMessage2")
        public String sendTopicMessage2() {
            String messageId = String.valueOf(UUID.randomUUID());
            String messageData = "message: woman is all ";
            String createTime = LocalDateTime.now().format(DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss"));
            Map<String, Object> womanMap = new HashMap<>();
            womanMap.put("messageId", messageId);
            womanMap.put("messageData", messageData);
            womanMap.put("createTime", createTime);
            rabbitTemplate.convertAndSend("topicExchange", "topic.woman", womanMap);
            return "ok";
        }
    }
    

    生产者这边已经完事,先不急着运行,在rabbitmq-consumer项目上,创建TopicManReceiver.java:

    import org.springframework.amqp.rabbit.annotation.RabbitHandler;
    import org.springframework.amqp.rabbit.annotation.RabbitListener;
    import org.springframework.stereotype.Component;
    import java.util.Map;
    
    /**
     * @Author : JCccc
     * @CreateTime : 2019/9/3
     * @Description :
     **/
    @Component
    @RabbitListener(queues = "topic.man")
    public class TopicManReceiver {
    
        @RabbitHandler
        public void process(Map testMessage) {
            System.out.println("TopicManReceiver消费者收到消息  : " + testMessage.toString());
        }
    }

    再创建一个TopicTotalReceiver.java:

    package com.elegant.rabbitmqconsumer.receiver;
    
    import org.springframework.amqp.rabbit.annotation.RabbitHandler;
    import org.springframework.amqp.rabbit.annotation.RabbitListener;
    import org.springframework.stereotype.Component;
    import java.util.Map;
    
    /**
     * @Author : JCccc
     * @CreateTime : 2019/9/3
     * @Description :
     **/
    
    @Component
    @RabbitListener(queues = "topic.woman")
    public class TopicTotalReceiver {
    
        @RabbitHandler
        public void process(Map testMessage) {
            System.out.println("TopicTotalReceiver消费者收到消息  : " + testMessage.toString());
        }
    }

    同样,加主题交换机的相关配置,TopicRabbitConfig.java(消费者一定要加这个配置吗? 不需要的其实,理由在前面已经说过了。):

    import org.springframework.amqp.core.Binding;
    import org.springframework.amqp.core.BindingBuilder;
    import org.springframework.amqp.core.Queue;
    import org.springframework.amqp.core.TopicExchange;
    import org.springframework.context.annotation.Bean;
    import org.springframework.context.annotation.Configuration;
    
    /**
     * @Author : JCccc
     * @CreateTime : 2019/9/3
     * @Description :
     **/
    
    @Configuration
    public class TopicRabbitConfig {
        //绑定键
        public final static String man = "topic.man";
        public final static String woman = "topic.woman";
    
        @Bean
        public Queue firstQueue() {
            return new Queue(TopicRabbitConfig.man);
        }
    
        @Bean
        public Queue secondQueue() {
            return new Queue(TopicRabbitConfig.woman);
        }
    
        @Bean
        TopicExchange exchange() {
            return new TopicExchange("topicExchange");
        }
    
    
        //将firstQueue和topicExchange绑定,而且绑定的键值为topic.man
        //这样只要是消息携带的路由键是topic.man,才会分发到该队列
        @Bean
        Binding bindingExchangeMessage() {
            return BindingBuilder.bind(firstQueue()).to(exchange()).with(man);
        }
    
        //将secondQueue和topicExchange绑定,而且绑定的键值为用上通配路由键规则topic.#
        // 这样只要是消息携带的路由键是以topic.开头,都会分发到该队列
        @Bean
        Binding bindingExchangeMessage2() {
            return BindingBuilder.bind(secondQueue()).to(exchange()).with("topic.#");
        }
    
    }


    然后把rabbitmq-provider,rabbitmq-consumer两个项目都跑起来,先调用/sendTopicMessage1  接口:

    然后看消费者rabbitmq-consumer的控制台输出情况:
    TopicManReceiver监听队列1,绑定键为:topic.man
    TopicTotalReceiver监听队列2,绑定键为:topic.#
    而当前推送的消息,携带的路由键为:topic.man  

    所以可以看到两个监听消费者receiver都成功消费到了消息,因为这两个recevier监听的队列的绑定键都能与这条消息携带的路由键匹配上。

    接下来调用接口/sendTopicMessage2:

    然后看消费者rabbitmq-consumer的控制台输出情况:
    TopicManReceiver监听队列1,绑定键为:topic.man
    TopicTotalReceiver监听队列2,绑定键为:topic.#
    而当前推送的消息,携带的路由键为:topic.woman

    所以可以看到两个监听消费者只有TopicTotalReceiver成功消费到了消息。

     

    接下来是使用Fanout Exchang 扇型交换机。

    同样地,先在rabbitmq-provider项目上创建FanoutRabbitConfig.java:

    import org.springframework.amqp.core.Binding;
    import org.springframework.amqp.core.BindingBuilder;
    import org.springframework.amqp.core.FanoutExchange;
    import org.springframework.amqp.core.Queue;
    import org.springframework.context.annotation.Bean;
    import org.springframework.context.annotation.Configuration;
    /**
     * @Author : JCccc
     * @CreateTime : 2019/9/3
     * @Description :
     **/
    
    @Configuration
    public class FanoutRabbitConfig {
    
        /**
         *  创建三个队列 :fanout.A   fanout.B  fanout.C
         *  将三个队列都绑定在交换机 fanoutExchange 上
         *  因为是扇型交换机, 路由键无需配置,配置也不起作用
         */
    
    
        @Bean
        public Queue queueA() {
            return new Queue("fanout.A");
        }
    
        @Bean
        public Queue queueB() {
            return new Queue("fanout.B");
        }
    
        @Bean
        public Queue queueC() {
            return new Queue("fanout.C");
        }
    
        @Bean
        FanoutExchange fanoutExchange() {
            return new FanoutExchange("fanoutExchange");
        }
    
        @Bean
        Binding bindingExchangeA() {
            return BindingBuilder.bind(queueA()).to(fanoutExchange());
        }
    
        @Bean
        Binding bindingExchangeB() {
            return BindingBuilder.bind(queueB()).to(fanoutExchange());
        }
    
        @Bean
        Binding bindingExchangeC() {
            return BindingBuilder.bind(queueC()).to(fanoutExchange());
        }
    }
    

    然后是写一个接口用于推送消息,
     

        @GetMapping("/sendFanoutMessage")
        public String sendFanoutMessage() {
            String messageId = String.valueOf(UUID.randomUUID());
            String messageData = "message: testFanoutMessage ";
            String createTime = LocalDateTime.now().format(DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss"));
            Map<String, Object> map = new HashMap<>();
            map.put("messageId", messageId);
            map.put("messageData", messageData);
            map.put("createTime", createTime);
            rabbitTemplate.convertAndSend("fanoutExchange", null, map);
            return "ok";
        }

    接着在rabbitmq-consumer项目里加上消息消费类,

    FanoutReceiverA.java:

    import org.springframework.amqp.rabbit.annotation.RabbitHandler;
    import org.springframework.amqp.rabbit.annotation.RabbitListener;
    import org.springframework.stereotype.Component;
    import java.util.Map;
    /**
     * @Author : JCccc
     * @CreateTime : 2019/9/3
     * @Description :
     **/
    @Component
    @RabbitListener(queues = "fanout.A")
    public class FanoutReceiverA {
    
        @RabbitHandler
        public void process(Map testMessage) {
            System.out.println("FanoutReceiverA消费者收到消息  : " +testMessage.toString());
        }
    
    }
    

    FanoutReceiverB.java:

    import org.springframework.amqp.rabbit.annotation.RabbitHandler;
    import org.springframework.amqp.rabbit.annotation.RabbitListener;
    import org.springframework.stereotype.Component;
    import java.util.Map;
    /**
     * @Author : JCccc
     * @CreateTime : 2019/9/3
     * @Description :
     **/
    @Component
    @RabbitListener(queues = "fanout.B")
    public class FanoutReceiverB {
    
        @RabbitHandler
        public void process(Map testMessage) {
            System.out.println("FanoutReceiverB消费者收到消息  : " +testMessage.toString());
        }
    
    }

    FanoutReceiverC.java:

    import org.springframework.amqp.rabbit.annotation.RabbitHandler;
    import org.springframework.amqp.rabbit.annotation.RabbitListener;
    import org.springframework.stereotype.Component;
    import java.util.Map;
    
    /**
     * @Author : JCccc
     * @CreateTime : 2019/9/3
     * @Description :
     **/
    @Component
    @RabbitListener(queues = "fanout.C")
    public class FanoutReceiverC {
    
        @RabbitHandler
        public void process(Map testMessage) {
            System.out.println("FanoutReceiverC消费者收到消息  : " +testMessage.toString());
        }
    
    }
    

    然后加上扇型交换机的配置类,FanoutRabbitConfig.java(消费者真的要加这个配置吗? 不需要的其实,理由在前面已经说过了):

    import org.springframework.amqp.core.Binding;
    import org.springframework.amqp.core.BindingBuilder;
    import org.springframework.amqp.core.FanoutExchange;
    import org.springframework.amqp.core.Queue;
    import org.springframework.context.annotation.Bean;
    import org.springframework.context.annotation.Configuration;
    /**
     * @Author : JCccc
     * @CreateTime : 2019/9/3
     * @Description :
     **/
    @Configuration
    public class FanoutRabbitConfig {
    
        /**
         *  创建三个队列 :fanout.A   fanout.B  fanout.C
         *  将三个队列都绑定在交换机 fanoutExchange 上
         *  因为是扇型交换机, 路由键无需配置,配置也不起作用
         */
    
    
        @Bean
        public Queue queueA() {
            return new Queue("fanout.A");
        }
    
        @Bean
        public Queue queueB() {
            return new Queue("fanout.B");
        }
    
        @Bean
        public Queue queueC() {
            return new Queue("fanout.C");
        }
    
        @Bean
        FanoutExchange fanoutExchange() {
            return new FanoutExchange("fanoutExchange");
        }
    
        @Bean
        Binding bindingExchangeA() {
            return BindingBuilder.bind(queueA()).to(fanoutExchange());
        }
    
        @Bean
        Binding bindingExchangeB() {
            return BindingBuilder.bind(queueB()).to(fanoutExchange());
        }
    
        @Bean
        Binding bindingExchangeC() {
            return BindingBuilder.bind(queueC()).to(fanoutExchange());
        }
    }
    

    最后将rabbitmq-provider和rabbitmq-consumer项目都跑起来,调用下接口/sendFanoutMessage :

    然后看看rabbitmq-consumer项目的控制台情况:

    可以看到只要发送到 fanoutExchange 这个扇型交换机的消息, 三个队列都绑定这个交换机,所以三个消息接收类都监听到了这条消息。


    到了这里其实三个常用的交换机的使用我们已经完毕了,那么接下来我们继续讲讲消息的回调,其实就是消息确认(生产者推送消息成功,消费者接收消息成功)。
     

    在rabbitmq-provider项目的application.yml文件上,加上消息确认的配置项后:
     

    ps: 本篇文章使用springboot版本为 2.1.7.RELEASE ; 
    如果你们在配置确认回调,测试发现无法触发回调函数,那么存在原因也许是因为版本导致的配置项不起效,
    可以把
    publisher-confirms: true 替换为  publisher-confirm-type: correlated

    server:
      port: 8021
    spring:
      #给项目来个名字
      application:
        name: rabbitmq-provider
      #配置rabbitMq 服务器
      rabbitmq:
        host: 127.0.0.1
        port: 5672
        username: root
        password: root
        #虚拟host 可以不设置,使用server默认host
        virtual-host: JCcccHost
        #消息确认配置项
    
        #确认消息已发送到交换机(Exchange)
        publisher-confirms: true
        #确认消息已发送到队列(Queue)
        publisher-returns: true

    然后是配置相关的消息确认回调函数,RabbitConfig.java:

    import org.springframework.amqp.core.Message;
    import org.springframework.amqp.rabbit.connection.ConnectionFactory;
    import org.springframework.amqp.rabbit.connection.CorrelationData;
    import org.springframework.amqp.rabbit.core.RabbitTemplate;
    import org.springframework.context.annotation.Bean;
    import org.springframework.context.annotation.Configuration;
    
    
    /**
     * @Author : JCccc
     * @CreateTime : 2019/9/3
     * @Description :
     **/
    @Configuration
    public class RabbitConfig {
    
        @Bean
        public RabbitTemplate createRabbitTemplate(ConnectionFactory connectionFactory){
            RabbitTemplate rabbitTemplate = new RabbitTemplate();
            rabbitTemplate.setConnectionFactory(connectionFactory);
            //设置开启Mandatory,才能触发回调函数,无论消息推送结果怎么样都强制调用回调函数
            rabbitTemplate.setMandatory(true);
    
            rabbitTemplate.setConfirmCallback(new RabbitTemplate.ConfirmCallback() {
                @Override
                public void confirm(CorrelationData correlationData, boolean ack, String cause) {
                    System.out.println("ConfirmCallback:     "+"相关数据:"+correlationData);
                    System.out.println("ConfirmCallback:     "+"确认情况:"+ack);
                    System.out.println("ConfirmCallback:     "+"原因:"+cause);
                }
            });
    
            rabbitTemplate.setReturnCallback(new RabbitTemplate.ReturnCallback() {
                @Override
                public void returnedMessage(Message message, int replyCode, String replyText, String exchange, String routingKey) {
                    System.out.println("ReturnCallback:     "+"消息:"+message);
                    System.out.println("ReturnCallback:     "+"回应码:"+replyCode);
                    System.out.println("ReturnCallback:     "+"回应信息:"+replyText);
                    System.out.println("ReturnCallback:     "+"交换机:"+exchange);
                    System.out.println("ReturnCallback:     "+"路由键:"+routingKey);
                }
            });
    
            return rabbitTemplate;
        }
    
    }

    到这里,生产者推送消息的消息确认调用回调函数已经完毕。
    可以看到上面写了两个回调函数,一个叫 ConfirmCallback ,一个叫 RetrunCallback;
    那么以上这两种回调函数都是在什么情况会触发呢?

    先从总体的情况分析,推送消息存在四种情况:

    ①消息推送到server,但是在server里找不到交换机
    ②消息推送到server,找到交换机了,但是没找到队列
    ③消息推送到sever,交换机和队列啥都没找到
    ④消息推送成功

    那么我先写几个接口来分别测试和认证下以上4种情况,消息确认触发回调函数的情况:

    ①消息推送到server,但是在server里找不到交换机
    写个测试接口,把消息推送到名为‘non-existent-exchange’的交换机上(这个交换机是没有创建没有配置的):

        @GetMapping("/TestMessageAck")
        public String TestMessageAck() {
            String messageId = String.valueOf(UUID.randomUUID());
            String messageData = "message: non-existent-exchange test message ";
            String createTime = LocalDateTime.now().format(DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss"));
            Map<String, Object> map = new HashMap<>();
            map.put("messageId", messageId);
            map.put("messageData", messageData);
            map.put("createTime", createTime);
            rabbitTemplate.convertAndSend("non-existent-exchange", "TestDirectRouting", map);
            return "ok";
        }
    

    调用接口,查看rabbitmq-provuder项目的控制台输出情况(原因里面有说,没有找到交换机'non-existent-exchange'):

    2019-09-04 09:37:45.197 ERROR 8172 --- [ 127.0.0.1:5672] o.s.a.r.c.CachingConnectionFactory       : Channel shutdown: channel error; protocol method: #method<channel.close>(reply-code=404, reply-text=NOT_FOUND - no exchange 'non-existent-exchange' in vhost 'JCcccHost', class-id=60, method-id=40)
    ConfirmCallback:     相关数据:null
    ConfirmCallback:     确认情况:false
    ConfirmCallback:     原因:channel error; protocol method: #method<channel.close>(reply-code=404, reply-text=NOT_FOUND - no exchange 'non-existent-exchange' in vhost 'JCcccHost', class-id=60, method-id=40)
    

        结论: ①这种情况触发的是 ConfirmCallback 回调函数。

     ②消息推送到server,找到交换机了,但是没找到队列  
    这种情况就是需要新增一个交换机,但是不给这个交换机绑定队列,我来简单地在DirectRabitConfig里面新增一个直连交换机,名叫‘lonelyDirectExchange’,但没给它做任何绑定配置操作:

        @Bean
        DirectExchange lonelyDirectExchange() {
            return new DirectExchange("lonelyDirectExchange");
        }

    然后写个测试接口,把消息推送到名为‘lonelyDirectExchange’的交换机上(这个交换机是没有任何队列配置的):

        @GetMapping("/TestMessageAck2")
        public String TestMessageAck2() {
            String messageId = String.valueOf(UUID.randomUUID());
            String messageData = "message: lonelyDirectExchange test message ";
            String createTime = LocalDateTime.now().format(DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss"));
            Map<String, Object> map = new HashMap<>();
            map.put("messageId", messageId);
            map.put("messageData", messageData);
            map.put("createTime", createTime);
            rabbitTemplate.convertAndSend("lonelyDirectExchange", "TestDirectRouting", map);
            return "ok";
        }

    调用接口,查看rabbitmq-provuder项目的控制台输出情况:

    ReturnCallback:     消息:(Body:'{createTime=2019-09-04 09:48:01, messageId=563077d9-0a77-4c27-8794-ecfb183eac80, messageData=message: lonelyDirectExchange test message }' MessageProperties [headers={}, contentType=application/x-java-serialized-object, contentLength=0, receivedDeliveryMode=PERSISTENT, priority=0, deliveryTag=0])
    ReturnCallback:     回应码:312
    ReturnCallback:     回应信息:NO_ROUTE
    ReturnCallback:     交换机:lonelyDirectExchange
    ReturnCallback:     路由键:TestDirectRouting
    ConfirmCallback:     相关数据:null
    ConfirmCallback:     确认情况:true
    ConfirmCallback:     原因:null

    可以看到这种情况,两个函数都被调用了;
    这种情况下,消息是推送成功到服务器了的,所以ConfirmCallback对消息确认情况是true;
    而在RetrunCallback回调函数的打印参数里面可以看到,消息是推送到了交换机成功了,但是在路由分发给队列的时候,找不到队列,所以报了错误 NO_ROUTE 。
      结论:②这种情况触发的是 ConfirmCallback和RetrunCallback两个回调函数。

    ③消息推送到sever,交换机和队列啥都没找到 
    这种情况其实一看就觉得跟①很像,没错 ,③和①情况回调是一致的,所以不做结果说明了。
      结论: ③这种情况触发的是 ConfirmCallback 回调函数。

     ④消息推送成功
    那么测试下,按照正常调用之前消息推送的接口就行,就调用下 /sendFanoutMessage接口,可以看到控制台输出:

    ConfirmCallback:     相关数据:null
    ConfirmCallback:     确认情况:true
    ConfirmCallback:     原因:null

    结论: ④这种情况触发的是 ConfirmCallback 回调函数。


    以上是生产者推送消息的消息确认 回调函数的使用介绍(可以在回调函数根据需求做对应的扩展或者业务数据处理)。

    接下来我们继续, 消费者接收到消息的消息确认机制。


    和生产者的消息确认机制不同,因为消息接收本来就是在监听消息,符合条件的消息就会消费下来。
    所以,消息接收的确认机制主要存在三种模式:

    自动确认, 这也是默认的消息确认情况。  AcknowledgeMode.NONE
    RabbitMQ成功将消息发出(即将消息成功写入TCP Socket)中立即认为本次投递已经被正确处理,不管消费者端是否成功处理本次投递。
    所以这种情况如果消费端消费逻辑抛出异常,也就是消费端没有处理成功这条消息,那么就相当于丢失了消息。
    一般这种情况我们都是使用try catch捕捉异常后,打印日志用于追踪数据,这样找出对应数据再做后续处理。

    ② 根据情况确认, 这个不做介绍
    手动确认 , 这个比较关键,也是我们配置接收消息确认机制时,多数选择的模式。
    消费者收到消息后,手动调用basic.ack/basic.nack/basic.reject后,RabbitMQ收到这些消息后,才认为本次投递成功。
    basic.ack用于肯定确认 
    basic.nack用于否定确认(注意:这是AMQP 0-9-1的RabbitMQ扩展) 
    basic.reject用于否定确认,但与basic.nack相比有一个限制:一次只能拒绝单条消息 

    消费者端以上的3个方法都表示消息已经被正确投递,但是basic.ack表示消息已经被正确处理。
    而basic.nack,basic.reject表示没有被正确处理:

    着重讲下reject,因为有时候一些场景是需要重新入列的。

    channel.basicReject(deliveryTag, true);  拒绝消费当前消息,如果第二参数传入true,就是将数据重新丢回队列里,那么下次还会消费这消息。设置false,就是告诉服务器,我已经知道这条消息数据了,因为一些原因拒绝它,而且服务器也把这个消息丢掉就行。 下次不想再消费这条消息了。

    使用拒绝后重新入列这个确认模式要谨慎,因为一般都是出现异常的时候,catch异常再拒绝入列,选择是否重入列。

    但是如果使用不当会导致一些每次都被你重入列的消息一直消费-入列-消费-入列这样循环,会导致消息积压。

     

    顺便也简单讲讲 nack,这个也是相当于设置不消费某条消息。

    channel.basicNack(deliveryTag, false, true);
    第一个参数依然是当前消息到的数据的唯一id;
    第二个参数是指是否针对多条消息;如果是true,也就是说一次性针对当前通道的消息的tagID小于当前这条消息的,都拒绝确认。
    第三个参数是指是否重新入列,也就是指不确认的消息是否重新丢回到队列里面去。

    同样使用不确认后重新入列这个确认模式要谨慎,因为这里也可能因为考虑不周出现消息一直被重新丢回去的情况,导致积压。

     


    看了上面这么多介绍,接下来我们一起配置下,看看一般的消息接收 手动确认是怎么样的。
    ​​​​​​

    在消费者项目里,
    新建MessageListenerConfig.java上添加代码相关的配置代码:

    
    import com.elegant.rabbitmqconsumer.receiver.MyAckReceiver;
    import org.springframework.amqp.core.AcknowledgeMode;
    import org.springframework.amqp.core.Queue;
    import org.springframework.amqp.rabbit.connection.CachingConnectionFactory;
    import org.springframework.amqp.rabbit.listener.SimpleMessageListenerContainer;
    import org.springframework.beans.factory.annotation.Autowired;
    import org.springframework.context.annotation.Bean;
    import org.springframework.context.annotation.Configuration;
    
    /**
     * @Author : JCccc
     * @CreateTime : 2019/9/4
     * @Description :
     **/
    @Configuration
    public class MessageListenerConfig {
    
        @Autowired
        private CachingConnectionFactory connectionFactory;
        @Autowired
        private MyAckReceiver myAckReceiver;//消息接收处理类
    
        @Bean
        public SimpleMessageListenerContainer simpleMessageListenerContainer() {
            SimpleMessageListenerContainer container = new SimpleMessageListenerContainer(connectionFactory);
            container.setConcurrentConsumers(1);
            container.setMaxConcurrentConsumers(1);
            container.setAcknowledgeMode(AcknowledgeMode.MANUAL); // RabbitMQ默认是自动确认,这里改为手动确认消息
            //设置一个队列
            container.setQueueNames("TestDirectQueue");
            //如果同时设置多个如下: 前提是队列都是必须已经创建存在的
            //  container.setQueueNames("TestDirectQueue","TestDirectQueue2","TestDirectQueue3");
    
    
            //另一种设置队列的方法,如果使用这种情况,那么要设置多个,就使用addQueues
            //container.setQueues(new Queue("TestDirectQueue",true));
            //container.addQueues(new Queue("TestDirectQueue2",true));
            //container.addQueues(new Queue("TestDirectQueue3",true));
            container.setMessageListener(myAckReceiver);
    
            return container;
        }
    
    
    }
    

    对应的手动确认消息监听类,MyAckReceiver.java(手动确认模式需要实现 ChannelAwareMessageListener):
    //之前的相关监听器可以先注释掉,以免造成多个同类型监听器都监听同一个队列。
    //这里的获取消息转换,只作参考,如果报数组越界可以自己根据格式去调整。

    import com.rabbitmq.client.Channel;
    import org.springframework.amqp.core.Message;
    import org.springframework.amqp.rabbit.listener.api.ChannelAwareMessageListener;
    import org.springframework.stereotype.Component;
    import java.util.HashMap;
    import java.util.Map;
    
    @Component
    
    public class MyAckReceiver implements ChannelAwareMessageListener {
    
        @Override
        public void onMessage(Message message, Channel channel) throws Exception {
            long deliveryTag = message.getMessageProperties().getDeliveryTag();
            try {
                //因为传递消息的时候用的map传递,所以将Map从Message内取出需要做些处理
                String msg = message.toString();
                String[] msgArray = msg.split("'");//可以点进Message里面看源码,单引号直接的数据就是我们的map消息数据
                Map<String, String> msgMap = mapStringToMap(msgArray[1].trim(),3);
                String messageId=msgMap.get("messageId");
                String messageData=msgMap.get("messageData");
                String createTime=msgMap.get("createTime");
                System.out.println("  MyAckReceiver  messageId:"+messageId+"  messageData:"+messageData+"  createTime:"+createTime);
                System.out.println("消费的主题消息来自:"+message.getMessageProperties().getConsumerQueue());
                channel.basicAck(deliveryTag, true); //第二个参数,手动确认可以被批处理,当该参数为 true 时,则可以一次性确认 delivery_tag 小于等于传入值的所有消息
    //			channel.basicReject(deliveryTag, true);//第二个参数,true会重新放回队列,所以需要自己根据业务逻辑判断什么时候使用拒绝
            } catch (Exception e) {
                channel.basicReject(deliveryTag, false);
                e.printStackTrace();
            }
        }
    
         //{key=value,key=value,key=value} 格式转换成map
        private Map<String, String> mapStringToMap(String str,int entryNum ) {
            str = str.substring(1, str.length() - 1);
            String[] strs = str.split(",",entryNum);
            Map<String, String> map = new HashMap<String, String>();
            for (String string : strs) {
                String key = string.split("=")[0].trim();
                String value = string.split("=")[1];
                map.put(key, value);
            }
            return map;
        }
    }
    

    这时,先调用接口/sendDirectMessage, 给直连交换机TestDirectExchange 的队列TestDirectQueue 推送一条消息,可以看到监听器正常消费了下来:

     

     

    到这里,我们其实已经掌握了怎么去使用消息消费的手动确认了。

    但是这个场景往往不够! 因为很多伙伴之前给我评论反应,他们需要这个消费者项目里面,监听的好几个队列都想变成手动确认模式,而且处理的消息业务逻辑不一样。

    没有问题,接下来看代码

    场景: 除了直连交换机的队列TestDirectQueue需要变成手动确认以外,我们还需要将一个其他的队列

    或者多个队列也变成手动确认,而且不同队列实现不同的业务处理。

     

    那么我们需要做的第一步,往SimpleMessageListenerContainer里添加多个队列:

    然后我们的手动确认消息监听类,MyAckReceiver.java 就可以同时将上面设置到的队列的消息都消费下来。

    但是我们需要做不用的业务逻辑处理,那么只需要  根据消息来自的队列名进行区分处理即可,如:

    import com.rabbitmq.client.Channel;
    import org.springframework.amqp.core.Message;
    import org.springframework.amqp.rabbit.listener.api.ChannelAwareMessageListener;
    import org.springframework.stereotype.Component;
    import java.util.HashMap;
    import java.util.Map;
    
    @Component
    public class MyAckReceiver implements ChannelAwareMessageListener {
    
        @Override
        public void onMessage(Message message, Channel channel) throws Exception {
            long deliveryTag = message.getMessageProperties().getDeliveryTag();
            try {
                //因为传递消息的时候用的map传递,所以将Map从Message内取出需要做些处理
                String msg = message.toString();
                String[] msgArray = msg.split("'");//可以点进Message里面看源码,单引号直接的数据就是我们的map消息数据
                Map<String, String> msgMap = mapStringToMap(msgArray[1].trim(),3);
                String messageId=msgMap.get("messageId");
                String messageData=msgMap.get("messageData");
                String createTime=msgMap.get("createTime");
                
                if ("TestDirectQueue".equals(message.getMessageProperties().getConsumerQueue())){
                    System.out.println("消费的消息来自的队列名为:"+message.getMessageProperties().getConsumerQueue());
                    System.out.println("消息成功消费到  messageId:"+messageId+"  messageData:"+messageData+"  createTime:"+createTime);
                    System.out.println("执行TestDirectQueue中的消息的业务处理流程......");
                    
                }
    
                if ("fanout.A".equals(message.getMessageProperties().getConsumerQueue())){
                    System.out.println("消费的消息来自的队列名为:"+message.getMessageProperties().getConsumerQueue());
                    System.out.println("消息成功消费到  messageId:"+messageId+"  messageData:"+messageData+"  createTime:"+createTime);
                    System.out.println("执行fanout.A中的消息的业务处理流程......");
    
                }
                
                channel.basicAck(deliveryTag, true);
    //			channel.basicReject(deliveryTag, true);//为true会重新放回队列
            } catch (Exception e) {
                channel.basicReject(deliveryTag, false);
                e.printStackTrace();
            }
        }
    
        //{key=value,key=value,key=value} 格式转换成map
        private Map<String, String> mapStringToMap(String str,int enNum) {
            str = str.substring(1, str.length() - 1);
            String[] strs = str.split(",",enNum);
            Map<String, String> map = new HashMap<String, String>();
            for (String string : strs) {
                String key = string.split("=")[0].trim();
                String value = string.split("=")[1];
                map.put(key, value);
            }
            return map;
        }
    }
    

    ok,这时候我们来分别往不同队列推送消息,看看效果:

    调用接口/sendDirectMessage  和 /sendFanoutMessage ,

     

    如果你还想新增其他的监听队列,也就是按照这种方式新增配置即可(或者完全可以分开多个消费者项目去监听处理)。 

     

     

    好,这篇Springboot整合rabbitMq教程就暂且到此。

     

     

     

     

    展开全文
  • SpringBoot整合RabbitMQ.zip

    2020-06-29 20:49:04
    SpringBoot整合RabbitMQ的详细过程 **1.该篇博文首先讲述了交换机和队列之间的绑定关系** ①direct、②fanout、③topic **2.然后讲消息的回调** 四种情况下,确认触发哪个回调函数: ①消息推送到server,但是在...
  • springboot整合RabbitMQ

    2019-05-04 02:05:41
    NULL 博文链接:https://zzc1684.iteye.com/blog/2433719
  • SpringBoot 整合RabbitMQ

    2021-12-09 18:15:37
    SpringBoot 整合RabbitMQ

    简介

    在Spring项目中,可以使用Spring-Rabbit去操作RabbitMQ,尤其是在spring boot项目中只需要引入对应的amqp启动器依赖即可,方便的使用发送消息,使用注解接收消息。
    一般在开发过程中:
    生产者工程:
    application.yml文件配置RabbitMQ相关信息;
    在生产者工程中编写配置类,用于创建交换机和队列,并进行绑定
    注入RabbitTemplate对象,通过RabbitTemplate对象发送消息到交换机
    消费者工程:
    application.yml文件配置RabbitMQ相关信息
    创建消息处理类,用于接收队列中的消息并进行处理

    搭建生产者工程

    创建工程

    创建生产者工程springboot_rabbitmq_producer,工程坐标如下:

    <artifactId>springboot_rabbitmq_producer</artifactId>
    	<groupId>cn.com.javakf</groupId>
    <version>1.0-SNAPSHOT</version>
    

    pom文件引入以下依赖

    <dependencies>
    	<dependency>
    		<groupId>org.springframework.boot</groupId>
    		<artifactId>spring-boot-starter-amqp</artifactId>
    	</dependency>
    </dependencies>
    

    启动类

    创建启动类cn.com.javakf.rabbitmq.Application,代码如下:

    @SpringBootApplication
    public class Application {
    	public static void main(String[] args) {
    		SpringApplication.run(Application.class, args);
    	}
    }
    

    配置RabbitMQ

    (1)application.yml配置文件
    创建application.yml,内容如下:

    spring:
      rabbitmq:
        host: 192.168.80.131
        port: 5672
        virtual-host: javakf
        username: admin
        password: admin
    

    (2)绑定交换机和队列
    创建RabbitMQ队列与交换机绑定的配置类

    @Configuration
    public class RabbitMQConfig {
    
    	/***
    	 * 声明交换机
    	 */
    	@Bean(name = "itemTopicExchange")
    	public Exchange topicExchange() {
    		return ExchangeBuilder.topicExchange("item_topic_exchange").durable(true).build();
    	}
    
    	/***
    	 * 声明队列
    	 */
    	@Bean(name = "itemQueue")
    	public Queue itemQueue() {
    		return QueueBuilder.durable("item_queue").build();
    	}
    
    	/***
    	 * 队列绑定到交换机上
    	 */
    	@Bean
    	public Binding itemQueueExchange(@Qualifier("itemQueue") Queue queue,
    			@Qualifier("itemTopicExchange") Exchange exchange) {
    		return BindingBuilder.bind(queue).to(exchange).with("item.#").noargs();
    	}
    }
    

    搭建消费者工程

    创建工程

    创建消费者工程springboot_rabbitmq_consumer,工程坐标如下:

    <artifactId>springboot_rabbitmq_consumer</artifactId>
    	<groupId>cn.com.javakf</groupId>
    <version>1.0-SNAPSHOT</version>
    

    pom文件引入以下依赖

    <dependencies>
    	<dependency>
    		<groupId>org.springframework.boot</groupId>
    		<artifactId>spring-boot-starter-amqp</artifactId>
    	</dependency>
    </dependencies>
    

    启动类

    @SpringBootApplication
    public class Application {
    	public static void main(String[] args) {
    		SpringApplication.run(Application.class, args);
    	}
    }
    

    配置RabbitMQ

    创建application.yml,内容如下:

    spring:
      rabbitmq:
        host: 192.168.80.131
        port: 5672
        virtual-host: javakf
        username: admin
        password: admin
    

    消息监听处理类

    @Component
    public class MessageListener {
    
    	/**
    	 * 监听某个队列的消息
    	 * 
    	 * @param message 接收到的消息
    	 */
    	@RabbitListener(queues = "item_queue")
    	public void myListener1(String message) {
    		System.out.println("消费者接收到的消息为:" + message);
    	}
    }
    

    测试

    在生产者工程springboot-rabbitmq-producer中创建测试类

    @RunWith(SpringRunner.class)
    @SpringBootTest
    public class RabbitMQTest {
    
    	// 用于发送MQ消息
    	@Autowired
    	private RabbitTemplate rabbitTemplate;
    
    	/***
    	 * 消息生产测试
    	 */
    	@Test
    	public void testCreateMessage() {
    		rabbitTemplate.convertAndSend("item_topic_exchange", "item.insert", "商品新增,routing key 为item.insert");
    		rabbitTemplate.convertAndSend("item_topic_exchange", "item.update", "商品修改,routing key 为item.update");
    		rabbitTemplate.convertAndSend("item_topic_exchange", "item.delete", "商品删除,routing key 为item.delete");
    	}
    }
    
    

    先运行上述测试程序(交换机和队列才能先被声明和绑定),然后启动消费者;在消费者工程springboot_rabbitmq_consumer中控制台查看是否接收到对应消息。

    另外;也可以在RabbitMQ的管理控制台中查看到交换机与队列的绑定:

    原文链接:https://blog.csdn.net/weixin_45730091/article/details/102779945

    展开全文
  • SpringBoot整合RabbitMQ

    2021-01-20 16:14:39
    一.RabbitMQ介绍 RabbitMQ是实现了高级消息队列协议(AMQP)的开源消息代理软件(亦称面向消息的中间件)。RabbitMQ服务器是用Erlang语言编写的,而集群和故障转移是构建在开放电信平台框架上的。RabbitMQ是一种消息...

    一.RabbitMQ介绍

    RabbitMQ是实现了高级消息队列协议(AMQP)的开源消息代理软件(亦称面向消息的中间件)。RabbitMQ服务器是用Erlang语言编写的,而集群和故障转移是构建在开放电信平台框架上的。RabbitMQ是一种消息中间件,用于处理来自客户端的异步消息。服务端将要发送的消息放入到队列池中。接收端可以根据RabbitMQ配置的转发机制接收服务端发来的消息。RabbitMQ依据指定的转发规则进行消息的转发、缓冲和持久化操作,主要用在多服务器间或单服务器的子系统间进行通信,是分布式系统标准的配置。

    Exchange交换机
    在上一节我们看到生产者将消息投递到Queue中,实际上这在RabbitMQ中这种事情永远都不会发生。实际的情况是,生产者将消息发送到Exchange(交换机),由Exchange将消息路由到一个或多个Queue中(或者丢弃)。

    routing key
    生产者在将消息发送给Exchange的时候,一般会指定一个routing key,来指定这个消息的路由规则,而这个routing key需要与Exchange Type及binding key联合使用才能最终生效。
    在Exchange Type与binding key固定的情况下(在正常使用时一般这些内容都是固定配置好的),我们的生产者就可以在发送消息给Exchange时,通过指定routing key来决定消息流向哪里。RabbitMQ为routing key设定的长度限制为255 bytes。

    Binding绑定
    RabbitMQ中通过Binding将Exchange与Queue关联起来,这样RabbitMQ就知道如何正确地将消息路由到指定的Queue了。

    Binding key
    在绑定(Binding)Exchange与Queue的同时,一般会指定一个binding key;生产者将消息发送给Exchange时,一般会指定一个routing key;当binding key与routing key相匹配时,消息将会被路由到对应的Queue中。这个将在Exchange Types章节会列举实际的例子加以说明。
    在绑定多个Queue到同一个Exchange的时候,这些Binding允许使用相同的binding key。binding key 并不是在所有情况下都生效,它依赖于Exchange Type,比如fanout类型的Exchange就会无视binding key,而是将消息路由到所有绑定到该Exchange的Queue。

    Exchange Types
    RabbitMQ常用的Exchange Type有fanout、direct、topic、headers这四种(AMQP规范里还提到两种Exchange Type,分别为system与自定义)。

    二.SpringBoot整合RabbitMQ

    1. 引入依赖:

    		<dependency>
                <groupId>org.springframework.boot</groupId>
                <artifactId>spring-boot-starter-amqp</artifactId>
            </dependency>
    

    2. application.yml文件配置:

    spring:
      rabbitmq:
        host: 192.168.157.129
        port: 5672
        username: admin
        password: admin
        virtual-host: /admin
    

    点对点方式:
    3. 简单队列模式:

    在这里插入图片描述
    在上图中,“ P”是我们的生产者,“ C”是我们的消费者。中间的红框是一个队列,代表消息缓冲区。

    • 配置类配置队列:
    @Configuration
    public class RabbitConfig {
        //simple模式
        @Bean
        public Queue getQueue(){
            return new Queue("springboot-simple-queue",true,false,false);
        }
    
    • 生产者发送消息:
    @Component
    public class SimpleSender {
    
        @Autowired
        private RabbitTemplate rabbitTemplate;
    
        public void send(String msg){
            rabbitTemplate.convertAndSend("","springboot-simple-queue",msg);
            System.out.println("消息发送成功");
        }
    }
    
    • 消费者接收消息
    @Component
    public class SimpleRec {
        @RabbitListener(queues = "springboot-simple-queue")
        public void process(String msg){
            System.out.println("接收到的消息为:"+msg);
      }
    }
    
    • 编写测试程序
    @SpringBootTest
    class RabbitmqSpringbootApplicationTests {
    
        @Autowired
        private SimpleSender sender;
    
        @Test
        void contextLoads() {
            sender.send("整合RabbitMQ成功!");
        }
    
    }
    

    发布/订阅方式:
    4. Fanout交换机模式:
    在这里插入图片描述
    生产者只能向交换机(Exchange)发送消息。交换机是一个非常简单的东西。一边接收来自生产者的消息,另一边将消息推送到队列。此处我们采用fanout交换机方式,fanout交换机非常简单。它只是将接收到的所有消息广播给它所知道的所有队列。

    • 配置类配置队列:
     //fanout发布订阅模式,声明fanout交换机
        @Bean
        public FanoutExchange getFanoutExchange(){
            return new FanoutExchange("springboot-fanout-exchange");
        }
    
     //声明队列
        @Bean
        public Queue getQueueOne(){
            return new Queue("springboot-fanout-queue1");
        }
    
        @Bean
        public Queue getQueueTwo(){
            return new Queue("springboot-fanout-queue2");
        }
    
     //将队列绑定到fanout交换机上
        @Bean
        public Binding getBindingOne(FanoutExchange getFanoutExchange,Queue getQueueOne){
            return BindingBuilder.bind(getQueueOne).to(getFanoutExchange);
        }
    
    • 生产者发送消息:
    
    @Component
    public class PublishSender {
    
        @Autowired
        private RabbitTemplate rabbitTemplate;
    
        public void send(String msg){
            rabbitTemplate.convertAndSend("springboot-fanout-exchange","",msg);
            System.out.println("消息发送成功!!!");
        }
    }
    
    
    • 消费者接收消息:
    @Component
    public class PublishReceiver {
        @RabbitListener(queues = "springboot-fanout-queue1")
        public void process(String msg) {
            System.out.println("消息队列1接收到的消息为:" + msg);
        }
    
        @RabbitListener(queues = "springboot-fanout-queue2")
        public void process2(String msg) {
    
            System.out.println("消息队列2接收到的消息为:" + msg);
        }
    }
    

    5. 路由模式
    在这里插入图片描述
    路由模式通过关键词来匹配,确定把消息发到哪个队列,消费者只订阅所有消息中的一部分。这里我们将用直连交换机(Direct exchange)。它背后的路由算法很简单——消息传递到bindingKey与routingKey完全匹配的队列。
    (实际代码略,着重讲一下之后的主题模式)

    6. 主题模式
    在这里插入图片描述
    我们采用Topics主题模式替代之前的路由模式,topics模式具有特殊的关键词规则,发送到Topic交换机的消息,它的的routingKey,必须是由点分隔的多个单词。

    *可以通配单个单词。
    #可以通配零个或多个单词。

    • 配置类配置队列:
     //topic发布订阅模式,声明交换机
        @Bean
        public TopicExchange getTopicExchange(){
            return new TopicExchange("springboot-topic-exchange");
        }
    
     //声明队列
        @Bean
        public Queue getQueueOne(){
            return new Queue("springboot-topic-queue1");
        }
        
    //将队列1绑定到topic交换机上
        @Bean
        public Binding getBindingTopic(TopicExchange getTopicExchange,Queue getQueueOne){
            return BindingBuilder.bind(getQueueOne).to(getTopicExchange).with("NBA.*");
        }
       
    
    • 生产者发送消息:
    @Component
    public class TopicSender {
    
        @Autowired
        private RabbitTemplate rabbitTemplate;
    
        public void send(String msg){
            rabbitTemplate.convertAndSend("springboot-topic-exchange","NBA.add",msg);
            System.out.println("消息发送成功");
        }
    }
    
    • 消费者接收消息:
    @Component
    public class TopicReceiver {
        @RabbitListener(queues = "springboot-topic-queue1")
        public void process(String msg) {
            System.out.println("消息队列1接收到的消息为:" + msg);
        }
    }
    
    
    展开全文
  • Springboot整合Rabbitmq

    2021-07-24 19:14:36
    -- springboot整合rabbitmq--> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-amqp</artifactId> </dependency> 配置 ...
      <!-- springboot整合rabbitmq-->
            <dependency>
                <groupId>org.springframework.boot</groupId>
                <artifactId>spring-boot-starter-amqp</artifactId>
            </dependency>
    

    配置

    spring.application.name=springboot_rabbitmq
    spring.rabbitmq.host=localhost
    spring.rabbitmq.virtual-host=/
    spring.rabbitmq.username=guest
    spring.rabbitmq.password=guest
    spring.rabbitmq.port=5672
    
    

    配置队列 交换机类型以及绑定等

    
    package com.test.springboot_rabbitmq.config;
    
    
    import org.springframework.amqp.core.*;
    import org.springframework.context.annotation.Bean;
    import org.springframework.context.annotation.Configuration;
    
    @Configuration
    public class RabbitmqConfig {
    
    
        @Bean
        public Queue myQueue() {
            return new Queue("myqueue");
        }
    
    
        @Bean
        public Exchange myExchange() {
    
            // new Exchange()
            // return new TopicExchange("topic.biz.ex", false, false, null);
    //         return new DirectExchange("direct.biz.ex", false, false, null);
            // return new FanoutExchange("fanout.biz.ex", false, false, null);
            // return new HeadersExchange("header.biz.ex", false, false, null);
            // 交换器名称,交换器类型(),是否是持久化的,是否自动删除,交换器属性 Map集合
            // return new CustomExchange("custom.biz.ex", ExchangeTypes.DIRECT, false, false, null);
            return new DirectExchange("myex", false, false, null);
        }
    
    
        @Bean
        public Binding myBining() {
    
            // 绑定的目的地,绑定的类型:到交换器还是到队列,交换器名称,路由key, 绑定的属性
            // new Binding("", Binding.DestinationType.EXCHANGE, "", "", null);
            // 绑定的目的地,绑定的类型:到交换器还是到队列,交换器名称,路由key, 绑定的属性
            // new Binding("", Binding.DestinationType.QUEUE, "", "", null);
            // 绑定了交换器direct.biz.ex到队列myqueue,路由key是 direct.biz.ex
            return new Binding("myqueue",
                    Binding.DestinationType.QUEUE, "myex",
                    "direct.biz.ex", null);
    
        }
    
    }
    
    
    

    发送消息

    
    
    @RestController
    @RequestMapping("test")
    public class RabbitmqTestController {
    
    
        @Autowired
        private AmqpTemplate rabbitTemplate;
    
        @RequestMapping("/send/{message}")
        public String sendMessage(@PathVariable String message) throws UnsupportedEncodingException {
    
            //消息属性
            MessageProperties messageProperties = MessagePropertiesBuilder.newInstance()
                    .setContentEncoding(MessageProperties.CONTENT_TYPE_JSON) //类型
                    .setHeader("key", "123").build();
    
            //消息编码
            Message build = MessageBuilder.withBody(message.getBytes("utf-8")).andProperties(messageProperties).build();
            rabbitTemplate.convertAndSend("myex", "direct.biz.ex", build);
            return "ok";
        }
    
    }
    

    监听消息

    
    
    @Component
    public class HelloConsumer {
    //
    //    @RabbitListener(queues = "myqueue")
    //    public void service(String message) {
    //        System.out.println("消息队列推送来的消息:" + message);
    //    }
    
    
        /**
         * 监听到的消息 以及设置的消息属性
         * @param message
         * @param value
         */
        @RabbitListener(queues = "myqueue")
        public void service(@Payload String message, @Header(name = "key")String value) {
            System.out.println("消息队列推送来的消息:" + message  + " value = " + value);
        }
    }
    
    
    
    展开全文
  • springboot整合rabbitmq

    2020-09-14 14:30:15
    springboot整合rabbitmq简单示例。交换机的模式使用的是Topic,消息发送使用RabbitTemplate,消息接收使用RabbitListener
  • Springboot整合RabbitMQ最简单demo
  • 实战:springboot整合rabbitMQ

    千次阅读 2021-11-04 11:24:56
    二、springboot整合rabbitMQ 1.新建springboot项目 2.pom:主要添加以下两个依赖 <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-
  • Springboot 整合RabbitMq

    2021-11-04 14:19:56
    该篇文章内容较多,包括有rabbitMq相关的一些简单理论介绍,provider消息推送实例,consumer消息消费实例,Direct、Topic、Fanout的使用,消息回调、手动确认等。 (但是关于rabbitMq的安装,就不介绍了) 在安装完...
  • 主要为大家详细介绍了SpringBootRabbitMq实现定时任务,文中示例代码介绍的非常详细,具有一定的参考价值,感兴趣的小伙伴们可以参考一下
  • 主要介绍了springboot集成rabbitMQ之对象传输的方法,小编觉得挺不错的,现在分享给大家,也给大家做个参考。一起跟随小编过来看看吧
  • Springboot+rabbitmq集成

    2018-10-25 10:15:58
    文件内包含了rabbit安装的必需文件以及springboot整合rabbitmq的完整代码,代码里包含了原生的rabbitmq使用代码和整合springboot后的使用代码,还有rabbit队列的所有消息队列模式,代码简单易懂,解压打开就可以使用
  • SpringBoot整合RabbitMQ之Spring事件驱动模型-系统源码数据库流程图 SpringBoot整合RabbitMQ实战视频教程:https://edu.csdn.net/course/detail/9314 (感兴趣也可以加QQ联系:1974544863)
  • 一、RabbitMQ入门程序 二、Work queues 工作模式 三、Publish / Subscribe 发布/订阅模式 四、Routing 路由模式 五、Topics 六、Header 七、RPC 八、Spring Data Elasticsearch 一、RabbitMQ入门程序 <...
  • springbootrabbitmq结合的实战、实例项目,有助于帮助你了解springboot中怎么使用rabbitmq。 获取资源:关注我!给我留言或者私信发邮箱~
  • } 配置文件 #配置MQ连接信息 spring.rabbitmq.addresses=192.168.235.128 spring.rabbitmq.port=5672 spring.rabbitmq.username=admin spring.rabbitmq.password=123456 #spring.rabbitmq.publisher-confirm-type=...
  • 配置就不写了,直接上SpringBoot使用rabbitmq,有人看在补配置 rabbitmq常用的6种工作模式 1.简单模式 2.工作模式 3.发布/订阅模式 4.路由模式 5.主题模式 6.RPC模式 rabbitmq使用的常量 队列,交换机,路由都用一个...

空空如也

空空如也

1 2 3 4 5 ... 20
收藏数 16,399
精华内容 6,559
关键字:

springboot整合rabbitmq

spring 订阅