SpringBoot教程(十五) | SpringBoot集成RabbitMq(死信队列、延迟队列)

本文主要是介绍SpringBoot教程(十五) | SpringBoot集成RabbitMq(死信队列、延迟队列),希望对大家解决编程问题提供一定的参考价值,需要的开发者们随着小编来一起学习吧!

SpringBoot教程(十五) | SpringBoot集成RabbitMq(死信队列、延迟队列)

  • (一)死信队列
    • 使用场景
    • 具体用法
    • 前提
    • 示例:
  • (二)延迟队列
    • 使用场景
    • 方法一:通过死亡队列实现
    • 方法二:通过延迟消息插件(rabbitmq_delayed_message_exchange)实现

(一)死信队列

死信队列是一个重要的概念,用于处理那些因各种原因无法被正常消费的消息。
它不是RabbitMQ直接提供的一个现成的方法或工具,而是通过特定的配置和机制来实现的。

使用场景

死信队列在多种场景下都非常有用,包括但不限于:

  1. 消息重试机制:当消息处理失败时,可以将其发送到死信队列进行重试。
  2. 异常消息处理:对于无法被正常处理的异常消息,可以将其存储在死信队列中,以便后续分析处理。
  3. 延迟消息处理:通过结合消息的TTL(Time-To-Live,生存时间)和死信队列,可以实现消息的延迟处理。
  4. 确保消息不丢失:在消息处理过程中,如果发生消费者崩溃或网络故障等情况,消息可能会丢失。通过死信队列,可以确保这些消息得到保留,并在系统恢复后重新处理。

具体用法

要在RabbitMQ中设置和使用死信队列,通常需要按照以下步骤进行:

  1. 定义死信交换机(DLX):首先,需要定义一个交换机作为死信交换机,它可以是任何类型的交换机(如direct、fanout、topic等)。
  2. 配置原队列:在声明原队列时,需要指定两个参数:x-dead-letter-exchange和x-dead-letter-routing-key。前者指定了当消息变成死信时应该发送到的交换机(即死信交换机),后者指定了发送到该交换机的路由键。
  3. 声明死信队列:接着,需要声明一个或多个死信队列,并将它们绑定到死信交换机上。这样,当死信消息被发送到死信交换机时,就可以根据路由键将其路由到相应的死信队列中。
  4. 处理死信消息:最后,需要编写消费者代码来监听死信队列中的消息,并对这些消息进行相应的处理。

前提

要想进入死信队列,得出现异常,出现异常后,会根据你的配置帮你放到死信队列中 所以异常不要被捕获。
如果实在要捕获的话,就得你在消费者这边去做“发送消息的”操作,自己把发送过来消息塞到死信队列中

示例:

消费者 mq的yml配置(重试机制)

spring:rabbitmq:host: 127.0.0.1port: 5672username: guestpassword: guestlistener:simple:# 重试机制retry:enabled: true #是否开启消费者重试

配置类:

package com.example.reactboot.config;import org.springframework.amqp.core.Binding;
import org.springframework.amqp.core.BindingBuilder;
import org.springframework.amqp.core.DirectExchange;
import org.springframework.amqp.core.Queue;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;import java.util.HashMap;
import java.util.Map;@Configuration
public class DirectExchangeConfig {//===========================普通===========================//定义队列的名称常量public static final String DIRECT_QUEUE = "directQueue";public static final String DIRECT_QUEUE2 = "directQueue2";//定义直接交换机的名称常量public static final String DIRECT_EXCHANGE = "directExchange";//定义路由键常量,用于交换机和队列之间的绑定public static final String DIRECT_ROUTING_KEY = "direct";//定义路由键常量,用于交换机和队列之间的绑定public static final String DIRECT_ROUTING_KEY_2= "direct2";//定义队列,名称为DIRECT_QUEUE//为普通队列 绑设置 死信参数@Beanpublic Queue directQueue() {//return new Queue(DIRECT_QUEUE, true);Map<String, Object> args = new HashMap<>();// 设置死信交换机args.put("x-dead-letter-exchange", DLX_EXCHANGE);// 设置发送到死信交换机的路由键args.put("x-dead-letter-routing-key", DLX_ROUTING_KEY);// 创建队列,设置为持久化、非排他、非自动删除,并附带死信参数return new Queue(DIRECT_QUEUE, true, false, false, args);}//定义直接交换机@Beanpublic DirectExchange directExchange() {return new DirectExchange(DIRECT_EXCHANGE, true, false);}//定义队列,名称为DIRECT_QUEUE2@Beanpublic Queue directQueue2() {return new Queue(DIRECT_QUEUE2, true);}//定义一个绑定,将directQueue队列绑定到directExchange交换机上,//使用direct作为路由键@Beanpublic Binding bindingDirectExchange(Queue directQueue, DirectExchange directExchange) {return BindingBuilder.bind(directQueue).to(directExchange).with(DIRECT_ROUTING_KEY);}// 定义一个绑定Bean,将directQueue2队列也绑定到directExchange交换机上,@Beanpublic Binding bindingDirectExchange2(Queue directQueue2, DirectExchange directExchange) {return BindingBuilder.bind(directQueue2).to(directExchange).with(DIRECT_ROUTING_KEY_2);}//===========================死信===========================// 定义死信交换机的名称public static final String DLX_EXCHANGE = "dlx_exchange";// 定义发送到死信交换机的路由键public static final String DLX_ROUTING_KEY = "dlx.routing.key";// 定义死信队列的名称public static final String DLX_QUEUE = "dlx_queue";/*** 声明死信交换机,这里使用Direct类型。* @return 返回一个配置好的DirectExchange对象。*/@BeanDirectExchange dlxExchange() {// 创建并返回Direct类型的交换机return new DirectExchange(DLX_EXCHANGE,true, false);}/*** 声明死信队列。* @return 返回一个配置好的Queue对象,用作死信队列。*/@BeanQueue dlxQueue() {// 创建并返回死信队列,设置为持久化return new Queue(DLX_QUEUE, true);}/*** 绑定死信队列到死信交换机,使用指定的路由键。*/@BeanBinding binding(Queue dlxQueue,DirectExchange dlxExchange) {return BindingBuilder.bind(dlxQueue).to(dlxExchange).with(DLX_ROUTING_KEY);}}

生产者发送消息:

package com.example.reactboot.controller;import com.example.reactboot.config.DirectExchangeConfig;
import org.springframework.amqp.core.Message;
import org.springframework.amqp.core.MessageBuilder;
import org.springframework.amqp.core.MessageDeliveryMode;
import org.springframework.amqp.rabbit.connection.CorrelationData;
import org.springframework.amqp.rabbit.core.RabbitTemplate;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.RestController;import java.nio.charset.StandardCharsets;@RestController
public class RabbitMqTest {@AutowiredRabbitTemplate rabbitTemplate;@RequestMapping("/sendMQ")public String sendMessage() {rabbitTemplate.convertAndSend(DirectExchangeConfig.DIRECT_EXCHANGE,DirectExchangeConfig.DIRECT_ROUTING_KEY, "发送一条测试消息:direct");return "direct消息发送成功!!";}}

消费者消费消息:

package com.example.reactboot.queueListener;import com.example.reactboot.config.DirectExchangeConfig;
import com.rabbitmq.client.Channel;
import lombok.extern.slf4j.Slf4j;
import org.springframework.amqp.core.Message;
import org.springframework.amqp.rabbit.annotation.RabbitHandler;
import org.springframework.amqp.rabbit.annotation.RabbitListener;
import org.springframework.stereotype.Component;import java.io.IOException;
import java.text.SimpleDateFormat;
import java.util.Date;/*** @className: DirectQueueListener* @description: 直连交换机的监听器* @author: sh.Liu* @date: 2021-08-23 16:03*/
@Slf4j
@Component
public class DirectQueueListener {//监听普通队列@RabbitHandler@RabbitListener(queues = DirectExchangeConfig.DIRECT_QUEUE)public void process(String xx){SimpleDateFormat sdf = new SimpleDateFormat("yyyy-MM-dd HH:mm:ss");System.out.println("DirectReceiver消费者收到消息1  : " + xx + " 接收时间:" + sdf.format(new Date()) + "\n");//先执行业务代码int i = 1 / 0;}//监听死信队列@RabbitHandler@RabbitListener(queues = DirectExchangeConfig.DLX_QUEUE)public void process3(String testMessage) {System.out.println("死信得列里面的  : " + testMessage + "\n");}}

(二)延迟队列

延迟队列是一种特殊的消息队列,其内部消息是有序的,并且具有延时属性。
在RabbitMQ中,虽然AMQP协议本身没有直接支持延迟队列,但可以通过一些变通的方法(如使用死信队列配合消息的TTL属性,或者使用RabbitMQ的延迟消息插件)来实现延迟队列的功能。

使用场景

延迟队列在多种业务场景中都有广泛的应用,包括但不限于:

  1. 订单超时未支付自动取消:用户下单后,如果在规定时间内未完成支付,系统可以自动取消订单。
  2. 退款超时通知:用户申请退款后,如果长时间未得到处理,系统可以自动通知相关运营人员介入。
  3. 新用户注册后的引导邮件:用户注册账号后,系统可以在一段时间后发送欢迎邮件或引导邮件。
  4. 会议提醒:在预定的会议开始前一段时间,系统自动发送提醒给参会人员。
  5. 任务调度:在指定时间后执行某项任务,如定时清理日志、执行批处理任务等。

方法一:通过死亡队列实现

以下是使用死信队列配合TTL属性实现延迟队列的基本步骤:

  1. 定义死信交换机(DLX, Dead-Letter Exchange)和死信队列(DLQ, Dead-Letter Queue)
  2. 设置普通队列的TTL和死信交换机:在创建普通队列时,可以为其设置TTL属性,指定消息在该队列中的最大存活时间。同时,需要将该队列的死信交换机设置为前面定义的DLX,以便消息在过期后能够被发送到DLQ。
  3. 生产者发送消息:生产者将消息发送到普通队列,并指定消息的TTL。消息在队列中等待,直到TTL过期。
  4. 消息过期并发送到死信队列:当消息的TTL过期后,RabbitMQ会自动将该消息发送到其配置的死信交换机,再由死信交换机根据路由键将其发送到DLQ。
  5. 消费者从死信队列消费消息:消费者监听DLQ,当有新消息到达时,进行消费处理。

就是:把普通队列的消息设置存活时间,目前有两者方式:
1.在队列上面设置消息的过期时间
2.直接在消息上面设置过期时间。

方式一(队列上面设置消息过期时间):

上面的关于 死信示例 完全可以复用进行测试

在以下的方法里面多加一行 args.put(“x-message-ttl”, 10000);

    //定义队列,名称为DIRECT_QUEUE//为普通队列 绑设置 死信参数@Beanpublic Queue directQueue() {//return new Queue(DIRECT_QUEUE, true);Map<String, Object> args = new HashMap<>();// 设置消息TTL为10秒args.put("x-message-ttl", 10000);// 设置死信交换机args.put("x-dead-letter-exchange", DLX_EXCHANGE);// 设置发送到死信交换机的路由键args.put("x-dead-letter-routing-key", DLX_ROUTING_KEY);// 创建队列,设置为持久化、非排他、非自动删除,并附带死信参数return new Queue(DIRECT_QUEUE, true, false, false, args);}  

你可以把 DirectQueueListener 里面的 process 方法注释掉(以免被消费掉)。
再执行生产者的 sendMessage 方法。
这个时候你就可以看到下面关于 监听死信队列 的方法 ,等10秒后就会打印你发的消息了

方式二(消息上面设置过期时间):

上面的关于 死信示例 完全可以复用进行测试

改一下 这个 生产者发送消息:

package com.example.reactboot.controller;import com.example.reactboot.config.DirectExchangeConfig;
import org.springframework.amqp.core.Message;
import org.springframework.amqp.core.MessageBuilder;
import org.springframework.amqp.core.MessageDeliveryMode;
import org.springframework.amqp.rabbit.connection.CorrelationData;
import org.springframework.amqp.rabbit.core.RabbitTemplate;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.RestController;import java.nio.charset.StandardCharsets;@RestController
public class RabbitMqTest {@AutowiredRabbitTemplate rabbitTemplate;@RequestMapping("/sendMQ")public String sendMessage() {SimpleDateFormat sdf = new SimpleDateFormat("yyyy-MM-dd HH:mm:ss");rabbitTemplate.convertAndSend(DirectExchangeConfig.DIRECT_EXCHANGE, DirectExchangeConfig.DIRECT_ROUTING_KEY, "发送一条测试消息:direct! "+sdf.format(new Date()), new MessagePostProcessor() {@Overridepublic Message postProcessMessage(Message message) throws AmqpException {//设置过期时间,超过5秒消息就会消失message.getMessageProperties().setExpiration("5000");//设置编码格式message.getMessageProperties().setContentEncoding("UTF-8");return message;}});     return "direct消息发送成功!!";}}

你可以把 DirectQueueListener 里面的 process 方法注释掉(以免被消费掉)。
再执行生产者的 sendMessage 方法。
这个时候你就可以看到下面关于 监听死信队列 的方法 ,等5秒后就会打印你发的消息了

到这里其实就结束了,剩下的就是监听到死信队列里面的消息后的业务操作了

方法二:通过延迟消息插件(rabbitmq_delayed_message_exchange)实现

后续在说

这篇关于SpringBoot教程(十五) | SpringBoot集成RabbitMq(死信队列、延迟队列)的文章就介绍到这儿,希望我们推荐的文章对编程师们有所帮助!



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

相关文章

Java使用Curator进行ZooKeeper操作的详细教程

《Java使用Curator进行ZooKeeper操作的详细教程》ApacheCurator是一个基于ZooKeeper的Java客户端库,它极大地简化了使用ZooKeeper的开发工作,在分布式系统... 目录1、简述2、核心功能2.1 CuratorFramework2.2 Recipes3、示例实践3

Springboot处理跨域的实现方式(附Demo)

《Springboot处理跨域的实现方式(附Demo)》:本文主要介绍Springboot处理跨域的实现方式(附Demo),具有很好的参考价值,希望对大家有所帮助,如有错误或未考虑完全的地方,望不... 目录Springboot处理跨域的方式1. 基本知识2. @CrossOrigin3. 全局跨域设置4.

springboot security使用jwt认证方式

《springbootsecurity使用jwt认证方式》:本文主要介绍springbootsecurity使用jwt认证方式,具有很好的参考价值,希望对大家有所帮助,如有错误或未考虑完全的地... 目录前言代码示例依赖定义mapper定义用户信息的实体beansecurity相关的类提供登录接口测试提供一

Spring Boot 3.4.3 基于 Spring WebFlux 实现 SSE 功能(代码示例)

《SpringBoot3.4.3基于SpringWebFlux实现SSE功能(代码示例)》SpringBoot3.4.3结合SpringWebFlux实现SSE功能,为实时数据推送提供... 目录1. SSE 简介1.1 什么是 SSE?1.2 SSE 的优点1.3 适用场景2. Spring WebFlu

基于SpringBoot实现文件秒传功能

《基于SpringBoot实现文件秒传功能》在开发Web应用时,文件上传是一个常见需求,然而,当用户需要上传大文件或相同文件多次时,会造成带宽浪费和服务器存储冗余,此时可以使用文件秒传技术通过识别重复... 目录前言文件秒传原理代码实现1. 创建项目基础结构2. 创建上传存储代码3. 创建Result类4.

Java利用JSONPath操作JSON数据的技术指南

《Java利用JSONPath操作JSON数据的技术指南》JSONPath是一种强大的工具,用于查询和操作JSON数据,类似于SQL的语法,它为处理复杂的JSON数据结构提供了简单且高效... 目录1、简述2、什么是 jsONPath?3、Java 示例3.1 基本查询3.2 过滤查询3.3 递归搜索3.4

Tomcat版本与Java版本的关系及说明

《Tomcat版本与Java版本的关系及说明》:本文主要介绍Tomcat版本与Java版本的关系及说明,具有很好的参考价值,希望对大家有所帮助,如有错误或未考虑完全的地方,望不吝赐教... 目录Tomcat版本与Java版本的关系Tomcat历史版本对应的Java版本Tomcat支持哪些版本的pythonJ

springboot security验证码的登录实例

《springbootsecurity验证码的登录实例》:本文主要介绍springbootsecurity验证码的登录实例,具有很好的参考价值,希望对大家有所帮助,如有错误或未考虑完全的地方,... 目录前言代码示例引入依赖定义验证码生成器定义获取验证码及认证接口测试获取验证码登录总结前言在spring

SpringBoot日志配置SLF4J和Logback的方法实现

《SpringBoot日志配置SLF4J和Logback的方法实现》日志记录是不可或缺的一部分,本文主要介绍了SpringBoot日志配置SLF4J和Logback的方法实现,文中通过示例代码介绍的非... 目录一、前言二、案例一:初识日志三、案例二:使用Lombok输出日志四、案例三:配置Logback一

springboot security快速使用示例详解

《springbootsecurity快速使用示例详解》:本文主要介绍springbootsecurity快速使用示例,具有很好的参考价值,希望对大家有所帮助,如有错误或未考虑完全的地方,望不吝... 目录创www.chinasem.cn建spring boot项目生成脚手架配置依赖接口示例代码项目结构启用s