springboot整合rabbitmq实现生产者消息确认、死信交换器、未路由到队列的消息

我不是女神ヾ 2024-02-19 18:27 56阅读 0赞

在上篇文章 springboot 整合 rabbitmq 中,我们实现了springboot 和rabbitmq的简单整合,这篇文章主要是对上篇文章功能的增强,主要完成如下功能。

需求:

生产者在启动的时候,自动创建好队列、绑定、交换器并设置好 死信交换器、备份交换器(alternate-exchange)。生产者发送消息后,生产者这边需要对发送的消息进行确认,确认RabbitMQ接收到了消息。为了测试未被路由的消息和死信消息,发送方,发送11条正常的,可以被路由到消息队列中的消息,发送一条不可路由到消息队列中的消息,使之进入 alternate-exchange 交换器中。接收方在接收到消息后,随机拒绝一些消息,使之进入 x-dead-letter-exchange 中。

实现如下功能:

1、使用@Bean方式自动实现队列、交换器、绑定的创建。
2、使用@RabbitListener实现队列消息的监听。
3、实现生产者消息确认。
4、实现死信交换器(过期的消息、basic.nack或basic.reject且requeue参数为false或队列满的消息将进入此交换器)。
5、实现备份交换器(alternate-exchange),未被正确路由的消息将会经过此交换器。

部分功能实现要点:

1、生产者消息确认

|- spring.rabbitmq.template.mandatory = true 设置成true

|- spring.rabbitmq.publisher-confirms = true 设置成true

|- 编写一个 java 类,实现 RabbitTemplate.ConfirmCallback 接口,在这个里面我们可以确认消息是否到达了RabbitMQ服务器。

2、死信交换器的实现

|- 申明队列的时候设置 x-dead-letter-exchange 参数

3、处理未被路由的消息

|- 申明交换器的时候设置 alternate-exchange 参数

实现步骤如下:

1、jar包的引入,生产者和消费者都一样

  1. <dependencies>
  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>
  10. <dependency>
  11. <groupId>org.projectlombok</groupId>
  12. <artifactId>lombok</artifactId>
  13. <optional>true</optional>
  14. </dependency>
  15. <dependency>
  16. <groupId>org.springframework.boot</groupId>
  17. <artifactId>spring-boot-starter-test</artifactId>
  18. <scope>test</scope>
  19. </dependency>
  20. </dependencies>

2、生产者 - 配置文件

  1. server:
  2. port: 9088
  3. spring:
  4. rabbitmq:
  5. host: 140.143.237.224
  6. port: 5672
  7. username: root
  8. password: root
  9. virtual-host: /
  10. connection-timeout: 10000
  11. template:
  12. mandatory: true
  13. publisher-confirms: true

注意:**此处需要将 mandatory和publisher-confirms参数设置成true**

3、生产者 - 生产者消息确认编写

  1. @Slf4j
  2. public class RabbitConfirmCallback implements RabbitTemplate.ConfirmCallback {
  3. @Override
  4. public void confirm(CorrelationData correlationData, boolean ack, String cause) {
  5. log.info("(start)生产者消息确认=========================");
  6. log.info("correlationData:[{}]", correlationData);
  7. log.info("ack:[{}]", ack);
  8. log.info("cause:[{}]", cause);
  9. if (!ack) {
  10. log.info("消息可能未到达rabbitmq服务器");
  11. }
  12. log.info("(end)生产者消息确认=========================");
  13. }
  14. }

注意:**此处需要实现 RabbitTemplate.ConfirmCallback接口,实现消息的确认**

4、生产者 - 生产者配置

  1. @Configuration
  2. public class RabbitmqConfiguration {
  3. @Autowired
  4. private RabbitTemplate rabbitTemplate;
  5. @PostConstruct
  6. public void initRabbitTemplate() {
  7. // 设置生产者消息确认
  8. rabbitTemplate.setConfirmCallback(new RabbitConfirmCallback());
  9. }
  10. /**
  11. * 申明队列
  12. *
  13. * @return
  14. */
  15. @Bean
  16. public Queue queue() {
  17. Map<String, Object> arguments = new HashMap<>(4);
  18. // 申明死信交换器
  19. arguments.put("x-dead-letter-exchange", "exchange-dlx");
  20. return new Queue("queue-rabbit-springboot-advance", true, false, false, arguments);
  21. }
  22. /**
  23. * 没有路由到的消息将进入此队列
  24. *
  25. * @return
  26. */
  27. @Bean
  28. public Queue unRouteQueue() {
  29. return new Queue("queue-unroute");
  30. }
  31. /**
  32. * 死信队列
  33. *
  34. * @return
  35. */
  36. @Bean
  37. public Queue dlxQueue() {
  38. return new Queue("dlx-queue");
  39. }
  40. /**
  41. * 申明交换器
  42. *
  43. * @return
  44. */
  45. @Bean
  46. public Exchange exchange() {
  47. Map<String, Object> arguments = new HashMap<>(4);
  48. // 当发往exchange-rabbit-springboot-advance的消息,routingKey和bindingKey没有匹配上时,将会由exchange-unroute交换器进行处理
  49. arguments.put("alternate-exchange", "exchange-unroute");
  50. return new DirectExchange("exchange-rabbit-springboot-advance", true, false, arguments);
  51. }
  52. @Bean
  53. public FanoutExchange unRouteExchange() {
  54. // 此处的交换器的名字要和 exchange() 方法中 alternate-exchange 参数的值一致
  55. return new FanoutExchange("exchange-unroute");
  56. }
  57. /**
  58. * 申明死信交换器
  59. *
  60. * @return
  61. */
  62. @Bean
  63. public FanoutExchange dlxExchange() {
  64. return new FanoutExchange("exchange-dlx");
  65. }
  66. /**
  67. * 申明绑定
  68. *
  69. * @return
  70. */
  71. @Bean
  72. public Binding binding() {
  73. return BindingBuilder.bind(queue()).to(exchange()).with("product").noargs();
  74. }
  75. @Bean
  76. public Binding unRouteBinding() {
  77. return BindingBuilder.bind(unRouteQueue()).to(unRouteExchange());
  78. }
  79. @Bean
  80. public Binding dlxBinding() {
  81. return BindingBuilder.bind(dlxQueue()).to(dlxExchange());
  82. }
  83. }

注意:**x-dead-letter-exchange和alternate-exchange参数的值和交换器中的值需要保持一致**

5、生产者 - 编写消息发送者

  1. @Component
  2. @Slf4j
  3. public class RabbitProducer implements ApplicationListener<ContextRefreshedEvent> {
  4. @Autowired
  5. private RabbitTemplate rabbitTemplate;
  6. @Override
  7. public void onApplicationEvent(ContextRefreshedEvent event) {
  8. String exchange = "exchange-rabbit-springboot-advance";
  9. String routingKey = "product";
  10. String unRoutingKey = "norProduct";
  11. // 1.发送一条正常的消息 CorrelationData唯一(可以在ConfirmListener中确认消息)
  12. IntStream.rangeClosed(0, 10).forEach(num -> {
  13. String message = LocalDateTime.now().toString() + "发送第" + (num + 1) + "条消息.";
  14. rabbitTemplate.convertAndSend(exchange, routingKey, message, new CorrelationData("routing" + UUID.randomUUID().toString()));
  15. log.info("发送一条消息,exchange:[{}],routingKey:[{}],message:[{}]", exchange, routingKey, message);
  16. });
  17. // 2.发送一条未被路由的消息,此消息将会进入备份交换器(alternate exchange)
  18. String message = LocalDateTime.now().toString() + "发送一条消息.";
  19. rabbitTemplate.convertAndSend(exchange, unRoutingKey, message, new CorrelationData("unRouting-" + UUID.randomUUID().toString()));
  20. log.info("发送一条消息,exchange:[{}],routingKey:[{}],message:[{}]", exchange, unRoutingKey, message);
  21. }
  22. }

注意:1、**此处发送了2中消息,一种消息可以被正确的路由到消息队列中,另一种由于routingKey是不存在的,因此不会路由到队列中,观察这条消息有没有路由到 alternate-exchange 绑定的队列中。**

2、CorrelationData数据需要唯一,此值可用于生产者确认消息。

6、生产者 - 启动类

  1. @SpringBootApplication
  2. public class ProducerApplication {
  3. public static void main(String[] args) {
  4. SpringApplication.run(ProducerApplication.class, args);
  5. }
  6. }

7、消费者 - 配置文件

  1. server:
  2. port: 9087
  3. spring:
  4. rabbitmq:
  5. host: 140.143.237.224
  6. port: 5672
  7. username: root
  8. password: root
  9. virtual-host: /
  10. connection-timeout: 10000
  11. listener:
  12. simple:
  13. acknowledge-mode: manual # 手动应答
  14. auto-startup: true
  15. default-requeue-rejected: false # 不重回队列
  16. concurrency: 5
  17. max-concurrency: 20
  18. prefetch: 1 # 每次只处理一个信息
  19. retry:
  20. enabled: true

8、消费者 - 消息接收

  1. @Component
  2. @Slf4j
  3. public class RabbitConsumer {
  4. /**
  5. * 监听 queue-rabbit-springboot-advance 队列
  6. *
  7. * @param receiveMessage 接收到的消息
  8. * @param message
  9. * @param channel
  10. */
  11. @RabbitListener(queues = "queue-rabbit-springboot-advance")
  12. public void receiveMessage(String receiveMessage, Message message, Channel channel) {
  13. try {
  14. // 手动签收
  15. log.info("接收到消息:[{}]", receiveMessage);
  16. if (new Random().nextInt(10) < 5) {
  17. log.warn("拒绝一条信息:[{}],此消息将会由死信交换器进行路由.", receiveMessage);
  18. channel.basicNack(message.getMessageProperties().getDeliveryTag(), false, false);
  19. } else {
  20. channel.basicAck(message.getMessageProperties().getDeliveryTag(), false);
  21. }
  22. } catch (Exception e) {
  23. log.info("接收到消息之后的处理发生异常.", e);
  24. try {
  25. channel.basicAck(message.getMessageProperties().getDeliveryTag(), false);
  26. } catch (IOException e1) {
  27. log.error("签收异常.", e1);
  28. }
  29. }
  30. }
  31. }

注意:**消费者接收中会随机拒绝几条消息,观察这个消息有没有进入 x-dead-letter-exchange 交换器绑定的队列中。**

9、执行结果aHR0cDovL2RsMi5pdGV5ZS5jb20vdXBsb2FkL2F0dGFjaG1lbnQvMDEzMS8wMDI3L2FmMGZlMjA5LTE1OWYtMzRiOS1hODcwLTkyZDlhNmUxMjNkNS5wbmc

完整代码:

代码如下:https://gitee.com/huan1993/rabbitmq/tree/master/rabbitmq-springboot-advanced

发表评论

表情:
评论列表 (有 0 条评论,56人围观)

还没有评论,来说两句吧...

相关阅读