Springboot整合RabbitMQ基本使用

[复制链接]
发表于 2022-11-19 15:58:59 | 显示全部楼层 |阅读模式

马上注册,结交更多好友,享用更多功能,让你轻松玩转社区。

您需要 登录 才可以下载或查看,没有账号?立即注册

×
1、依赖
  1.         <dependency>
  2.                         <groupId>org.springframework.boot</groupId>
  3.                         <artifactId>spring-boot-starter-amqp</artifactId>
  4.                 </dependency>
复制代码
2、rabbitmq链接配置
  1. spring:
  2.   rabbitmq:
  3.     host: 127.0.0.1
  4.     port: 5672
  5.     username: wq
  6.     password: qifeng
  7.     virtual-host: /
  8.     #开启ack
  9.     listener:
  10.       direct:
  11.         acknowledge-mode: manual
  12.         prefetch: 1 # 限制一次拉取消息的数量
  13.       simple:
  14.         acknowledge-mode: manual #采取手动应答
  15.         #concurrency: 1 # 指定最小的消费者数量
  16.         #max-concurrency: 1 #指定最大的消费者数量
  17.         retry:
  18.           enabled: true # 是否支持重试
复制代码
3、生产者

  1. package com.wanqi.mq;
  2. import org.springframework.amqp.core.*;
  3. import org.springframework.beans.factory.annotation.Qualifier;
  4. import org.springframework.context.annotation.Bean;
  5. import org.springframework.context.annotation.Configuration;
  6. /**
  7. * @Description TODO
  8. * @Version 1.0.0
  9. * @Date 2022/11/19
  10. * @Author wandaren
  11. */
  12. @Configuration
  13. public class RabbitMQConfig {
  14.     public static final String EXCHANGE_NAME = "boot_topic_exchange";
  15.     public static final String QUEUE_NAME = "boot_queue";
  16.     //交换机
  17.     @Bean("bootExchange")
  18.     public Exchange bootExchange(){
  19.         return ExchangeBuilder.topicExchange(EXCHANGE_NAME).durable(true)
  20.                 .build();
  21.     }
  22.     //
  23.     @Bean("bootQueue")
  24.     public Queue bootQueue(){
  25.         return QueueBuilder.durable(QUEUE_NAME).build();
  26.     }
  27.     //队列和交换机绑定关系
  28.     /*
  29.         知道哪个队列
  30.         知道哪个交换机
  31.         知道routing key
  32.      */
  33.     @Bean
  34.     public Binding bindingQueueExchange(@Qualifier("bootQueue") Queue queue,
  35.                                         @Qualifier("bootExchange") Exchange exchange){
  36.         return BindingBuilder.bind(queue).to(exchange).with("boot.#").noargs();
  37.     }
  38. }
复制代码

  • 发送消息
  1. package com.wanqi;
  2. import com.wanqi.mq.RabbitMQConfig;
  3. import org.junit.jupiter.api.Test;
  4. import org.springframework.amqp.rabbit.core.RabbitTemplate;
  5. import org.springframework.beans.factory.annotation.Autowired;
  6. import org.springframework.boot.test.context.SpringBootTest;
  7. @SpringBootTest
  8. class SpringbootMqProducersApplicationTests {
  9.         @Autowired
  10.         private RabbitTemplate rabbitTemplate;
  11.         @Test
  12.         void contextLoads() {
  13.                 rabbitTemplate.convertAndSend(RabbitMQConfig.EXCHANGE_NAME, "boot.test", "boot mq hello~~~");
  14.         }
  15. }
复制代码
4、消费者
  1. package com.wanqi.listener;
  2. import com.rabbitmq.client.Channel;
  3. import org.springframework.amqp.core.Message;
  4. import org.springframework.amqp.rabbit.annotation.RabbitListener;
  5. import org.springframework.stereotype.Component;
  6. import java.io.IOException;
  7. /**
  8. * @Description TODO
  9. * @Version 1.0.0
  10. * @Date 2022/11/19
  11. * @Author wandaren
  12. */
  13. @Component
  14. public class RabbitMQListener {
  15.     @RabbitListener(queues = {"boot_queue"})
  16.     public void listenerQueue(Object msg, Message message, Channel channel) {
  17.         final long deliveryTag = message.getMessageProperties().getDeliveryTag();
  18.         try {
  19.             System.out.println(msg.toString());
  20.             System.out.println(new String(message.getBody()));
  21. //            int x = 3/0;
  22.             channel.basicAck(deliveryTag, true);
  23.         } catch (Exception e) {
  24.             try {
  25.                 channel.basicNack(deliveryTag, true, true);
  26.             } catch (IOException ex) {
  27.                 throw new RuntimeException(ex);
  28.             }
  29.         }
  30.     }
  31. }
复制代码
5、死信队列

消息成为死信的三种情况
1、队列消息长度到达限制;
2、消费者拒接消费消息,并且不重回队列;
3、原队列存在消息过期设置,消息到达超时时间未被消费;
5.1、声明交换机、队列、绑定队列


  • topic模式
  1. package com.wanqi.mq;
  2. import org.springframework.amqp.core.*;
  3. import org.springframework.beans.factory.annotation.Qualifier;
  4. import org.springframework.context.annotation.Bean;
  5. import org.springframework.context.annotation.Configuration;
  6. /**
  7. * @Description TODO
  8. * @Version 1.0.0
  9. * @Date 2022/11/19
  10. * @Author wandaren
  11. */
  12. @Configuration
  13. public class QDRabbitMQConfig {
  14.     //交换机名称
  15.     public static final String ITEM_EXCHANGE = "item_exchange";
  16.     public static final String DEAD_EXCHANGE = "dead_exchange";
  17.     //队列名称
  18.     public static final String ITEM_QUEUE = "item_queue";
  19.     public static final String DEAD_QUEUE = "dead_queue";
  20.     //声明业务交换机
  21.     @Bean("itemExchange")
  22.     public Exchange itemExchange(){
  23.         return ExchangeBuilder.topicExchange(ITEM_EXCHANGE).durable(true).build();
  24.     }
  25.     //声明死信交换机
  26.     @Bean("deadExchange")
  27.     public Exchange deadExchange(){
  28.         return ExchangeBuilder.topicExchange(DEAD_EXCHANGE).durable(true).build();
  29.     }
  30.     /**
  31.      * 声明普通队列,设置队列消息过期时间,队列长度,绑定的死信队列
  32.      */
  33.     @Bean("itemQueue")
  34.     public Queue itemQueue(){
  35.         return QueueBuilder
  36.                 .durable(ITEM_QUEUE)
  37.                 // 队列消息过期时间
  38.                 .ttl(10000)
  39.                 // 队列长度
  40.                 .maxLength(15)
  41.                 // 声明当前队列绑定的死信交换机
  42.                 .deadLetterExchange(DEAD_EXCHANGE)
  43.                 // 声明当前队列死信转发的路由key
  44.                 .deadLetterRoutingKey("infoDead.haha")
  45.                 .build();
  46.     }
  47.     /**
  48.      * 死信队列,消费者需要监听的队列
  49.      */
  50.     @Bean("deadQueue")
  51.     public Queue deadQueue(){
  52.         return QueueBuilder
  53.                 .durable(DEAD_QUEUE)
  54.                 .build();
  55.     }
  56.     //绑定队列和交换机(业务)
  57.     @Bean
  58.     public Binding itemQueueExchange(@Qualifier("itemQueue") Queue queue,
  59.                                      @Qualifier("itemExchange") Exchange exchange){
  60.         return BindingBuilder.bind(queue).to(exchange).with("infoRouting.#").noargs();
  61.     }
  62.     //绑定队列和交换机(死信)
  63.     @Bean
  64.     public Binding deadQueueExchange(@Qualifier("deadQueue") Queue queue,
  65.                                      @Qualifier("deadExchange") Exchange exchange){
  66.         return BindingBuilder.bind(queue).to(exchange).with("infoDead.#").noargs();
  67.     }
  68. }
复制代码
5.2、发送测试消息
  1. package com.wanqi;
  2. import com.wanqi.mq.RabbitMQConfig;
  3. import com.wanqi.mq.RabbitMQConfig2;
  4. import org.junit.jupiter.api.Test;
  5. import org.springframework.amqp.core.Message;
  6. import org.springframework.amqp.core.MessageProperties;
  7. import org.springframework.amqp.rabbit.core.RabbitTemplate;
  8. import org.springframework.beans.factory.annotation.Autowired;
  9. import org.springframework.boot.test.context.SpringBootTest;
  10. @SpringBootTest
  11. class SpringbootMqProducersApplicationTests {
  12.         @Autowired
  13.         private RabbitTemplate rabbitTemplate;
  14.    
  15.         // 原队列存在消息过期设置,消息到达超时时间未被消费;
  16.         @Test
  17.         void contextLoads3() {
  18.                 MessageProperties messageProperties = new MessageProperties();
  19.                 // 设置过期时间,单位:毫秒
  20.                 messageProperties.setExpiration("5000");
  21.                 byte[] msgBytes = "rabbitmq ttl message ...".getBytes();
  22.                 Message message = new Message(msgBytes, messageProperties);
  23.                 //发送消息
  24.                 rabbitTemplate.convertAndSend(QDRabbitMQConfig.ITEM_EXCHANGE,"infoRouting.hehe",message);
  25.                 System.out.println("发送消息成功");
  26.         }
  27.     // 模拟队列消息长度到达限制
  28.         @Test
  29.         void contextLoads4() {
  30.                 for (int i = 0; i < 20; i++) {
  31.                         rabbitTemplate.convertAndSend(QDRabbitMQConfig.ITEM_EXCHANGE,"infoRouting.hehe",i + "---message");
  32.                 }
  33.         }
  34. }
复制代码
5.3、消费者
  1. package com.wanqi.listener;
  2. import com.rabbitmq.client.Channel;
  3. import org.springframework.amqp.core.Message;
  4. import org.springframework.amqp.rabbit.annotation.RabbitListener;
  5. import org.springframework.stereotype.Component;
  6. import java.io.IOException;
  7. /**
  8. * @Description TODO
  9. * @Version 1.0.0
  10. * @Date 2022/11/19
  11. * @Author wandaren
  12. */
  13. @Component
  14. public class RabbitMQListener {
  15.     @RabbitListener(queues = {"dead_queue"})
  16.     public void listenerQueue2(Object msg, Message message, Channel channel) {
  17.         final long deliveryTag = message.getMessageProperties().getDeliveryTag();
  18.         try {
  19.             System.out.println(msg.toString());
  20.             channel.basicAck(deliveryTag, true);
  21.         } catch (Exception e) {
  22.             try {
  23.                 channel.basicNack(deliveryTag, true, true);
  24.             } catch (IOException ex) {
  25.                 throw new RuntimeException(ex);
  26.             }
  27.         }
  28.     }
  29. }
复制代码
6、延迟队列(TTL + 死信队列)

6.1、声明交换机、队列、绑定队列


  • direct模式
  1. package com.wanqi.mq;
  2. import org.springframework.amqp.core.*;
  3. import org.springframework.beans.factory.annotation.Qualifier;
  4. import org.springframework.context.annotation.Bean;
  5. import org.springframework.context.annotation.Configuration;
  6. /**
  7. * @Description TODO
  8. * @Version 1.0.0
  9. * @Date 2022/11/19
  10. * @Author wandaren
  11. */
  12. @Configuration
  13. public class TtlQueueConfig {
  14.     //普通交换机名称
  15.     public static final String X_CHANGE = "X_Exchange";
  16.     //死信交换机名称
  17.     public static final String Y_DEAD_CHANGE = "Y_Exchange";
  18.     //普通队列
  19.     public static final String QUEUE_A = "QA_QUEUE";
  20.     public static final String QUEUE_B = "QB_QUEUE";
  21.     //死信队列
  22.     public static final String DEAD_QUEUE_D = "QD_QUEUE";
  23.     //声明普通交换机
  24.     @Bean("xExchange")
  25.     public DirectExchange xExchange() {
  26.         return ExchangeBuilder.directExchange(X_CHANGE).durable(true)
  27.                 .build();
  28.     }
  29.     //声明死信交换机
  30.     @Bean("yExchange")
  31.     public DirectExchange yExchange() {
  32.         return ExchangeBuilder.directExchange(Y_DEAD_CHANGE).durable(true)
  33.                 .build();
  34.     }
  35.     /**
  36.      * 声明队列,延迟10秒
  37.      */
  38.     @Bean("queueA")
  39.     public Queue queueA() {
  40.         return QueueBuilder.durable(QUEUE_A)
  41.                 .deadLetterExchange(Y_DEAD_CHANGE) //死信交换机
  42.                 .deadLetterRoutingKey("YD")  //死信RoutingKey
  43.                 .ttl(10000)  //消息过期时间
  44.                 .build();
  45.     }
  46.     /**
  47.      * 声明队列,延迟40秒
  48.      *
  49.      * @return
  50.      */
  51.     @Bean("queueB")
  52.     public Queue queueB() {
  53.         return QueueBuilder.durable(QUEUE_B)
  54.                 .deadLetterExchange(Y_DEAD_CHANGE) //死信交换机
  55.                 .deadLetterRoutingKey("YD")  //死信RoutingKey
  56.                 .ttl(40000)  //消息过期时间
  57.                 .build();
  58.     }
  59.     /**
  60.      * 死信队列,消费者需要监听的队列
  61.      */
  62.     @Bean("queueD")
  63.     public Queue queueD() {
  64.         return QueueBuilder.durable(DEAD_QUEUE_D).build();
  65.     }
  66.     //绑定  X_CHANGE绑定queueA
  67.     @Bean
  68.     public Binding queueABindingX(@Qualifier("queueA") Queue queueA, @Qualifier("xExchange") DirectExchange xExchange) {
  69.         return BindingBuilder.bind(queueA).to(xExchange).with("XA");
  70.     }
  71.     //绑定  X_CHANGE绑定queueB
  72.     @Bean
  73.     public Binding queueBBindingX(@Qualifier("queueB") Queue queueB, @Qualifier("xExchange") DirectExchange xExchange) {
  74.         return BindingBuilder.bind(queueB).to(xExchange).with("XB");
  75.     }
  76.     //绑定  Y_CHANGE绑定queueD
  77.     @Bean
  78.     public Binding queueDBindingY(@Qualifier("queueD") Queue queueD, @Qualifier("yExchange") DirectExchange yExchange) {
  79.         return BindingBuilder.bind(queueD).to(yExchange).with("YD");
  80.     }
  81. }
复制代码
6.2、发送消息
  1.     @Test
  2.     void contextLoads5() {
  3.         String message = "延迟队列消息";
  4.         rabbitTemplate.convertAndSend(TtlQueueConfig.X_CHANGE, "XA", "TTL=10s的队列:" + message);
  5.         rabbitTemplate.convertAndSend(TtlQueueConfig.X_CHANGE, "XB", "TTL=40s的队列:" + message);
  6.     }
复制代码
6.3、消费
  1. package com.wanqi.listener;
  2. import com.rabbitmq.client.Channel;
  3. import org.springframework.amqp.core.Message;
  4. import org.springframework.amqp.rabbit.annotation.RabbitListener;
  5. import org.springframework.stereotype.Component;
  6. import java.io.IOException;
  7. import java.util.Date;
  8. /**
  9. * @Description TODO
  10. * @Version 1.0.0
  11. * @Date 2022/11/19
  12. * @Author wandaren
  13. */
  14. @Component
  15. public class RabbitMQListener {
  16.     @RabbitListener(queues = "QD_QUEUE")
  17.     public void receiveMessage(Object msg,Message message, Channel channel){
  18.         final long deliveryTag = message.getMessageProperties().getDeliveryTag();
  19.         try {
  20.             System.out.println(msg.toString());
  21.             System.out.print("当前时间: " + new Date());
  22.             System.out.println("  收到死信队列的消息:" + msg);
  23.             channel.basicAck(deliveryTag, true);
  24.         } catch (Exception e) {
  25.             try {
  26.                 channel.basicNack(deliveryTag, true, true);
  27.             } catch (IOException ex) {
  28.                 throw new RuntimeException(ex);
  29.             }
  30.         }
  31.     }
  32. }
复制代码
回复

使用道具 举报

登录后关闭弹窗

登录参与点评抽奖  加入IT实名职场社区
去登录
快速回复 返回顶部 返回列表