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 正则表达式URL 匹配与源码全解析

《Java正则表达式URL匹配与源码全解析》在Web应用开发中,我们经常需要对URL进行格式验证,今天我们结合Java的Pattern和Matcher类,深入理解正则表达式在实际应用中... 目录1.正则表达式分解:2. 添加域名匹配 (2)3. 添加路径和查询参数匹配 (3) 4. 最终优化版本5.设计思

C#实现将Excel表格转换为图片(JPG/ PNG)

《C#实现将Excel表格转换为图片(JPG/PNG)》Excel表格可能会因为不同设备或字体缺失等问题,导致格式错乱或数据显示异常,转换为图片后,能确保数据的排版等保持一致,下面我们看看如何使用C... 目录通过C# 转换Excel工作表到图片通过C# 转换指定单元格区域到图片知识扩展C# 将 Excel

Java使用ANTLR4对Lua脚本语法校验详解

《Java使用ANTLR4对Lua脚本语法校验详解》ANTLR是一个强大的解析器生成器,用于读取、处理、执行或翻译结构化文本或二进制文件,下面就跟随小编一起看看Java如何使用ANTLR4对Lua脚本... 目录什么是ANTLR?第一个例子ANTLR4 的工作流程Lua脚本语法校验准备一个Lua Gramm

Java字符串操作技巧之语法、示例与应用场景分析

《Java字符串操作技巧之语法、示例与应用场景分析》在Java算法题和日常开发中,字符串处理是必备的核心技能,本文全面梳理Java中字符串的常用操作语法,结合代码示例、应用场景和避坑指南,可快速掌握字... 目录引言1. 基础操作1.1 创建字符串1.2 获取长度1.3 访问字符2. 字符串处理2.1 子字

Java Optional的使用技巧与最佳实践

《JavaOptional的使用技巧与最佳实践》在Java中,Optional是用于优雅处理null的容器类,其核心目标是显式提醒开发者处理空值场景,避免NullPointerExce... 目录一、Optional 的核心用途二、使用技巧与最佳实践三、常见误区与反模式四、替代方案与扩展五、总结在 Java

基于Java实现回调监听工具类

《基于Java实现回调监听工具类》这篇文章主要为大家详细介绍了如何基于Java实现一个回调监听工具类,文中的示例代码讲解详细,感兴趣的小伙伴可以跟随小编一起学习一下... 目录监听接口类 Listenable实际用法打印结果首先,会用到 函数式接口 Consumer, 通过这个可以解耦回调方法,下面先写一个

使用Java将DOCX文档解析为Markdown文档的代码实现

《使用Java将DOCX文档解析为Markdown文档的代码实现》在现代文档处理中,Markdown(MD)因其简洁的语法和良好的可读性,逐渐成为开发者、技术写作者和内容创作者的首选格式,然而,许多文... 目录引言1. 工具和库介绍2. 安装依赖库3. 使用Apache POI解析DOCX文档4. 将解析

Qt中QGroupBox控件的实现

《Qt中QGroupBox控件的实现》QGroupBox是Qt框架中一个非常有用的控件,它主要用于组织和管理一组相关的控件,本文主要介绍了Qt中QGroupBox控件的实现,具有一定的参考价值,感兴趣... 目录引言一、基本属性二、常用方法2.1 构造函数 2.2 设置标题2.3 设置复选框模式2.4 是否

Java字符串处理全解析(String、StringBuilder与StringBuffer)

《Java字符串处理全解析(String、StringBuilder与StringBuffer)》:本文主要介绍Java字符串处理全解析(String、StringBuilder与StringBu... 目录Java字符串处理全解析:String、StringBuilder与StringBuffer一、St

C++使用printf语句实现进制转换的示例代码

《C++使用printf语句实现进制转换的示例代码》在C语言中,printf函数可以直接实现部分进制转换功能,通过格式说明符(formatspecifier)快速输出不同进制的数值,下面给大家分享C+... 目录一、printf 原生支持的进制转换1. 十进制、八进制、十六进制转换2. 显示进制前缀3. 指