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

相关文章

Oracle查询优化之高效实现仅查询前10条记录的方法与实践

《Oracle查询优化之高效实现仅查询前10条记录的方法与实践》:本文主要介绍Oracle查询优化之高效实现仅查询前10条记录的相关资料,包括使用ROWNUM、ROW_NUMBER()函数、FET... 目录1. 使用 ROWNUM 查询2. 使用 ROW_NUMBER() 函数3. 使用 FETCH FI

Python脚本实现自动删除C盘临时文件夹

《Python脚本实现自动删除C盘临时文件夹》在日常使用电脑的过程中,临时文件夹往往会积累大量的无用数据,占用宝贵的磁盘空间,下面我们就来看看Python如何通过脚本实现自动删除C盘临时文件夹吧... 目录一、准备工作二、python脚本编写三、脚本解析四、运行脚本五、案例演示六、注意事项七、总结在日常使用

Java实现Excel与HTML互转

《Java实现Excel与HTML互转》Excel是一种电子表格格式,而HTM则是一种用于创建网页的标记语言,虽然两者在用途上存在差异,但有时我们需要将数据从一种格式转换为另一种格式,下面我们就来看看... Excel是一种电子表格格式,广泛用于数据处理和分析,而HTM则是一种用于创建网页的标记语言。虽然两

java图像识别工具类(ImageRecognitionUtils)使用实例详解

《java图像识别工具类(ImageRecognitionUtils)使用实例详解》:本文主要介绍如何在Java中使用OpenCV进行图像识别,包括图像加载、预处理、分类、人脸检测和特征提取等步骤... 目录前言1. 图像识别的背景与作用2. 设计目标3. 项目依赖4. 设计与实现 ImageRecogni

Java中Springboot集成Kafka实现消息发送和接收功能

《Java中Springboot集成Kafka实现消息发送和接收功能》Kafka是一个高吞吐量的分布式发布-订阅消息系统,主要用于处理大规模数据流,它由生产者、消费者、主题、分区和代理等组件构成,Ka... 目录一、Kafka 简介二、Kafka 功能三、POM依赖四、配置文件五、生产者六、消费者一、Kaf

Java访问修饰符public、private、protected及默认访问权限详解

《Java访问修饰符public、private、protected及默认访问权限详解》:本文主要介绍Java访问修饰符public、private、protected及默认访问权限的相关资料,每... 目录前言1. public 访问修饰符特点:示例:适用场景:2. private 访问修饰符特点:示例:

详解Java如何向http/https接口发出请求

《详解Java如何向http/https接口发出请求》这篇文章主要为大家详细介绍了Java如何实现向http/https接口发出请求,文中的示例代码讲解详细,感兴趣的小伙伴可以跟随小编一起学习一下... 用Java发送web请求所用到的包都在java.net下,在具体使用时可以用如下代码,你可以把它封装成一

使用Python实现在Word中添加或删除超链接

《使用Python实现在Word中添加或删除超链接》在Word文档中,超链接是一种将文本或图像连接到其他文档、网页或同一文档中不同部分的功能,本文将为大家介绍一下Python如何实现在Word中添加或... 在Word文档中,超链接是一种将文本或图像连接到其他文档、网页或同一文档中不同部分的功能。通过添加超

windos server2022里的DFS配置的实现

《windosserver2022里的DFS配置的实现》DFS是WindowsServer操作系统提供的一种功能,用于在多台服务器上集中管理共享文件夹和文件的分布式存储解决方案,本文就来介绍一下wi... 目录什么是DFS?优势:应用场景:DFS配置步骤什么是DFS?DFS指的是分布式文件系统(Distr

NFS实现多服务器文件的共享的方法步骤

《NFS实现多服务器文件的共享的方法步骤》NFS允许网络中的计算机之间共享资源,客户端可以透明地读写远端NFS服务器上的文件,本文就来介绍一下NFS实现多服务器文件的共享的方法步骤,感兴趣的可以了解一... 目录一、简介二、部署1、准备1、服务端和客户端:安装nfs-utils2、服务端:创建共享目录3、服