转自:https://blog.csdn.net/qq_35387940/article/details/100514134

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

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

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

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

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

JCccc-RabbitMqRabbitMq -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依赖:

  1.  
    <!–rabbitmq–>
  2.  
    <dependency>
  3.  
    <groupId>org.springframework.boot</groupId>
  4.  
    <artifactId>spring-boot-starter-amqp</artifactId>
  5.  
    </dependency>
  6.  
    <dependency>
  7.  
    <groupId>org.springframework.boot</groupId>
  8.  
    <artifactId>spring-boot-starter-web</artifactId>
  9.  
    </dependency>

然后application.yml:

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

  1.  
    server:
  2.  
    port: 8021
  3.  
    spring:
  4.  
    #给项目来个名字
  5.  
    application:
  6.  
    name: rabbitmq-provider
  7.  
    #配置rabbitMq 服务器
  8.  
    rabbitmq:
  9.  
    host: 127.0.0.1
  10.  
    port: 5672
  11.  
    username: root
  12.  
    password: root
  13.  
    #虚拟host 可以不设置,使用server默认host
  14.  
    virtual-host: JCcccHost

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

  1.  
    import org.springframework.amqp.core.Binding;
  2.  
    import org.springframework.amqp.core.BindingBuilder;
  3.  
    import org.springframework.amqp.core.DirectExchange;
  4.  
    import org.springframework.amqp.core.Queue;
  5.  
    import org.springframework.context.annotation.Bean;
  6.  
    import org.springframework.context.annotation.Configuration;
  7.  
     
  8.  
    /**
  9.  
    * @Author : JCccc
  10.  
    * @CreateTime : 2019/9/3
  11.  
    * @Description :
  12.  
    **/
  13.  
    @Configuration
  14.  
    public class DirectRabbitConfig {
  15.  
     
  16.  
    //队列 起名:TestDirectQueue
  17.  
    @Bean
  18.  
    public Queue TestDirectQueue() {
  19.  
    // durable:是否持久化,默认是false,持久化队列:会被存储在磁盘上,当消息代理重启时仍然存在,暂存队列:当前连接有效
  20.  
    // exclusive:默认也是false,只能被当前创建的连接使用,而且当连接关闭后队列即被删除。此参考优先级高于durable
  21.  
    // autoDelete:是否自动删除,当没有生产者或者消费者使用此队列,该队列会自动删除。
  22.  
    // return new Queue(“TestDirectQueue”,true,true,false);
  23.  
     
  24.  
    //一般设置一下队列的持久化就好,其余两个就是默认false
  25.  
    return new Queue(“TestDirectQueue”,true);
  26.  
    }
  27.  
     
  28.  
    //Direct交换机 起名:TestDirectExchange
  29.  
    @Bean
  30.  
    DirectExchange TestDirectExchange() {
  31.  
    // return new DirectExchange(“TestDirectExchange”,true,true);
  32.  
    return new DirectExchange(“TestDirectExchange”,true,false);
  33.  
    }
  34.  
     
  35.  
    //绑定 将队列和交换机绑定, 并设置用于匹配键:TestDirectRouting
  36.  
    @Bean
  37.  
    Binding bindingDirect() {
  38.  
    return BindingBuilder.bind(TestDirectQueue()).to(TestDirectExchange()).with(“TestDirectRouting”);
  39.  
    }
  40.  
     
  41.  
     
  42.  
     
  43.  
    @Bean
  44.  
    DirectExchange lonelyDirectExchange() {
  45.  
    return new DirectExchange(“lonelyDirectExchange”);
  46.  
    }
  47.  
     
  48.  
     
  49.  
     
  50.  
    }

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

  1.  
    import org.springframework.amqp.rabbit.core.RabbitTemplate;
  2.  
    import org.springframework.beans.factory.annotation.Autowired;
  3.  
    import org.springframework.web.bind.annotation.GetMapping;
  4.  
    import org.springframework.web.bind.annotation.RestController;
  5.  
    import java.time.LocalDateTime;
  6.  
    import java.time.format.DateTimeFormatter;
  7.  
    import java.util.HashMap;
  8.  
    import java.util.Map;
  9.  
    import java.util.UUID;
  10.  
     
  11.  
    /**
  12.  
    * @Author : JCccc
  13.  
    * @CreateTime : 2019/9/3
  14.  
    * @Description :
  15.  
    **/
  16.  
    @RestController
  17.  
    public class SendMessageController {
  18.  
     
  19.  
    @Autowired
  20.  
    RabbitTemplate rabbitTemplate; //使用RabbitTemplate,这提供了接收/发送等等方法
  21.  
     
  22.  
    @GetMapping(“/sendDirectMessage”)
  23.  
    public String sendDirectMessage() {
  24.  
    String messageId = String.valueOf(UUID.randomUUID());
  25.  
    String messageData = “test message, hello!”;
  26.  
    String createTime = LocalDateTime.now().format(DateTimeFormatter.ofPattern(“yyyy-MM-dd HH:mm:ss”));
  27.  
    Map<String,Object> map=new HashMap<>();
  28.  
    map.put(“messageId”,messageId);
  29.  
    map.put(“messageData”,messageData);
  30.  
    map.put(“createTime”,createTime);
  31.  
    //将消息携带绑定键值:TestDirectRouting 发送到交换机TestDirectExchange
  32.  
    rabbitTemplate.convertAndSend(“TestDirectExchange”, “TestDirectRouting”, map);
  33.  
    return “ok”;
  34.  
    }
  35.  
     
  36.  
     
  37.  
    }

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

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


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

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

 


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

pom.xml里的jar依赖:

  1.  
    <!–rabbitmq–>
  2.  
    <dependency>
  3.  
    <groupId>org.springframework.boot</groupId>
  4.  
    <artifactId>spring-boot-starter-amqp</artifactId>
  5.  
    </dependency>
  6.  
    <dependency>
  7.  
    <groupId>org.springframework.boot</groupId>
  8.  
    <artifactId>spring-boot-starter</artifactId>
  9.  
    </dependency>

然后是 application.yml:

  1.  
     
  2.  
    server:
  3.  
    port: 8022
  4.  
    spring:
  5.  
    #给项目来个名字
  6.  
    application:
  7.  
    name: rabbitmq-consumer
  8.  
    #配置rabbitMq 服务器
  9.  
    rabbitmq:
  10.  
    host: 127.0.0.1
  11.  
    port: 5672
  12.  
    username: root
  13.  
    password: root
  14.  
    #虚拟host 可以不设置,使用server默认host
  15.  
    virtual-host: JCcccHost

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

  1.  
    import org.springframework.amqp.core.Binding;
  2.  
    import org.springframework.amqp.core.BindingBuilder;
  3.  
    import org.springframework.amqp.core.DirectExchange;
  4.  
    import org.springframework.amqp.core.Queue;
  5.  
    import org.springframework.context.annotation.Bean;
  6.  
    import org.springframework.context.annotation.Configuration;
  7.  
     
  8.  
    /**
  9.  
    * @Author : JCccc
  10.  
    * @CreateTime : 2019/9/3
  11.  
    * @Description :
  12.  
    **/
  13.  
    @Configuration
  14.  
    public class DirectRabbitConfig {
  15.  
     
  16.  
    //队列 起名:TestDirectQueue
  17.  
    @Bean
  18.  
    public Queue TestDirectQueue() {
  19.  
    return new Queue(“TestDirectQueue”,true);
  20.  
    }
  21.  
     
  22.  
    //Direct交换机 起名:TestDirectExchange
  23.  
    @Bean
  24.  
    DirectExchange TestDirectExchange() {
  25.  
    return new DirectExchange(“TestDirectExchange”);
  26.  
    }
  27.  
     
  28.  
    //绑定 将队列和交换机绑定, 并设置用于匹配键:TestDirectRouting
  29.  
    @Bean
  30.  
    Binding bindingDirect() {
  31.  
    return BindingBuilder.bind(TestDirectQueue()).to(TestDirectExchange()).with(“TestDirectRouting”);
  32.  
    }
  33.  
    }

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

  1.  
    @Component
  2.  
    @RabbitListener(queues = “TestDirectQueue”)//监听的队列名称 TestDirectQueue
  3.  
    public class DirectReceiver {
  4.  
     
  5.  
    @RabbitHandler
  6.  
    public void process(Map testMessage) {
  7.  
    System.out.println(“DirectReceiver消费者收到消息 : “ + testMessage.toString());
  8.  
    }
  9.  
     
  10.  
    }

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

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

 

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

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

 

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

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


  1.  
    import org.springframework.amqp.core.Binding;
  2.  
    import org.springframework.amqp.core.BindingBuilder;
  3.  
    import org.springframework.amqp.core.Queue;
  4.  
    import org.springframework.amqp.core.TopicExchange;
  5.  
    import org.springframework.context.annotation.Bean;
  6.  
    import org.springframework.context.annotation.Configuration;
  7.  
     
  8.  
    /**
  9.  
    * @Author : JCccc
  10.  
    * @CreateTime : 2019/9/3
  11.  
    * @Description :
  12.  
    **/
  13.  
     
  14.  
    @Configuration
  15.  
    public class TopicRabbitConfig {
  16.  
    //绑定键
  17.  
    public final static String man = “topic.man”;
  18.  
    public final static String woman = “topic.woman”;
  19.  
     
  20.  
    @Bean
  21.  
    public Queue firstQueue() {
  22.  
    return new Queue(TopicRabbitConfig.man);
  23.  
    }
  24.  
     
  25.  
    @Bean
  26.  
    public Queue secondQueue() {
  27.  
    return new Queue(TopicRabbitConfig.woman);
  28.  
    }
  29.  
     
  30.  
    @Bean
  31.  
    TopicExchange exchange() {
  32.  
    return new TopicExchange(“topicExchange”);
  33.  
    }
  34.  
     
  35.  
     
  36.  
    //将firstQueue和topicExchange绑定,而且绑定的键值为topic.man
  37.  
    //这样只要是消息携带的路由键是topic.man,才会分发到该队列
  38.  
    @Bean
  39.  
    Binding bindingExchangeMessage() {
  40.  
    return BindingBuilder.bind(firstQueue()).to(exchange()).with(man);
  41.  
    }
  42.  
     
  43.  
    //将secondQueue和topicExchange绑定,而且绑定的键值为用上通配路由键规则topic.#
  44.  
    // 这样只要是消息携带的路由键是以topic.开头,都会分发到该队列
  45.  
    @Bean
  46.  
    Binding bindingExchangeMessage2() {
  47.  
    return BindingBuilder.bind(secondQueue()).to(exchange()).with(“topic.#”);
  48.  
    }
  49.  
     
  50.  
    }

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

  1.  
    @GetMapping(“/sendTopicMessage1”)
  2.  
    public String sendTopicMessage1() {
  3.  
    String messageId = String.valueOf(UUID.randomUUID());
  4.  
    String messageData = “message: M A N “;
  5.  
    String createTime = LocalDateTime.now().format(DateTimeFormatter.ofPattern(“yyyy-MM-dd HH:mm:ss”));
  6.  
    Map<String, Object> manMap = new HashMap<>();
  7.  
    manMap.put(“messageId”, messageId);
  8.  
    manMap.put(“messageData”, messageData);
  9.  
    manMap.put(“createTime”, createTime);
  10.  
    rabbitTemplate.convertAndSend(“topicExchange”, “topic.man”, manMap);
  11.  
    return “ok”;
  12.  
    }
  13.  
     
  14.  
    @GetMapping(“/sendTopicMessage2”)
  15.  
    public String sendTopicMessage2() {
  16.  
    String messageId = String.valueOf(UUID.randomUUID());
  17.  
    String messageData = “message: woman is all “;
  18.  
    String createTime = LocalDateTime.now().format(DateTimeFormatter.ofPattern(“yyyy-MM-dd HH:mm:ss”));
  19.  
    Map<String, Object> womanMap = new HashMap<>();
  20.  
    womanMap.put(“messageId”, messageId);
  21.  
    womanMap.put(“messageData”, messageData);
  22.  
    womanMap.put(“createTime”, createTime);
  23.  
    rabbitTemplate.convertAndSend(“topicExchange”, “topic.woman”, womanMap);
  24.  
    return “ok”;
  25.  
    }
  26.  
    }

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

  1.  
    import org.springframework.amqp.rabbit.annotation.RabbitHandler;
  2.  
    import org.springframework.amqp.rabbit.annotation.RabbitListener;
  3.  
    import org.springframework.stereotype.Component;
  4.  
    import java.util.Map;
  5.  
     
  6.  
    /**
  7.  
    * @Author : JCccc
  8.  
    * @CreateTime : 2019/9/3
  9.  
    * @Description :
  10.  
    **/
  11.  
    @Component
  12.  
    @RabbitListener(queues = “topic.man”)
  13.  
    public class TopicManReceiver {
  14.  
     
  15.  
    @RabbitHandler
  16.  
    public void process(Map testMessage) {
  17.  
    System.out.println(“TopicManReceiver消费者收到消息 : “ + testMessage.toString());
  18.  
    }
  19.  
    }

再创建一个TopicTotalReceiver.java:

  1.  
    package com.elegant.rabbitmqconsumer.receiver;
  2.  
     
  3.  
    import org.springframework.amqp.rabbit.annotation.RabbitHandler;
  4.  
    import org.springframework.amqp.rabbit.annotation.RabbitListener;
  5.  
    import org.springframework.stereotype.Component;
  6.  
    import java.util.Map;
  7.  
     
  8.  
    /**
  9.  
    * @Author : JCccc
  10.  
    * @CreateTime : 2019/9/3
  11.  
    * @Description :
  12.  
    **/
  13.  
     
  14.  
    @Component
  15.  
    @RabbitListener(queues = “topic.woman”)
  16.  
    public class TopicTotalReceiver {
  17.  
     
  18.  
    @RabbitHandler
  19.  
    public void process(Map testMessage) {
  20.  
    System.out.println(“TopicTotalReceiver消费者收到消息 : “ + testMessage.toString());
  21.  
    }
  22.  
    }

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

  1.  
    import org.springframework.amqp.core.Binding;
  2.  
    import org.springframework.amqp.core.BindingBuilder;
  3.  
    import org.springframework.amqp.core.Queue;
  4.  
    import org.springframework.amqp.core.TopicExchange;
  5.  
    import org.springframework.context.annotation.Bean;
  6.  
    import org.springframework.context.annotation.Configuration;
  7.  
     
  8.  
    /**
  9.  
    * @Author : JCccc
  10.  
    * @CreateTime : 2019/9/3
  11.  
    * @Description :
  12.  
    **/
  13.  
     
  14.  
    @Configuration
  15.  
    public class TopicRabbitConfig {
  16.  
    //绑定键
  17.  
    public final static String man = “topic.man”;
  18.  
    public final static String woman = “topic.woman”;
  19.  
     
  20.  
    @Bean
  21.  
    public Queue firstQueue() {
  22.  
    return new Queue(TopicRabbitConfig.man);
  23.  
    }
  24.  
     
  25.  
    @Bean
  26.  
    public Queue secondQueue() {
  27.  
    return new Queue(TopicRabbitConfig.woman);
  28.  
    }
  29.  
     
  30.  
    @Bean
  31.  
    TopicExchange exchange() {
  32.  
    return new TopicExchange(“topicExchange”);
  33.  
    }
  34.  
     
  35.  
     
  36.  
    //将firstQueue和topicExchange绑定,而且绑定的键值为topic.man
  37.  
    //这样只要是消息携带的路由键是topic.man,才会分发到该队列
  38.  
    @Bean
  39.  
    Binding bindingExchangeMessage() {
  40.  
    return BindingBuilder.bind(firstQueue()).to(exchange()).with(man);
  41.  
    }
  42.  
     
  43.  
    //将secondQueue和topicExchange绑定,而且绑定的键值为用上通配路由键规则topic.#
  44.  
    // 这样只要是消息携带的路由键是以topic.开头,都会分发到该队列
  45.  
    @Bean
  46.  
    Binding bindingExchangeMessage2() {
  47.  
    return BindingBuilder.bind(secondQueue()).to(exchange()).with(“topic.#”);
  48.  
    }
  49.  
     
  50.  
    }


然后把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:

  1.  
    import org.springframework.amqp.core.Binding;
  2.  
    import org.springframework.amqp.core.BindingBuilder;
  3.  
    import org.springframework.amqp.core.FanoutExchange;
  4.  
    import org.springframework.amqp.core.Queue;
  5.  
    import org.springframework.context.annotation.Bean;
  6.  
    import org.springframework.context.annotation.Configuration;
  7.  
    /**
  8.  
    * @Author : JCccc
  9.  
    * @CreateTime : 2019/9/3
  10.  
    * @Description :
  11.  
    **/
  12.  
     
  13.  
    @Configuration
  14.  
    public class FanoutRabbitConfig {
  15.  
     
  16.  
    /**
  17.  
    * 创建三个队列 :fanout.A fanout.B fanout.C
  18.  
    * 将三个队列都绑定在交换机 fanoutExchange 上
  19.  
    * 因为是扇型交换机, 路由键无需配置,配置也不起作用
  20.  
    */
  21.  
     
  22.  
     
  23.  
    @Bean
  24.  
    public Queue queueA() {
  25.  
    return new Queue(“fanout.A”);
  26.  
    }
  27.  
     
  28.  
    @Bean
  29.  
    public Queue queueB() {
  30.  
    return new Queue(“fanout.B”);
  31.  
    }
  32.  
     
  33.  
    @Bean
  34.  
    public Queue queueC() {
  35.  
    return new Queue(“fanout.C”);
  36.  
    }
  37.  
     
  38.  
    @Bean
  39.  
    FanoutExchange fanoutExchange() {
  40.  
    return new FanoutExchange(“fanoutExchange”);
  41.  
    }
  42.  
     
  43.  
    @Bean
  44.  
    Binding bindingExchangeA() {
  45.  
    return BindingBuilder.bind(queueA()).to(fanoutExchange());
  46.  
    }
  47.  
     
  48.  
    @Bean
  49.  
    Binding bindingExchangeB() {
  50.  
    return BindingBuilder.bind(queueB()).to(fanoutExchange());
  51.  
    }
  52.  
     
  53.  
    @Bean
  54.  
    Binding bindingExchangeC() {
  55.  
    return BindingBuilder.bind(queueC()).to(fanoutExchange());
  56.  
    }
  57.  
    }

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

  1.  
    @GetMapping(“/sendFanoutMessage”)
  2.  
    public String sendFanoutMessage() {
  3.  
    String messageId = String.valueOf(UUID.randomUUID());
  4.  
    String messageData = “message: testFanoutMessage “;
  5.  
    String createTime = LocalDateTime.now().format(DateTimeFormatter.ofPattern(“yyyy-MM-dd HH:mm:ss”));
  6.  
    Map<String, Object> map = new HashMap<>();
  7.  
    map.put(“messageId”, messageId);
  8.  
    map.put(“messageData”, messageData);
  9.  
    map.put(“createTime”, createTime);
  10.  
    rabbitTemplate.convertAndSend(“fanoutExchange”, null, map);
  11.  
    return “ok”;
  12.  
    }

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

FanoutReceiverA.java:

  1.  
    import org.springframework.amqp.rabbit.annotation.RabbitHandler;
  2.  
    import org.springframework.amqp.rabbit.annotation.RabbitListener;
  3.  
    import org.springframework.stereotype.Component;
  4.  
    import java.util.Map;
  5.  
    /**
  6.  
    * @Author : JCccc
  7.  
    * @CreateTime : 2019/9/3
  8.  
    * @Description :
  9.  
    **/
  10.  
    @Component
  11.  
    @RabbitListener(queues = “fanout.A”)
  12.  
    public class FanoutReceiverA {
  13.  
     
  14.  
    @RabbitHandler
  15.  
    public void process(Map testMessage) {
  16.  
    System.out.println(“FanoutReceiverA消费者收到消息 : “ +testMessage.toString());
  17.  
    }
  18.  
     
  19.  
    }

FanoutReceiverB.java:

  1.  
    import org.springframework.amqp.rabbit.annotation.RabbitHandler;
  2.  
    import org.springframework.amqp.rabbit.annotation.RabbitListener;
  3.  
    import org.springframework.stereotype.Component;
  4.  
    import java.util.Map;
  5.  
    /**
  6.  
    * @Author : JCccc
  7.  
    * @CreateTime : 2019/9/3
  8.  
    * @Description :
  9.  
    **/
  10.  
    @Component
  11.  
    @RabbitListener(queues = “fanout.B”)
  12.  
    public class FanoutReceiverB {
  13.  
     
  14.  
    @RabbitHandler
  15.  
    public void process(Map testMessage) {
  16.  
    System.out.println(“FanoutReceiverB消费者收到消息 : “ +testMessage.toString());
  17.  
    }
  18.  
     
  19.  
    }

FanoutReceiverC.java:

  1.  
    import org.springframework.amqp.rabbit.annotation.RabbitHandler;
  2.  
    import org.springframework.amqp.rabbit.annotation.RabbitListener;
  3.  
    import org.springframework.stereotype.Component;
  4.  
    import java.util.Map;
  5.  
     
  6.  
    /**
  7.  
    * @Author : JCccc
  8.  
    * @CreateTime : 2019/9/3
  9.  
    * @Description :
  10.  
    **/
  11.  
    @Component
  12.  
    @RabbitListener(queues = “fanout.C”)
  13.  
    public class FanoutReceiverC {
  14.  
     
  15.  
    @RabbitHandler
  16.  
    public void process(Map testMessage) {
  17.  
    System.out.println(“FanoutReceiverC消费者收到消息 : “ +testMessage.toString());
  18.  
    }
  19.  
     
  20.  
    }

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

  1.  
    import org.springframework.amqp.core.Binding;
  2.  
    import org.springframework.amqp.core.BindingBuilder;
  3.  
    import org.springframework.amqp.core.FanoutExchange;
  4.  
    import org.springframework.amqp.core.Queue;
  5.  
    import org.springframework.context.annotation.Bean;
  6.  
    import org.springframework.context.annotation.Configuration;
  7.  
    /**
  8.  
    * @Author : JCccc
  9.  
    * @CreateTime : 2019/9/3
  10.  
    * @Description :
  11.  
    **/
  12.  
    @Configuration
  13.  
    public class FanoutRabbitConfig {
  14.  
     
  15.  
    /**
  16.  
    * 创建三个队列 :fanout.A fanout.B fanout.C
  17.  
    * 将三个队列都绑定在交换机 fanoutExchange 上
  18.  
    * 因为是扇型交换机, 路由键无需配置,配置也不起作用
  19.  
    */
  20.  
     
  21.  
     
  22.  
    @Bean
  23.  
    public Queue queueA() {
  24.  
    return new Queue(“fanout.A”);
  25.  
    }
  26.  
     
  27.  
    @Bean
  28.  
    public Queue queueB() {
  29.  
    return new Queue(“fanout.B”);
  30.  
    }
  31.  
     
  32.  
    @Bean
  33.  
    public Queue queueC() {
  34.  
    return new Queue(“fanout.C”);
  35.  
    }
  36.  
     
  37.  
    @Bean
  38.  
    FanoutExchange fanoutExchange() {
  39.  
    return new FanoutExchange(“fanoutExchange”);
  40.  
    }
  41.  
     
  42.  
    @Bean
  43.  
    Binding bindingExchangeA() {
  44.  
    return BindingBuilder.bind(queueA()).to(fanoutExchange());
  45.  
    }
  46.  
     
  47.  
    @Bean
  48.  
    Binding bindingExchangeB() {
  49.  
    return BindingBuilder.bind(queueB()).to(fanoutExchange());
  50.  
    }
  51.  
     
  52.  
    @Bean
  53.  
    Binding bindingExchangeC() {
  54.  
    return BindingBuilder.bind(queueC()).to(fanoutExchange());
  55.  
    }
  56.  
    }

最后将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

  1.  
    server:
  2.  
    port: 8021
  3.  
    spring:
  4.  
    #给项目来个名字
  5.  
    application:
  6.  
    name: rabbitmq-provider
  7.  
    #配置rabbitMq 服务器
  8.  
    rabbitmq:
  9.  
    host: 127.0.0.1
  10.  
    port: 5672
  11.  
    username: root
  12.  
    password: root
  13.  
    #虚拟host 可以不设置,使用server默认host
  14.  
    virtual-host: JCcccHost
  15.  
    #消息确认配置项
  16.  
     
  17.  
    #确认消息已发送到交换机(Exchange)
  18.  
    publisher-confirms: true
  19.  
    #确认消息已发送到队列(Queue)
  20.  
    publisher-returns: true

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

  1.  
    import org.springframework.amqp.core.Message;
  2.  
    import org.springframework.amqp.rabbit.connection.ConnectionFactory;
  3.  
    import org.springframework.amqp.rabbit.connection.CorrelationData;
  4.  
    import org.springframework.amqp.rabbit.core.RabbitTemplate;
  5.  
    import org.springframework.context.annotation.Bean;
  6.  
    import org.springframework.context.annotation.Configuration;
  7.  
     
  8.  
     
  9.  
    /**
  10.  
    * @Author : JCccc
  11.  
    * @CreateTime : 2019/9/3
  12.  
    * @Description :
  13.  
    **/
  14.  
    @Configuration
  15.  
    public class RabbitConfig {
  16.  
     
  17.  
    @Bean
  18.  
    public RabbitTemplate createRabbitTemplate(ConnectionFactory connectionFactory){
  19.  
    RabbitTemplate rabbitTemplate = new RabbitTemplate();
  20.  
    rabbitTemplate.setConnectionFactory(connectionFactory);
  21.  
    //设置开启Mandatory,才能触发回调函数,无论消息推送结果怎么样都强制调用回调函数
  22.  
    rabbitTemplate.setMandatory(true);
  23.  
     
  24.  
    rabbitTemplate.setConfirmCallback(new RabbitTemplate.ConfirmCallback() {
  25.  
    @Override
  26.  
    public void confirm(CorrelationData correlationData, boolean ack, String cause) {
  27.  
    System.out.println(“ConfirmCallback: “+“相关数据:”+correlationData);
  28.  
    System.out.println(“ConfirmCallback: “+“确认情况:”+ack);
  29.  
    System.out.println(“ConfirmCallback: “+“原因:”+cause);
  30.  
    }
  31.  
    });
  32.  
     
  33.  
    rabbitTemplate.setReturnCallback(new RabbitTemplate.ReturnCallback() {
  34.  
    @Override
  35.  
    public void returnedMessage(Message message, int replyCode, String replyText, String exchange, String routingKey) {
  36.  
    System.out.println(“ReturnCallback: “+“消息:”+message);
  37.  
    System.out.println(“ReturnCallback: “+“回应码:”+replyCode);
  38.  
    System.out.println(“ReturnCallback: “+“回应信息:”+replyText);
  39.  
    System.out.println(“ReturnCallback: “+“交换机:”+exchange);
  40.  
    System.out.println(“ReturnCallback: “+“路由键:”+routingKey);
  41.  
    }
  42.  
    });
  43.  
     
  44.  
    return rabbitTemplate;
  45.  
    }
  46.  
     
  47.  
    }

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

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

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

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

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

  1.  
    @GetMapping(“/TestMessageAck”)
  2.  
    public String TestMessageAck() {
  3.  
    String messageId = String.valueOf(UUID.randomUUID());
  4.  
    String messageData = “message: non-existent-exchange test message “;
  5.  
    String createTime = LocalDateTime.now().format(DateTimeFormatter.ofPattern(“yyyy-MM-dd HH:mm:ss”));
  6.  
    Map<String, Object> map = new HashMap<>();
  7.  
    map.put(“messageId”, messageId);
  8.  
    map.put(“messageData”, messageData);
  9.  
    map.put(“createTime”, createTime);
  10.  
    rabbitTemplate.convertAndSend(“non-existent-exchange”, “TestDirectRouting”, map);
  11.  
    return “ok”;
  12.  
    }

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

  1.  
    20190904 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’, classid=60, method-id=40)
  2.  
    ConfirmCallback: 相关数据:null
  3.  
    ConfirmCallback: 确认情况:false
  4.  
    ConfirmCallback: 原因:channel error; protocol method: #method<channel.close>(reply-code=404, reply-text=NOT_FOUND – no exchange ‘non-existent-exchange’ in vhost ‘JCcccHost’, classid=60, method-id=40)

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

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

  1.  
    @Bean
  2.  
    DirectExchange lonelyDirectExchange() {
  3.  
    return new DirectExchange(“lonelyDirectExchange”);
  4.  
    }

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

  1.  
    @GetMapping(“/TestMessageAck2”)
  2.  
    public String TestMessageAck2() {
  3.  
    String messageId = String.valueOf(UUID.randomUUID());
  4.  
    String messageData = “message: lonelyDirectExchange test message “;
  5.  
    String createTime = LocalDateTime.now().format(DateTimeFormatter.ofPattern(“yyyy-MM-dd HH:mm:ss”));
  6.  
    Map<String, Object> map = new HashMap<>();
  7.  
    map.put(“messageId”, messageId);
  8.  
    map.put(“messageData”, messageData);
  9.  
    map.put(“createTime”, createTime);
  10.  
    rabbitTemplate.convertAndSend(“lonelyDirectExchange”, “TestDirectRouting”, map);
  11.  
    return “ok”;
  12.  
    }

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

  1.  
    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])
  2.  
    ReturnCallback: 回应码:312
  3.  
    ReturnCallback: 回应信息:NO_ROUTE
  4.  
    ReturnCallback: 交换机:lonelyDirectExchange
  5.  
    ReturnCallback: 路由键:TestDirectRouting
  1.  
    ConfirmCallback: 相关数据:null
  2.  
    ConfirmCallback: 确认情况:true
  3.  
    ConfirmCallback: 原因:null

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

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

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

  1.  
    ConfirmCallback: 相关数据:null
  2.  
    ConfirmCallback: 确认情况:true
  3.  
    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上添加代码相关的配置代码:

  1.  
     
  2.  
    import com.elegant.rabbitmqconsumer.receiver.MyAckReceiver;
  3.  
    import org.springframework.amqp.core.AcknowledgeMode;
  4.  
    import org.springframework.amqp.core.Queue;
  5.  
    import org.springframework.amqp.rabbit.connection.CachingConnectionFactory;
  6.  
    import org.springframework.amqp.rabbit.listener.SimpleMessageListenerContainer;
  7.  
    import org.springframework.beans.factory.annotation.Autowired;
  8.  
    import org.springframework.context.annotation.Bean;
  9.  
    import org.springframework.context.annotation.Configuration;
  10.  
     
  11.  
    /**
  12.  
    * @Author : JCccc
  13.  
    * @CreateTime : 2019/9/4
  14.  
    * @Description :
  15.  
    **/
  16.  
    @Configuration
  17.  
    public class MessageListenerConfig {
  18.  
     
  19.  
    @Autowired
  20.  
    private CachingConnectionFactory connectionFactory;
  21.  
    @Autowired
  22.  
    private MyAckReceiver myAckReceiver;//消息接收处理类
  23.  
     
  24.  
    @Bean
  25.  
    public SimpleMessageListenerContainer simpleMessageListenerContainer() {
  26.  
    SimpleMessageListenerContainer container = new SimpleMessageListenerContainer(connectionFactory);
  27.  
    container.setConcurrentConsumers(1);
  28.  
    container.setMaxConcurrentConsumers(1);
  29.  
    container.setAcknowledgeMode(AcknowledgeMode.MANUAL); // RabbitMQ默认是自动确认,这里改为手动确认消息
  30.  
    //设置一个队列
  31.  
    container.setQueueNames(“TestDirectQueue”);
  32.  
    //如果同时设置多个如下: 前提是队列都是必须已经创建存在的
  33.  
    // container.setQueueNames(“TestDirectQueue”,”TestDirectQueue2″,”TestDirectQueue3″);
  34.  
     
  35.  
     
  36.  
    //另一种设置队列的方法,如果使用这种情况,那么要设置多个,就使用addQueues
  37.  
    //container.setQueues(new Queue(“TestDirectQueue”,true));
  38.  
    //container.addQueues(new Queue(“TestDirectQueue2”,true));
  39.  
    //container.addQueues(new Queue(“TestDirectQueue3”,true));
  40.  
    container.setMessageListener(myAckReceiver);
  41.  
     
  42.  
    return container;
  43.  
    }
  44.  
     
  45.  
     
  46.  
    }

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

  1.  
    import com.rabbitmq.client.Channel;
  2.  
    import org.springframework.amqp.core.Message;
  3.  
    import org.springframework.amqp.rabbit.listener.api.ChannelAwareMessageListener;
  4.  
    import org.springframework.stereotype.Component;
  5.  
    import java.util.HashMap;
  6.  
    import java.util.Map;
  7.  
     
  8.  
    @Component
  9.  
     
  10.  
    public class MyAckReceiver implements ChannelAwareMessageListener {
  11.  
     
  12.  
    @Override
  13.  
    public void onMessage(Message message, Channel channel) throws Exception {
  14.  
    long deliveryTag = message.getMessageProperties().getDeliveryTag();
  15.  
    try {
  16.  
    //因为传递消息的时候用的map传递,所以将Map从Message内取出需要做些处理
  17.  
    String msg = message.toString();
  18.  
    String[] msgArray = msg.split(“‘”);//可以点进Message里面看源码,单引号直接的数据就是我们的map消息数据
  19.  
    Map<String, String> msgMap = mapStringToMap(msgArray[1].trim(),3);
  20.  
    String messageId=msgMap.get(“messageId”);
  21.  
    String messageData=msgMap.get(“messageData”);
  22.  
    String createTime=msgMap.get(“createTime”);
  23.  
    System.out.println(” MyAckReceiver messageId:”+messageId+” messageData:”+messageData+” createTime:”+createTime);
  24.  
    System.out.println(“消费的主题消息来自:”+message.getMessageProperties().getConsumerQueue());
  25.  
    channel.basicAck(deliveryTag, true); //第二个参数,手动确认可以被批处理,当该参数为 true 时,则可以一次性确认 delivery_tag 小于等于传入值的所有消息
  26.  
    // channel.basicReject(deliveryTag, true);//第二个参数,true会重新放回队列,所以需要自己根据业务逻辑判断什么时候使用拒绝
  27.  
    } catch (Exception e) {
  28.  
    channel.basicReject(deliveryTag, false);
  29.  
    e.printStackTrace();
  30.  
    }
  31.  
    }
  32.  
     
  33.  
    //{key=value,key=value,key=value} 格式转换成map
  34.  
    private Map<String, String> mapStringToMap(String str,int entryNum ) {
  35.  
    str = str.substring(1, str.length() – 1);
  36.  
    String[] strs = str.split(“,”,entryNum);
  37.  
    Map<String, String> map = new HashMap<String, String>();
  38.  
    for (String string : strs) {
  39.  
    String key = string.split(“=”)[0].trim();
  40.  
    String value = string.split(“=”)[1];
  41.  
    map.put(key, value);
  42.  
    }
  43.  
    return map;
  44.  
    }
  45.  
    }

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

 

 

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

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

没有问题,接下来看代码

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

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

 

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

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

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

  1.  
    import com.rabbitmq.client.Channel;
  2.  
    import org.springframework.amqp.core.Message;
  3.  
    import org.springframework.amqp.rabbit.listener.api.ChannelAwareMessageListener;
  4.  
    import org.springframework.stereotype.Component;
  5.  
    import java.util.HashMap;
  6.  
    import java.util.Map;
  7.  
     
  8.  
    @Component
  9.  
    public class MyAckReceiver implements ChannelAwareMessageListener {
  10.  
     
  11.  
    @Override
  12.  
    public void onMessage(Message message, Channel channel) throws Exception {
  13.  
    long deliveryTag = message.getMessageProperties().getDeliveryTag();
  14.  
    try {
  15.  
    //因为传递消息的时候用的map传递,所以将Map从Message内取出需要做些处理
  16.  
    String msg = message.toString();
  17.  
    String[] msgArray = msg.split(“‘”);//可以点进Message里面看源码,单引号直接的数据就是我们的map消息数据
  18.  
    Map<String, String> msgMap = mapStringToMap(msgArray[1].trim(),3);
  19.  
    String messageId=msgMap.get(“messageId”);
  20.  
    String messageData=msgMap.get(“messageData”);
  21.  
    String createTime=msgMap.get(“createTime”);
  22.  
     
  23.  
    if (“TestDirectQueue”.equals(message.getMessageProperties().getConsumerQueue())){
  24.  
    System.out.println(“消费的消息来自的队列名为:”+message.getMessageProperties().getConsumerQueue());
  25.  
    System.out.println(“消息成功消费到 messageId:”+messageId+” messageData:”+messageData+” createTime:”+createTime);
  26.  
    System.out.println(“执行TestDirectQueue中的消息的业务处理流程……”);
  27.  
     
  28.  
    }
  29.  
     
  30.  
    if (“fanout.A”.equals(message.getMessageProperties().getConsumerQueue())){
  31.  
    System.out.println(“消费的消息来自的队列名为:”+message.getMessageProperties().getConsumerQueue());
  32.  
    System.out.println(“消息成功消费到 messageId:”+messageId+” messageData:”+messageData+” createTime:”+createTime);
  33.  
    System.out.println(“执行fanout.A中的消息的业务处理流程……”);
  34.  
     
  35.  
    }
  36.  
     
  37.  
    channel.basicAck(deliveryTag, true);
  38.  
    // channel.basicReject(deliveryTag, true);//为true会重新放回队列
  39.  
    } catch (Exception e) {
  40.  
    channel.basicReject(deliveryTag, false);
  41.  
    e.printStackTrace();
  42.  
    }
  43.  
    }
  44.  
     
  45.  
    //{key=value,key=value,key=value} 格式转换成map
  46.  
    private Map<String, String> mapStringToMap(String str,int enNum) {
  47.  
    str = str.substring(1, str.length() – 1);
  48.  
    String[] strs = str.split(“,”,enNum);
  49.  
    Map<String, String> map = new HashMap<String, String>();
  50.  
    for (String string : strs) {
  51.  
    String key = string.split(“=”)[0].trim();
  52.  
    String value = string.split(“=”)[1];
  53.  
    map.put(key, value);
  54.  
    }
  55.  
    return map;
  56.  
    }
  57.  
    }

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

调用接口/sendDirectMessage  和 /sendFanoutMessage ,

 

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

 

 

 

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

 

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

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

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

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

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

JCccc-RabbitMqRabbitMq -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依赖:

  1.  
    <!–rabbitmq–>
  2.  
    <dependency>
  3.  
    <groupId>org.springframework.boot</groupId>
  4.  
    <artifactId>spring-boot-starter-amqp</artifactId>
  5.  
    </dependency>
  6.  
    <dependency>
  7.  
    <groupId>org.springframework.boot</groupId>
  8.  
    <artifactId>spring-boot-starter-web</artifactId>
  9.  
    </dependency>

然后application.yml:

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

  1.  
    server:
  2.  
    port: 8021
  3.  
    spring:
  4.  
    #给项目来个名字
  5.  
    application:
  6.  
    name: rabbitmq-provider
  7.  
    #配置rabbitMq 服务器
  8.  
    rabbitmq:
  9.  
    host: 127.0.0.1
  10.  
    port: 5672
  11.  
    username: root
  12.  
    password: root
  13.  
    #虚拟host 可以不设置,使用server默认host
  14.  
    virtual-host: JCcccHost

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

  1.  
    import org.springframework.amqp.core.Binding;
  2.  
    import org.springframework.amqp.core.BindingBuilder;
  3.  
    import org.springframework.amqp.core.DirectExchange;
  4.  
    import org.springframework.amqp.core.Queue;
  5.  
    import org.springframework.context.annotation.Bean;
  6.  
    import org.springframework.context.annotation.Configuration;
  7.  
     
  8.  
    /**
  9.  
    * @Author : JCccc
  10.  
    * @CreateTime : 2019/9/3
  11.  
    * @Description :
  12.  
    **/
  13.  
    @Configuration
  14.  
    public class DirectRabbitConfig {
  15.  
     
  16.  
    //队列 起名:TestDirectQueue
  17.  
    @Bean
  18.  
    public Queue TestDirectQueue() {
  19.  
    // durable:是否持久化,默认是false,持久化队列:会被存储在磁盘上,当消息代理重启时仍然存在,暂存队列:当前连接有效
  20.  
    // exclusive:默认也是false,只能被当前创建的连接使用,而且当连接关闭后队列即被删除。此参考优先级高于durable
  21.  
    // autoDelete:是否自动删除,当没有生产者或者消费者使用此队列,该队列会自动删除。
  22.  
    // return new Queue(“TestDirectQueue”,true,true,false);
  23.  
     
  24.  
    //一般设置一下队列的持久化就好,其余两个就是默认false
  25.  
    return new Queue(“TestDirectQueue”,true);
  26.  
    }
  27.  
     
  28.  
    //Direct交换机 起名:TestDirectExchange
  29.  
    @Bean
  30.  
    DirectExchange TestDirectExchange() {
  31.  
    // return new DirectExchange(“TestDirectExchange”,true,true);
  32.  
    return new DirectExchange(“TestDirectExchange”,true,false);
  33.  
    }
  34.  
     
  35.  
    //绑定 将队列和交换机绑定, 并设置用于匹配键:TestDirectRouting
  36.  
    @Bean
  37.  
    Binding bindingDirect() {
  38.  
    return BindingBuilder.bind(TestDirectQueue()).to(TestDirectExchange()).with(“TestDirectRouting”);
  39.  
    }
  40.  
     
  41.  
     
  42.  
     
  43.  
    @Bean
  44.  
    DirectExchange lonelyDirectExchange() {
  45.  
    return new DirectExchange(“lonelyDirectExchange”);
  46.  
    }
  47.  
     
  48.  
     
  49.  
     
  50.  
    }

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

  1.  
    import org.springframework.amqp.rabbit.core.RabbitTemplate;
  2.  
    import org.springframework.beans.factory.annotation.Autowired;
  3.  
    import org.springframework.web.bind.annotation.GetMapping;
  4.  
    import org.springframework.web.bind.annotation.RestController;
  5.  
    import java.time.LocalDateTime;
  6.  
    import java.time.format.DateTimeFormatter;
  7.  
    import java.util.HashMap;
  8.  
    import java.util.Map;
  9.  
    import java.util.UUID;
  10.  
     
  11.  
    /**
  12.  
    * @Author : JCccc
  13.  
    * @CreateTime : 2019/9/3
  14.  
    * @Description :
  15.  
    **/
  16.  
    @RestController
  17.  
    public class SendMessageController {
  18.  
     
  19.  
    @Autowired
  20.  
    RabbitTemplate rabbitTemplate; //使用RabbitTemplate,这提供了接收/发送等等方法
  21.  
     
  22.  
    @GetMapping(“/sendDirectMessage”)
  23.  
    public String sendDirectMessage() {
  24.  
    String messageId = String.valueOf(UUID.randomUUID());
  25.  
    String messageData = “test message, hello!”;
  26.  
    String createTime = LocalDateTime.now().format(DateTimeFormatter.ofPattern(“yyyy-MM-dd HH:mm:ss”));
  27.  
    Map<String,Object> map=new HashMap<>();
  28.  
    map.put(“messageId”,messageId);
  29.  
    map.put(“messageData”,messageData);
  30.  
    map.put(“createTime”,createTime);
  31.  
    //将消息携带绑定键值:TestDirectRouting 发送到交换机TestDirectExchange
  32.  
    rabbitTemplate.convertAndSend(“TestDirectExchange”, “TestDirectRouting”, map);
  33.  
    return “ok”;
  34.  
    }
  35.  
     
  36.  
     
  37.  
    }

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

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


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

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

 


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

pom.xml里的jar依赖:

  1.  
    <!–rabbitmq–>
  2.  
    <dependency>
  3.  
    <groupId>org.springframework.boot</groupId>
  4.  
    <artifactId>spring-boot-starter-amqp</artifactId>
  5.  
    </dependency>
  6.  
    <dependency>
  7.  
    <groupId>org.springframework.boot</groupId>
  8.  
    <artifactId>spring-boot-starter</artifactId>
  9.  
    </dependency>

然后是 application.yml:

  1.  
     
  2.  
    server:
  3.  
    port: 8022
  4.  
    spring:
  5.  
    #给项目来个名字
  6.  
    application:
  7.  
    name: rabbitmq-consumer
  8.  
    #配置rabbitMq 服务器
  9.  
    rabbitmq:
  10.  
    host: 127.0.0.1
  11.  
    port: 5672
  12.  
    username: root
  13.  
    password: root
  14.  
    #虚拟host 可以不设置,使用server默认host
  15.  
    virtual-host: JCcccHost

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

  1.  
    import org.springframework.amqp.core.Binding;
  2.  
    import org.springframework.amqp.core.BindingBuilder;
  3.  
    import org.springframework.amqp.core.DirectExchange;
  4.  
    import org.springframework.amqp.core.Queue;
  5.  
    import org.springframework.context.annotation.Bean;
  6.  
    import org.springframework.context.annotation.Configuration;
  7.  
     
  8.  
    /**
  9.  
    * @Author : JCccc
  10.  
    * @CreateTime : 2019/9/3
  11.  
    * @Description :
  12.  
    **/
  13.  
    @Configuration
  14.  
    public class DirectRabbitConfig {
  15.  
     
  16.  
    //队列 起名:TestDirectQueue
  17.  
    @Bean
  18.  
    public Queue TestDirectQueue() {
  19.  
    return new Queue(“TestDirectQueue”,true);
  20.  
    }
  21.  
     
  22.  
    //Direct交换机 起名:TestDirectExchange
  23.  
    @Bean
  24.  
    DirectExchange TestDirectExchange() {
  25.  
    return new DirectExchange(“TestDirectExchange”);
  26.  
    }
  27.  
     
  28.  
    //绑定 将队列和交换机绑定, 并设置用于匹配键:TestDirectRouting
  29.  
    @Bean
  30.  
    Binding bindingDirect() {
  31.  
    return BindingBuilder.bind(TestDirectQueue()).to(TestDirectExchange()).with(“TestDirectRouting”);
  32.  
    }
  33.  
    }

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

  1.  
    @Component
  2.  
    @RabbitListener(queues = “TestDirectQueue”)//监听的队列名称 TestDirectQueue
  3.  
    public class DirectReceiver {
  4.  
     
  5.  
    @RabbitHandler
  6.  
    public void process(Map testMessage) {
  7.  
    System.out.println(“DirectReceiver消费者收到消息 : “ + testMessage.toString());
  8.  
    }
  9.  
     
  10.  
    }

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

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

 

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

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

 

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

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


  1.  
    import org.springframework.amqp.core.Binding;
  2.  
    import org.springframework.amqp.core.BindingBuilder;
  3.  
    import org.springframework.amqp.core.Queue;
  4.  
    import org.springframework.amqp.core.TopicExchange;
  5.  
    import org.springframework.context.annotation.Bean;
  6.  
    import org.springframework.context.annotation.Configuration;
  7.  
     
  8.  
    /**
  9.  
    * @Author : JCccc
  10.  
    * @CreateTime : 2019/9/3
  11.  
    * @Description :
  12.  
    **/
  13.  
     
  14.  
    @Configuration
  15.  
    public class TopicRabbitConfig {
  16.  
    //绑定键
  17.  
    public final static String man = “topic.man”;
  18.  
    public final static String woman = “topic.woman”;
  19.  
     
  20.  
    @Bean
  21.  
    public Queue firstQueue() {
  22.  
    return new Queue(TopicRabbitConfig.man);
  23.  
    }
  24.  
     
  25.  
    @Bean
  26.  
    public Queue secondQueue() {
  27.  
    return new Queue(TopicRabbitConfig.woman);
  28.  
    }
  29.  
     
  30.  
    @Bean
  31.  
    TopicExchange exchange() {
  32.  
    return new TopicExchange(“topicExchange”);
  33.  
    }
  34.  
     
  35.  
     
  36.  
    //将firstQueue和topicExchange绑定,而且绑定的键值为topic.man
  37.  
    //这样只要是消息携带的路由键是topic.man,才会分发到该队列
  38.  
    @Bean
  39.  
    Binding bindingExchangeMessage() {
  40.  
    return BindingBuilder.bind(firstQueue()).to(exchange()).with(man);
  41.  
    }
  42.  
     
  43.  
    //将secondQueue和topicExchange绑定,而且绑定的键值为用上通配路由键规则topic.#
  44.  
    // 这样只要是消息携带的路由键是以topic.开头,都会分发到该队列
  45.  
    @Bean
  46.  
    Binding bindingExchangeMessage2() {
  47.  
    return BindingBuilder.bind(secondQueue()).to(exchange()).with(“topic.#”);
  48.  
    }
  49.  
     
  50.  
    }

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

  1.  
    @GetMapping(“/sendTopicMessage1”)
  2.  
    public String sendTopicMessage1() {
  3.  
    String messageId = String.valueOf(UUID.randomUUID());
  4.  
    String messageData = “message: M A N “;
  5.  
    String createTime = LocalDateTime.now().format(DateTimeFormatter.ofPattern(“yyyy-MM-dd HH:mm:ss”));
  6.  
    Map<String, Object> manMap = new HashMap<>();
  7.  
    manMap.put(“messageId”, messageId);
  8.  
    manMap.put(“messageData”, messageData);
  9.  
    manMap.put(“createTime”, createTime);
  10.  
    rabbitTemplate.convertAndSend(“topicExchange”, “topic.man”, manMap);
  11.  
    return “ok”;
  12.  
    }
  13.  
     
  14.  
    @GetMapping(“/sendTopicMessage2”)
  15.  
    public String sendTopicMessage2() {
  16.  
    String messageId = String.valueOf(UUID.randomUUID());
  17.  
    String messageData = “message: woman is all “;
  18.  
    String createTime = LocalDateTime.now().format(DateTimeFormatter.ofPattern(“yyyy-MM-dd HH:mm:ss”));
  19.  
    Map<String, Object> womanMap = new HashMap<>();
  20.  
    womanMap.put(“messageId”, messageId);
  21.  
    womanMap.put(“messageData”, messageData);
  22.  
    womanMap.put(“createTime”, createTime);
  23.  
    rabbitTemplate.convertAndSend(“topicExchange”, “topic.woman”, womanMap);
  24.  
    return “ok”;
  25.  
    }
  26.  
    }

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

  1.  
    import org.springframework.amqp.rabbit.annotation.RabbitHandler;
  2.  
    import org.springframework.amqp.rabbit.annotation.RabbitListener;
  3.  
    import org.springframework.stereotype.Component;
  4.  
    import java.util.Map;
  5.  
     
  6.  
    /**
  7.  
    * @Author : JCccc
  8.  
    * @CreateTime : 2019/9/3
  9.  
    * @Description :
  10.  
    **/
  11.  
    @Component
  12.  
    @RabbitListener(queues = “topic.man”)
  13.  
    public class TopicManReceiver {
  14.  
     
  15.  
    @RabbitHandler
  16.  
    public void process(Map testMessage) {
  17.  
    System.out.println(“TopicManReceiver消费者收到消息 : “ + testMessage.toString());
  18.  
    }
  19.  
    }

再创建一个TopicTotalReceiver.java:

  1.  
    package com.elegant.rabbitmqconsumer.receiver;
  2.  
     
  3.  
    import org.springframework.amqp.rabbit.annotation.RabbitHandler;
  4.  
    import org.springframework.amqp.rabbit.annotation.RabbitListener;
  5.  
    import org.springframework.stereotype.Component;
  6.  
    import java.util.Map;
  7.  
     
  8.  
    /**
  9.  
    * @Author : JCccc
  10.  
    * @CreateTime : 2019/9/3
  11.  
    * @Description :
  12.  
    **/
  13.  
     
  14.  
    @Component
  15.  
    @RabbitListener(queues = “topic.woman”)
  16.  
    public class TopicTotalReceiver {
  17.  
     
  18.  
    @RabbitHandler
  19.  
    public void process(Map testMessage) {
  20.  
    System.out.println(“TopicTotalReceiver消费者收到消息 : “ + testMessage.toString());
  21.  
    }
  22.  
    }

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

  1.  
    import org.springframework.amqp.core.Binding;
  2.  
    import org.springframework.amqp.core.BindingBuilder;
  3.  
    import org.springframework.amqp.core.Queue;
  4.  
    import org.springframework.amqp.core.TopicExchange;
  5.  
    import org.springframework.context.annotation.Bean;
  6.  
    import org.springframework.context.annotation.Configuration;
  7.  
     
  8.  
    /**
  9.  
    * @Author : JCccc
  10.  
    * @CreateTime : 2019/9/3
  11.  
    * @Description :
  12.  
    **/
  13.  
     
  14.  
    @Configuration
  15.  
    public class TopicRabbitConfig {
  16.  
    //绑定键
  17.  
    public final static String man = “topic.man”;
  18.  
    public final static String woman = “topic.woman”;
  19.  
     
  20.  
    @Bean
  21.  
    public Queue firstQueue() {
  22.  
    return new Queue(TopicRabbitConfig.man);
  23.  
    }
  24.  
     
  25.  
    @Bean
  26.  
    public Queue secondQueue() {
  27.  
    return new Queue(TopicRabbitConfig.woman);
  28.  
    }
  29.  
     
  30.  
    @Bean
  31.  
    TopicExchange exchange() {
  32.  
    return new TopicExchange(“topicExchange”);
  33.  
    }
  34.  
     
  35.  
     
  36.  
    //将firstQueue和topicExchange绑定,而且绑定的键值为topic.man
  37.  
    //这样只要是消息携带的路由键是topic.man,才会分发到该队列
  38.  
    @Bean
  39.  
    Binding bindingExchangeMessage() {
  40.  
    return BindingBuilder.bind(firstQueue()).to(exchange()).with(man);
  41.  
    }
  42.  
     
  43.  
    //将secondQueue和topicExchange绑定,而且绑定的键值为用上通配路由键规则topic.#
  44.  
    // 这样只要是消息携带的路由键是以topic.开头,都会分发到该队列
  45.  
    @Bean
  46.  
    Binding bindingExchangeMessage2() {
  47.  
    return BindingBuilder.bind(secondQueue()).to(exchange()).with(“topic.#”);
  48.  
    }
  49.  
     
  50.  
    }


然后把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:

  1.  
    import org.springframework.amqp.core.Binding;
  2.  
    import org.springframework.amqp.core.BindingBuilder;
  3.  
    import org.springframework.amqp.core.FanoutExchange;
  4.  
    import org.springframework.amqp.core.Queue;
  5.  
    import org.springframework.context.annotation.Bean;
  6.  
    import org.springframework.context.annotation.Configuration;
  7.  
    /**
  8.  
    * @Author : JCccc
  9.  
    * @CreateTime : 2019/9/3
  10.  
    * @Description :
  11.  
    **/
  12.  
     
  13.  
    @Configuration
  14.  
    public class FanoutRabbitConfig {
  15.  
     
  16.  
    /**
  17.  
    * 创建三个队列 :fanout.A fanout.B fanout.C
  18.  
    * 将三个队列都绑定在交换机 fanoutExchange 上
  19.  
    * 因为是扇型交换机, 路由键无需配置,配置也不起作用
  20.  
    */
  21.  
     
  22.  
     
  23.  
    @Bean
  24.  
    public Queue queueA() {
  25.  
    return new Queue(“fanout.A”);
  26.  
    }
  27.  
     
  28.  
    @Bean
  29.  
    public Queue queueB() {
  30.  
    return new Queue(“fanout.B”);
  31.  
    }
  32.  
     
  33.  
    @Bean
  34.  
    public Queue queueC() {
  35.  
    return new Queue(“fanout.C”);
  36.  
    }
  37.  
     
  38.  
    @Bean
  39.  
    FanoutExchange fanoutExchange() {
  40.  
    return new FanoutExchange(“fanoutExchange”);
  41.  
    }
  42.  
     
  43.  
    @Bean
  44.  
    Binding bindingExchangeA() {
  45.  
    return BindingBuilder.bind(queueA()).to(fanoutExchange());
  46.  
    }
  47.  
     
  48.  
    @Bean
  49.  
    Binding bindingExchangeB() {
  50.  
    return BindingBuilder.bind(queueB()).to(fanoutExchange());
  51.  
    }
  52.  
     
  53.  
    @Bean
  54.  
    Binding bindingExchangeC() {
  55.  
    return BindingBuilder.bind(queueC()).to(fanoutExchange());
  56.  
    }
  57.  
    }

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

  1.  
    @GetMapping(“/sendFanoutMessage”)
  2.  
    public String sendFanoutMessage() {
  3.  
    String messageId = String.valueOf(UUID.randomUUID());
  4.  
    String messageData = “message: testFanoutMessage “;
  5.  
    String createTime = LocalDateTime.now().format(DateTimeFormatter.ofPattern(“yyyy-MM-dd HH:mm:ss”));
  6.  
    Map<String, Object> map = new HashMap<>();
  7.  
    map.put(“messageId”, messageId);
  8.  
    map.put(“messageData”, messageData);
  9.  
    map.put(“createTime”, createTime);
  10.  
    rabbitTemplate.convertAndSend(“fanoutExchange”, null, map);
  11.  
    return “ok”;
  12.  
    }

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

FanoutReceiverA.java:

  1.  
    import org.springframework.amqp.rabbit.annotation.RabbitHandler;
  2.  
    import org.springframework.amqp.rabbit.annotation.RabbitListener;
  3.  
    import org.springframework.stereotype.Component;
  4.  
    import java.util.Map;
  5.  
    /**
  6.  
    * @Author : JCccc
  7.  
    * @CreateTime : 2019/9/3
  8.  
    * @Description :
  9.  
    **/
  10.  
    @Component
  11.  
    @RabbitListener(queues = “fanout.A”)
  12.  
    public class FanoutReceiverA {
  13.  
     
  14.  
    @RabbitHandler
  15.  
    public void process(Map testMessage) {
  16.  
    System.out.println(“FanoutReceiverA消费者收到消息 : “ +testMessage.toString());
  17.  
    }
  18.  
     
  19.  
    }

FanoutReceiverB.java:

  1.  
    import org.springframework.amqp.rabbit.annotation.RabbitHandler;
  2.  
    import org.springframework.amqp.rabbit.annotation.RabbitListener;
  3.  
    import org.springframework.stereotype.Component;
  4.  
    import java.util.Map;
  5.  
    /**
  6.  
    * @Author : JCccc
  7.  
    * @CreateTime : 2019/9/3
  8.  
    * @Description :
  9.  
    **/
  10.  
    @Component
  11.  
    @RabbitListener(queues = “fanout.B”)
  12.  
    public class FanoutReceiverB {
  13.  
     
  14.  
    @RabbitHandler
  15.  
    public void process(Map testMessage) {
  16.  
    System.out.println(“FanoutReceiverB消费者收到消息 : “ +testMessage.toString());
  17.  
    }
  18.  
     
  19.  
    }

FanoutReceiverC.java:

  1.  
    import org.springframework.amqp.rabbit.annotation.RabbitHandler;
  2.  
    import org.springframework.amqp.rabbit.annotation.RabbitListener;
  3.  
    import org.springframework.stereotype.Component;
  4.  
    import java.util.Map;
  5.  
     
  6.  
    /**
  7.  
    * @Author : JCccc
  8.  
    * @CreateTime : 2019/9/3
  9.  
    * @Description :
  10.  
    **/
  11.  
    @Component
  12.  
    @RabbitListener(queues = “fanout.C”)
  13.  
    public class FanoutReceiverC {
  14.  
     
  15.  
    @RabbitHandler
  16.  
    public void process(Map testMessage) {
  17.  
    System.out.println(“FanoutReceiverC消费者收到消息 : “ +testMessage.toString());
  18.  
    }
  19.  
     
  20.  
    }

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

  1.  
    import org.springframework.amqp.core.Binding;
  2.  
    import org.springframework.amqp.core.BindingBuilder;
  3.  
    import org.springframework.amqp.core.FanoutExchange;
  4.  
    import org.springframework.amqp.core.Queue;
  5.  
    import org.springframework.context.annotation.Bean;
  6.  
    import org.springframework.context.annotation.Configuration;
  7.  
    /**
  8.  
    * @Author : JCccc
  9.  
    * @CreateTime : 2019/9/3
  10.  
    * @Description :
  11.  
    **/
  12.  
    @Configuration
  13.  
    public class FanoutRabbitConfig {
  14.  
     
  15.  
    /**
  16.  
    * 创建三个队列 :fanout.A fanout.B fanout.C
  17.  
    * 将三个队列都绑定在交换机 fanoutExchange 上
  18.  
    * 因为是扇型交换机, 路由键无需配置,配置也不起作用
  19.  
    */
  20.  
     
  21.  
     
  22.  
    @Bean
  23.  
    public Queue queueA() {
  24.  
    return new Queue(“fanout.A”);
  25.  
    }
  26.  
     
  27.  
    @Bean
  28.  
    public Queue queueB() {
  29.  
    return new Queue(“fanout.B”);
  30.  
    }
  31.  
     
  32.  
    @Bean
  33.  
    public Queue queueC() {
  34.  
    return new Queue(“fanout.C”);
  35.  
    }
  36.  
     
  37.  
    @Bean
  38.  
    FanoutExchange fanoutExchange() {
  39.  
    return new FanoutExchange(“fanoutExchange”);
  40.  
    }
  41.  
     
  42.  
    @Bean
  43.  
    Binding bindingExchangeA() {
  44.  
    return BindingBuilder.bind(queueA()).to(fanoutExchange());
  45.  
    }
  46.  
     
  47.  
    @Bean
  48.  
    Binding bindingExchangeB() {
  49.  
    return BindingBuilder.bind(queueB()).to(fanoutExchange());
  50.  
    }
  51.  
     
  52.  
    @Bean
  53.  
    Binding bindingExchangeC() {
  54.  
    return BindingBuilder.bind(queueC()).to(fanoutExchange());
  55.  
    }
  56.  
    }

最后将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

  1.  
    server:
  2.  
    port: 8021
  3.  
    spring:
  4.  
    #给项目来个名字
  5.  
    application:
  6.  
    name: rabbitmq-provider
  7.  
    #配置rabbitMq 服务器
  8.  
    rabbitmq:
  9.  
    host: 127.0.0.1
  10.  
    port: 5672
  11.  
    username: root
  12.  
    password: root
  13.  
    #虚拟host 可以不设置,使用server默认host
  14.  
    virtual-host: JCcccHost
  15.  
    #消息确认配置项
  16.  
     
  17.  
    #确认消息已发送到交换机(Exchange)
  18.  
    publisher-confirms: true
  19.  
    #确认消息已发送到队列(Queue)
  20.  
    publisher-returns: true

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

  1.  
    import org.springframework.amqp.core.Message;
  2.  
    import org.springframework.amqp.rabbit.connection.ConnectionFactory;
  3.  
    import org.springframework.amqp.rabbit.connection.CorrelationData;
  4.  
    import org.springframework.amqp.rabbit.core.RabbitTemplate;
  5.  
    import org.springframework.context.annotation.Bean;
  6.  
    import org.springframework.context.annotation.Configuration;
  7.  
     
  8.  
     
  9.  
    /**
  10.  
    * @Author : JCccc
  11.  
    * @CreateTime : 2019/9/3
  12.  
    * @Description :
  13.  
    **/
  14.  
    @Configuration
  15.  
    public class RabbitConfig {
  16.  
     
  17.  
    @Bean
  18.  
    public RabbitTemplate createRabbitTemplate(ConnectionFactory connectionFactory){
  19.  
    RabbitTemplate rabbitTemplate = new RabbitTemplate();
  20.  
    rabbitTemplate.setConnectionFactory(connectionFactory);
  21.  
    //设置开启Mandatory,才能触发回调函数,无论消息推送结果怎么样都强制调用回调函数
  22.  
    rabbitTemplate.setMandatory(true);
  23.  
     
  24.  
    rabbitTemplate.setConfirmCallback(new RabbitTemplate.ConfirmCallback() {
  25.  
    @Override
  26.  
    public void confirm(CorrelationData correlationData, boolean ack, String cause) {
  27.  
    System.out.println(“ConfirmCallback: “+“相关数据:”+correlationData);
  28.  
    System.out.println(“ConfirmCallback: “+“确认情况:”+ack);
  29.  
    System.out.println(“ConfirmCallback: “+“原因:”+cause);
  30.  
    }
  31.  
    });
  32.  
     
  33.  
    rabbitTemplate.setReturnCallback(new RabbitTemplate.ReturnCallback() {
  34.  
    @Override
  35.  
    public void returnedMessage(Message message, int replyCode, String replyText, String exchange, String routingKey) {
  36.  
    System.out.println(“ReturnCallback: “+“消息:”+message);
  37.  
    System.out.println(“ReturnCallback: “+“回应码:”+replyCode);
  38.  
    System.out.println(“ReturnCallback: “+“回应信息:”+replyText);
  39.  
    System.out.println(“ReturnCallback: “+“交换机:”+exchange);
  40.  
    System.out.println(“ReturnCallback: “+“路由键:”+routingKey);
  41.  
    }
  42.  
    });
  43.  
     
  44.  
    return rabbitTemplate;
  45.  
    }
  46.  
     
  47.  
    }

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

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

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

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

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

  1.  
    @GetMapping(“/TestMessageAck”)
  2.  
    public String TestMessageAck() {
  3.  
    String messageId = String.valueOf(UUID.randomUUID());
  4.  
    String messageData = “message: non-existent-exchange test message “;
  5.  
    String createTime = LocalDateTime.now().format(DateTimeFormatter.ofPattern(“yyyy-MM-dd HH:mm:ss”));
  6.  
    Map<String, Object> map = new HashMap<>();
  7.  
    map.put(“messageId”, messageId);
  8.  
    map.put(“messageData”, messageData);
  9.  
    map.put(“createTime”, createTime);
  10.  
    rabbitTemplate.convertAndSend(“non-existent-exchange”, “TestDirectRouting”, map);
  11.  
    return “ok”;
  12.  
    }

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

  1.  
    20190904 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’, classid=60, method-id=40)
  2.  
    ConfirmCallback: 相关数据:null
  3.  
    ConfirmCallback: 确认情况:false
  4.  
    ConfirmCallback: 原因:channel error; protocol method: #method<channel.close>(reply-code=404, reply-text=NOT_FOUND – no exchange ‘non-existent-exchange’ in vhost ‘JCcccHost’, classid=60, method-id=40)

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

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

  1.  
    @Bean
  2.  
    DirectExchange lonelyDirectExchange() {
  3.  
    return new DirectExchange(“lonelyDirectExchange”);
  4.  
    }

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

  1.  
    @GetMapping(“/TestMessageAck2”)
  2.  
    public String TestMessageAck2() {
  3.  
    String messageId = String.valueOf(UUID.randomUUID());
  4.  
    String messageData = “message: lonelyDirectExchange test message “;
  5.  
    String createTime = LocalDateTime.now().format(DateTimeFormatter.ofPattern(“yyyy-MM-dd HH:mm:ss”));
  6.  
    Map<String, Object> map = new HashMap<>();
  7.  
    map.put(“messageId”, messageId);
  8.  
    map.put(“messageData”, messageData);
  9.  
    map.put(“createTime”, createTime);
  10.  
    rabbitTemplate.convertAndSend(“lonelyDirectExchange”, “TestDirectRouting”, map);
  11.  
    return “ok”;
  12.  
    }

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

  1.  
    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])
  2.  
    ReturnCallback: 回应码:312
  3.  
    ReturnCallback: 回应信息:NO_ROUTE
  4.  
    ReturnCallback: 交换机:lonelyDirectExchange
  5.  
    ReturnCallback: 路由键:TestDirectRouting
  1.  
    ConfirmCallback: 相关数据:null
  2.  
    ConfirmCallback: 确认情况:true
  3.  
    ConfirmCallback: 原因:null

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

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

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

  1.  
    ConfirmCallback: 相关数据:null
  2.  
    ConfirmCallback: 确认情况:true
  3.  
    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上添加代码相关的配置代码:

  1.  
     
  2.  
    import com.elegant.rabbitmqconsumer.receiver.MyAckReceiver;
  3.  
    import org.springframework.amqp.core.AcknowledgeMode;
  4.  
    import org.springframework.amqp.core.Queue;
  5.  
    import org.springframework.amqp.rabbit.connection.CachingConnectionFactory;
  6.  
    import org.springframework.amqp.rabbit.listener.SimpleMessageListenerContainer;
  7.  
    import org.springframework.beans.factory.annotation.Autowired;
  8.  
    import org.springframework.context.annotation.Bean;
  9.  
    import org.springframework.context.annotation.Configuration;
  10.  
     
  11.  
    /**
  12.  
    * @Author : JCccc
  13.  
    * @CreateTime : 2019/9/4
  14.  
    * @Description :
  15.  
    **/
  16.  
    @Configuration
  17.  
    public class MessageListenerConfig {
  18.  
     
  19.  
    @Autowired
  20.  
    private CachingConnectionFactory connectionFactory;
  21.  
    @Autowired
  22.  
    private MyAckReceiver myAckReceiver;//消息接收处理类
  23.  
     
  24.  
    @Bean
  25.  
    public SimpleMessageListenerContainer simpleMessageListenerContainer() {
  26.  
    SimpleMessageListenerContainer container = new SimpleMessageListenerContainer(connectionFactory);
  27.  
    container.setConcurrentConsumers(1);
  28.  
    container.setMaxConcurrentConsumers(1);
  29.  
    container.setAcknowledgeMode(AcknowledgeMode.MANUAL); // RabbitMQ默认是自动确认,这里改为手动确认消息
  30.  
    //设置一个队列
  31.  
    container.setQueueNames(“TestDirectQueue”);
  32.  
    //如果同时设置多个如下: 前提是队列都是必须已经创建存在的
  33.  
    // container.setQueueNames(“TestDirectQueue”,”TestDirectQueue2″,”TestDirectQueue3″);
  34.  
     
  35.  
     
  36.  
    //另一种设置队列的方法,如果使用这种情况,那么要设置多个,就使用addQueues
  37.  
    //container.setQueues(new Queue(“TestDirectQueue”,true));
  38.  
    //container.addQueues(new Queue(“TestDirectQueue2”,true));
  39.  
    //container.addQueues(new Queue(“TestDirectQueue3”,true));
  40.  
    container.setMessageListener(myAckReceiver);
  41.  
     
  42.  
    return container;
  43.  
    }
  44.  
     
  45.  
     
  46.  
    }

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

  1.  
    import com.rabbitmq.client.Channel;
  2.  
    import org.springframework.amqp.core.Message;
  3.  
    import org.springframework.amqp.rabbit.listener.api.ChannelAwareMessageListener;
  4.  
    import org.springframework.stereotype.Component;
  5.  
    import java.util.HashMap;
  6.  
    import java.util.Map;
  7.  
     
  8.  
    @Component
  9.  
     
  10.  
    public class MyAckReceiver implements ChannelAwareMessageListener {
  11.  
     
  12.  
    @Override
  13.  
    public void onMessage(Message message, Channel channel) throws Exception {
  14.  
    long deliveryTag = message.getMessageProperties().getDeliveryTag();
  15.  
    try {
  16.  
    //因为传递消息的时候用的map传递,所以将Map从Message内取出需要做些处理
  17.  
    String msg = message.toString();
  18.  
    String[] msgArray = msg.split(“‘”);//可以点进Message里面看源码,单引号直接的数据就是我们的map消息数据
  19.  
    Map<String, String> msgMap = mapStringToMap(msgArray[1].trim(),3);
  20.  
    String messageId=msgMap.get(“messageId”);
  21.  
    String messageData=msgMap.get(“messageData”);
  22.  
    String createTime=msgMap.get(“createTime”);
  23.  
    System.out.println(” MyAckReceiver messageId:”+messageId+” messageData:”+messageData+” createTime:”+createTime);
  24.  
    System.out.println(“消费的主题消息来自:”+message.getMessageProperties().getConsumerQueue());
  25.  
    channel.basicAck(deliveryTag, true); //第二个参数,手动确认可以被批处理,当该参数为 true 时,则可以一次性确认 delivery_tag 小于等于传入值的所有消息
  26.  
    // channel.basicReject(deliveryTag, true);//第二个参数,true会重新放回队列,所以需要自己根据业务逻辑判断什么时候使用拒绝
  27.  
    } catch (Exception e) {
  28.  
    channel.basicReject(deliveryTag, false);
  29.  
    e.printStackTrace();
  30.  
    }
  31.  
    }
  32.  
     
  33.  
    //{key=value,key=value,key=value} 格式转换成map
  34.  
    private Map<String, String> mapStringToMap(String str,int entryNum ) {
  35.  
    str = str.substring(1, str.length() – 1);
  36.  
    String[] strs = str.split(“,”,entryNum);
  37.  
    Map<String, String> map = new HashMap<String, String>();
  38.  
    for (String string : strs) {
  39.  
    String key = string.split(“=”)[0].trim();
  40.  
    String value = string.split(“=”)[1];
  41.  
    map.put(key, value);
  42.  
    }
  43.  
    return map;
  44.  
    }
  45.  
    }

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

 

 

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

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

没有问题,接下来看代码

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

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

 

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

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

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

  1.  
    import com.rabbitmq.client.Channel;
  2.  
    import org.springframework.amqp.core.Message;
  3.  
    import org.springframework.amqp.rabbit.listener.api.ChannelAwareMessageListener;
  4.  
    import org.springframework.stereotype.Component;
  5.  
    import java.util.HashMap;
  6.  
    import java.util.Map;
  7.  
     
  8.  
    @Component
  9.  
    public class MyAckReceiver implements ChannelAwareMessageListener {
  10.  
     
  11.  
    @Override
  12.  
    public void onMessage(Message message, Channel channel) throws Exception {
  13.  
    long deliveryTag = message.getMessageProperties().getDeliveryTag();
  14.  
    try {
  15.  
    //因为传递消息的时候用的map传递,所以将Map从Message内取出需要做些处理
  16.  
    String msg = message.toString();
  17.  
    String[] msgArray = msg.split(“‘”);//可以点进Message里面看源码,单引号直接的数据就是我们的map消息数据
  18.  
    Map<String, String> msgMap = mapStringToMap(msgArray[1].trim(),3);
  19.  
    String messageId=msgMap.get(“messageId”);
  20.  
    String messageData=msgMap.get(“messageData”);
  21.  
    String createTime=msgMap.get(“createTime”);
  22.  
     
  23.  
    if (“TestDirectQueue”.equals(message.getMessageProperties().getConsumerQueue())){
  24.  
    System.out.println(“消费的消息来自的队列名为:”+message.getMessageProperties().getConsumerQueue());
  25.  
    System.out.println(“消息成功消费到 messageId:”+messageId+” messageData:”+messageData+” createTime:”+createTime);
  26.  
    System.out.println(“执行TestDirectQueue中的消息的业务处理流程……”);
  27.  
     
  28.  
    }
  29.  
     
  30.  
    if (“fanout.A”.equals(message.getMessageProperties().getConsumerQueue())){
  31.  
    System.out.println(“消费的消息来自的队列名为:”+message.getMessageProperties().getConsumerQueue());
  32.  
    System.out.println(“消息成功消费到 messageId:”+messageId+” messageData:”+messageData+” createTime:”+createTime);
  33.  
    System.out.println(“执行fanout.A中的消息的业务处理流程……”);
  34.  
     
  35.  
    }
  36.  
     
  37.  
    channel.basicAck(deliveryTag, true);
  38.  
    // channel.basicReject(deliveryTag, true);//为true会重新放回队列
  39.  
    } catch (Exception e) {
  40.  
    channel.basicReject(deliveryTag, false);
  41.  
    e.printStackTrace();
  42.  
    }
  43.  
    }
  44.  
     
  45.  
    //{key=value,key=value,key=value} 格式转换成map
  46.  
    private Map<String, String> mapStringToMap(String str,int enNum) {
  47.  
    str = str.substring(1, str.length() – 1);
  48.  
    String[] strs = str.split(“,”,enNum);
  49.  
    Map<String, String> map = new HashMap<String, String>();
  50.  
    for (String string : strs) {
  51.  
    String key = string.split(“=”)[0].trim();
  52.  
    String value = string.split(“=”)[1];
  53.  
    map.put(key, value);
  54.  
    }
  55.  
    return map;
  56.  
    }
  57.  
    }

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

调用接口/sendDirectMessage  和 /sendFanoutMessage ,

 

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

 

 

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