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

相关文章

线上Java OOM问题定位与解决方案超详细解析

《线上JavaOOM问题定位与解决方案超详细解析》OOM是JVM抛出的错误,表示内存分配失败,:本文主要介绍线上JavaOOM问题定位与解决方案的相关资料,文中通过代码介绍的非常详细,需要的朋... 目录一、OOM问题核心认知1.1 OOM定义与技术定位1.2 OOM常见类型及技术特征二、OOM问题定位工具

Python的Darts库实现时间序列预测

《Python的Darts库实现时间序列预测》Darts一个集统计、机器学习与深度学习模型于一体的Python时间序列预测库,本文主要介绍了Python的Darts库实现时间序列预测,感兴趣的可以了解... 目录目录一、什么是 Darts?二、安装与基本配置安装 Darts导入基础模块三、时间序列数据结构与

基于 Cursor 开发 Spring Boot 项目详细攻略

《基于Cursor开发SpringBoot项目详细攻略》Cursor是集成GPT4、Claude3.5等LLM的VSCode类AI编程工具,支持SpringBoot项目开发全流程,涵盖环境配... 目录cursor是什么?基于 Cursor 开发 Spring Boot 项目完整指南1. 环境准备2. 创建

Python使用FastAPI实现大文件分片上传与断点续传功能

《Python使用FastAPI实现大文件分片上传与断点续传功能》大文件直传常遇到超时、网络抖动失败、失败后只能重传的问题,分片上传+断点续传可以把大文件拆成若干小块逐个上传,并在中断后从已完成分片继... 目录一、接口设计二、服务端实现(FastAPI)2.1 运行环境2.2 目录结构建议2.3 serv

C#实现千万数据秒级导入的代码

《C#实现千万数据秒级导入的代码》在实际开发中excel导入很常见,现代社会中很容易遇到大数据处理业务,所以本文我就给大家分享一下千万数据秒级导入怎么实现,文中有详细的代码示例供大家参考,需要的朋友可... 目录前言一、数据存储二、处理逻辑优化前代码处理逻辑优化后的代码总结前言在实际开发中excel导入很

Spring Security简介、使用与最佳实践

《SpringSecurity简介、使用与最佳实践》SpringSecurity是一个能够为基于Spring的企业应用系统提供声明式的安全访问控制解决方案的安全框架,本文给大家介绍SpringSec... 目录一、如何理解 Spring Security?—— 核心思想二、如何在 Java 项目中使用?——

SpringBoot+RustFS 实现文件切片极速上传的实例代码

《SpringBoot+RustFS实现文件切片极速上传的实例代码》本文介绍利用SpringBoot和RustFS构建高性能文件切片上传系统,实现大文件秒传、断点续传和分片上传等功能,具有一定的参考... 目录一、为什么选择 RustFS + SpringBoot?二、环境准备与部署2.1 安装 RustF

Nginx部署HTTP/3的实现步骤

《Nginx部署HTTP/3的实现步骤》本文介绍了在Nginx中部署HTTP/3的详细步骤,文中通过示例代码介绍的非常详细,对大家的学习或者工作具有一定的参考学习价值,需要的朋友们下面随着小编来一起学... 目录前提条件第一步:安装必要的依赖库第二步:获取并构建 BoringSSL第三步:获取 Nginx

springboot中使用okhttp3的小结

《springboot中使用okhttp3的小结》OkHttp3是一个JavaHTTP客户端,可以处理各种请求类型,比如GET、POST、PUT等,并且支持高效的HTTP连接池、请求和响应缓存、以及异... 在 Spring Boot 项目中使用 OkHttp3 进行 HTTP 请求是一个高效且流行的方式。

java.sql.SQLTransientConnectionException连接超时异常原因及解决方案

《java.sql.SQLTransientConnectionException连接超时异常原因及解决方案》:本文主要介绍java.sql.SQLTransientConnectionExcep... 目录一、引言二、异常信息分析三、可能的原因3.1 连接池配置不合理3.2 数据库负载过高3.3 连接泄漏