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

本文主要是介绍springboot整合rabbitmq实现生产者消息确认、死信交换器、未路由到队列的消息,希望对大家解决编程问题提供一定的参考价值,需要的开发者们随着小编来一起学习吧!

    在上篇文章  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包的引入,生产者和消费者都一样

<dependencies><dependency><groupId>org.springframework.boot</groupId><artifactId>spring-boot-starter-amqp</artifactId></dependency><dependency><groupId>org.springframework.boot</groupId><artifactId>spring-boot-starter-web</artifactId></dependency><dependency><groupId>org.projectlombok</groupId><artifactId>lombok</artifactId><optional>true</optional></dependency><dependency><groupId>org.springframework.boot</groupId><artifactId>spring-boot-starter-test</artifactId><scope>test</scope></dependency></dependencies>

 

 2、生产者 - 配置文件

server:port: 9088
spring:rabbitmq:host: 140.143.237.224port: 5672username: rootpassword: rootvirtual-host: /connection-timeout: 10000template:mandatory: truepublisher-confirms: true

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

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

@Slf4j
public class RabbitConfirmCallback implements RabbitTemplate.ConfirmCallback {@Overridepublic void confirm(CorrelationData correlationData, boolean ack, String cause) {log.info("(start)生产者消息确认=========================");log.info("correlationData:[{}]", correlationData);log.info("ack:[{}]", ack);log.info("cause:[{}]", cause);if (!ack) {log.info("消息可能未到达rabbitmq服务器");}log.info("(end)生产者消息确认=========================");}
}

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

  4、生产者 - 生产者配置

@Configuration
public class RabbitmqConfiguration {@Autowiredprivate RabbitTemplate rabbitTemplate;@PostConstructpublic void initRabbitTemplate() {// 设置生产者消息确认rabbitTemplate.setConfirmCallback(new RabbitConfirmCallback());}/*** 申明队列** @return*/@Beanpublic Queue queue() {Map<String, Object> arguments = new HashMap<>(4);// 申明死信交换器arguments.put("x-dead-letter-exchange", "exchange-dlx");return new Queue("queue-rabbit-springboot-advance", true, false, false, arguments);}/*** 没有路由到的消息将进入此队列** @return*/@Beanpublic Queue unRouteQueue() {return new Queue("queue-unroute");}/*** 死信队列** @return*/@Beanpublic Queue dlxQueue() {return new Queue("dlx-queue");}/*** 申明交换器** @return*/@Beanpublic Exchange exchange() {Map<String, Object> arguments = new HashMap<>(4);// 当发往exchange-rabbit-springboot-advance的消息,routingKey和bindingKey没有匹配上时,将会由exchange-unroute交换器进行处理arguments.put("alternate-exchange", "exchange-unroute");return new DirectExchange("exchange-rabbit-springboot-advance", true, false, arguments);}@Beanpublic FanoutExchange unRouteExchange() {// 此处的交换器的名字要和 exchange() 方法中 alternate-exchange 参数的值一致return new FanoutExchange("exchange-unroute");}/*** 申明死信交换器** @return*/@Beanpublic FanoutExchange dlxExchange() {return new FanoutExchange("exchange-dlx");}/*** 申明绑定** @return*/@Beanpublic Binding binding() {return BindingBuilder.bind(queue()).to(exchange()).with("product").noargs();}@Beanpublic Binding unRouteBinding() {return BindingBuilder.bind(unRouteQueue()).to(unRouteExchange());}@Beanpublic Binding dlxBinding() {return BindingBuilder.bind(dlxQueue()).to(dlxExchange());}
}

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

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

@Component
@Slf4j
public class RabbitProducer implements ApplicationListener<ContextRefreshedEvent> {@Autowiredprivate RabbitTemplate rabbitTemplate;@Overridepublic void onApplicationEvent(ContextRefreshedEvent event) {String exchange = "exchange-rabbit-springboot-advance";String routingKey = "product";String unRoutingKey = "norProduct";// 1.发送一条正常的消息 CorrelationData唯一(可以在ConfirmListener中确认消息)IntStream.rangeClosed(0, 10).forEach(num -> {String message = LocalDateTime.now().toString() + "发送第" + (num + 1) + "条消息.";rabbitTemplate.convertAndSend(exchange, routingKey, message, new CorrelationData("routing" + UUID.randomUUID().toString()));log.info("发送一条消息,exchange:[{}],routingKey:[{}],message:[{}]", exchange, routingKey, message);});// 2.发送一条未被路由的消息,此消息将会进入备份交换器(alternate exchange)String message = LocalDateTime.now().toString() + "发送一条消息.";rabbitTemplate.convertAndSend(exchange, unRoutingKey, message, new CorrelationData("unRouting-" + UUID.randomUUID().toString()));log.info("发送一条消息,exchange:[{}],routingKey:[{}],message:[{}]", exchange, unRoutingKey, message);}
}

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

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

  6、生产者 - 启动类

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

 

  7、消费者  - 配置文件

server:port: 9087
spring:rabbitmq:host: 140.143.237.224port: 5672username: rootpassword: rootvirtual-host: /connection-timeout: 10000listener:simple:acknowledge-mode: manual # 手动应答auto-startup: truedefault-requeue-rejected: false # 不重回队列concurrency: 5max-concurrency: 20prefetch: 1 # 每次只处理一个信息retry:enabled: true

 

  8、消费者 - 消息接收

@Component
@Slf4j
public class RabbitConsumer {/*** 监听 queue-rabbit-springboot-advance 队列** @param receiveMessage 接收到的消息* @param message* @param channel*/@RabbitListener(queues = "queue-rabbit-springboot-advance")public void receiveMessage(String receiveMessage, Message message, Channel channel) {try {// 手动签收log.info("接收到消息:[{}]", receiveMessage);if (new Random().nextInt(10) < 5) {log.warn("拒绝一条信息:[{}],此消息将会由死信交换器进行路由.", receiveMessage);channel.basicNack(message.getMessageProperties().getDeliveryTag(), false, false);} else {channel.basicAck(message.getMessageProperties().getDeliveryTag(), false);}} catch (Exception e) {log.info("接收到消息之后的处理发生异常.", e);try {channel.basicAck(message.getMessageProperties().getDeliveryTag(), false);} catch (IOException e1) {log.error("签收异常.", e1);}}}
}

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

  9、执行结果
 

完整代码:

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

这篇关于springboot整合rabbitmq实现生产者消息确认、死信交换器、未路由到队列的消息的文章就介绍到这儿,希望我们推荐的文章对编程师们有所帮助!



http://www.chinasem.cn/article/157555

相关文章

Spring事务传播机制最佳实践

《Spring事务传播机制最佳实践》Spring的事务传播机制为我们提供了优雅的解决方案,本文将带您深入理解这一机制,掌握不同场景下的最佳实践,感兴趣的朋友一起看看吧... 目录1. 什么是事务传播行为2. Spring支持的七种事务传播行为2.1 REQUIRED(默认)2.2 SUPPORTS2

怎样通过分析GC日志来定位Java进程的内存问题

《怎样通过分析GC日志来定位Java进程的内存问题》:本文主要介绍怎样通过分析GC日志来定位Java进程的内存问题,具有很好的参考价值,希望对大家有所帮助,如有错误或未考虑完全的地方,望不吝赐教... 目录一、GC 日志基础配置1. 启用详细 GC 日志2. 不同收集器的日志格式二、关键指标与分析维度1.

Java进程异常故障定位及排查过程

《Java进程异常故障定位及排查过程》:本文主要介绍Java进程异常故障定位及排查过程,具有很好的参考价值,希望对大家有所帮助,如有错误或未考虑完全的地方,望不吝赐教... 目录一、故障发现与初步判断1. 监控系统告警2. 日志初步分析二、核心排查工具与步骤1. 进程状态检查2. CPU 飙升问题3. 内存

Python实现对阿里云OSS对象存储的操作详解

《Python实现对阿里云OSS对象存储的操作详解》这篇文章主要为大家详细介绍了Python实现对阿里云OSS对象存储的操作相关知识,包括连接,上传,下载,列举等功能,感兴趣的小伙伴可以了解下... 目录一、直接使用代码二、详细使用1. 环境准备2. 初始化配置3. bucket配置创建4. 文件上传到os

java中新生代和老生代的关系说明

《java中新生代和老生代的关系说明》:本文主要介绍java中新生代和老生代的关系说明,具有很好的参考价值,希望对大家有所帮助,如有错误或未考虑完全的地方,望不吝赐教... 目录一、内存区域划分新生代老年代二、对象生命周期与晋升流程三、新生代与老年代的协作机制1. 跨代引用处理2. 动态年龄判定3. 空间分

Java设计模式---迭代器模式(Iterator)解读

《Java设计模式---迭代器模式(Iterator)解读》:本文主要介绍Java设计模式---迭代器模式(Iterator),具有很好的参考价值,希望对大家有所帮助,如有错误或未考虑完全的地方,... 目录1、迭代器(Iterator)1.1、结构1.2、常用方法1.3、本质1、解耦集合与遍历逻辑2、统一

Java内存分配与JVM参数详解(推荐)

《Java内存分配与JVM参数详解(推荐)》本文详解JVM内存结构与参数调整,涵盖堆分代、元空间、GC选择及优化策略,帮助开发者提升性能、避免内存泄漏,本文给大家介绍Java内存分配与JVM参数详解,... 目录引言JVM内存结构JVM参数概述堆内存分配年轻代与老年代调整堆内存大小调整年轻代与老年代比例元空

深度解析Java DTO(最新推荐)

《深度解析JavaDTO(最新推荐)》DTO(DataTransferObject)是一种用于在不同层(如Controller层、Service层)之间传输数据的对象设计模式,其核心目的是封装数据,... 目录一、什么是DTO?DTO的核心特点:二、为什么需要DTO?(对比Entity)三、实际应用场景解析

Java 线程安全与 volatile与单例模式问题及解决方案

《Java线程安全与volatile与单例模式问题及解决方案》文章主要讲解线程安全问题的五个成因(调度随机、变量修改、非原子操作、内存可见性、指令重排序)及解决方案,强调使用volatile关键字... 目录什么是线程安全线程安全问题的产生与解决方案线程的调度是随机的多个线程对同一个变量进行修改线程的修改操

关于集合与数组转换实现方法

《关于集合与数组转换实现方法》:本文主要介绍关于集合与数组转换实现方法,具有很好的参考价值,希望对大家有所帮助,如有错误或未考虑完全的地方,望不吝赐教... 目录1、Arrays.asList()1.1、方法作用1.2、内部实现1.3、修改元素的影响1.4、注意事项2、list.toArray()2.1、方