赞
踩
RabbitMQ 的死信队列(Dead-Letter-Exchanges,简称 DLX)是一个强大的特性,它允许在消息在队列中无法被正常消费(例如,消息被拒绝并且没有设置重新入队,或者消息过期)时,将这些消息转发到另一个交换机。这个特性在很多场景下都非常有用,比如重试机制、延迟队列等。
要设置死信队列,你需要在队列声明时指定几个参数:
x-dead-letter-exchange:指定消息在变为死信后要发送到的交换机。
x-dead-letter-routing-key(可选):指定消息在变为死信后使用的路由键。如果未设置,则使用原消息的路由键。
message-ttl 或 x-message-ttl(可选):设置消息的生存时间(TTL)。当消息在队列中的时间超过此值后,它将成为死信。
x-max-length(可选):设置队列的最大长度。当队列中的消息数量超过此值时,最早的消息将成为死信。
这些参数可以通过 RabbitMQ 的管理界面、命令行工具或编程 API 设置。
当一个消息在队列中由于某些原因(如过期、被拒绝且未设置重新入队、队列达到最大长度等)成为死信时。
RabbitMQ 会检查该队列是否配置了 x-dead-letter-exchange。
如果配置了,RabbitMQ 会将死信发送到指定的死信交换机。
死信交换机再根据配置的路由键或原消息的路由键将消息路由到相应的队列。
RabbitMQ是建立在强大的Erlang OTP平台上,因此安装Rabbit MQ的前提是安装Erlang。
因为RabbitMQ服务器是用Erlang语言编写的, 所以,你需要去查看rabbitMq适应Erlang的版本,因为不同的rabbitMq版本对应不同的Erlang版本,可以点击如下该链接查看版本匹配度:
https://www.rabbitmq.com/which-erlang.html#compatibility-matrix
下载地址:Erlang
推荐使用链接: https://download.csdn.net/download/weixin_42123075/89064540,这里包含Erlang和对应版本的RabbitMQ。
下载完成后先安装Erlang。
下载地址: https://github.com/rabbitmq/rabbitmq-server/releases?page=7
rabbitmq-plugins list
如下图所示
rabbitmq-plugins enable rabbitmq_management
如下图所示
安装rabbitMq的目录(我的是D:\Software\rabbitmq\rabbitmq_server-3.8.15) -> sbin目录 -> 双击rabbitmq-server.bat,我如下图所示:
rabbitmq-server -detached
创建普通队列myQueue和普通交换机myExchange,交换机类型为Topic,并指定死信队列
<!-- RabbitMQ -->
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-amqp</artifactId>
</dependency>
spring:
rabbitmq:
addresses: localhost:5672
connection-timeout: 15000
password: guest
username: guest
# 使用启用消息确认模式
# publisher-confirms: true
virtual-host: /
import org.springframework.amqp.core.*;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
@Configuration
public class RabbitMQConfig {
// 普通队列
public final static String QUEUE_NAME = "myQueue";
// 普通交换机
public final static String EXCHANGE_NAME = "myExchange";
// 普通队列路由
public final static String ROUTING_KEY = "myRoutingKey";
// 死信交换机
public static final String DEAD_LETTER_EXCHANGE = "dead-letter-exchange";
// 死信队列
public static final String DEAD_LETTER_QUEUE = "dead-letter-queue";
// 死信路由
public static final String DEAD_LETTER_ROUTING_KEY = "dead-letter-key";
@Bean
TopicExchange myExchange() {
return new TopicExchange(EXCHANGE_NAME);
}
@Bean
Binding binding(Queue myQueue, TopicExchange myExchange) {
return BindingBuilder.bind(myQueue).to(myExchange).with(ROUTING_KEY);
}
/**
* 定义死信交换机
* @return DirectExchange
*/
@Bean
public DirectExchange deadLetterExchange() {
return new DirectExchange(DEAD_LETTER_EXCHANGE);
}
/**
* 定义死信队列
* @return Queue
*/
@Bean
public Queue deadLetterQueue() {
return new Queue(DEAD_LETTER_QUEUE,true,false,false,null);
}
/**
* 死信队列绑定死信交换机
* @param deadLetterQueue 死信队列
* @param deadLetterExchange 死信交换机
* @return Binding
*/
@Bean
public Binding deadLetterBinding(Queue deadLetterQueue, DirectExchange deadLetterExchange) {
return BindingBuilder.bind(deadLetterQueue).to(deadLetterExchange).with(DEAD_LETTER_ROUTING_KEY);
}
/**
* 普通队列声明指定死信交换机
* @return
*/
@Bean
public Queue myQueue() {
return QueueBuilder.durable(QUEUE_NAME)
// 设置死信交换机
.withArgument("x-dead-letter-exchange", DEAD_LETTER_EXCHANGE)
// 设置死信路由键
.withArgument("x-dead-letter-routing-key", DEAD_LETTER_ROUTING_KEY)
// 设置队列最大长度
// .withArgument("x-max-length", 5)
.build();
}
}
import com.ruoyi.quartz.config.RabbitMQConfig;
import org.springframework.amqp.core.AmqpTemplate;
import org.springframework.amqp.core.Message;
import org.springframework.amqp.core.MessageBuilder;
import org.springframework.amqp.core.MessageProperties;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service;
@Service
public class RabbitMQService {
@Autowired
private AmqpTemplate rabbitTemplate;
public void send(String message) {
rabbitTemplate.convertAndSend(RabbitMQConfig.EXCHANGE_NAME, RabbitMQConfig.ROUTING_KEY, message);
}
/**
* 发送消息并设置过期时间
* @param exchange 交换机
* @param routingKey 路由
* @param messageBody 消息体
* @param expirationTimeInMillis 过期时间,单位:毫秒
*/
public void sendMessageWithExpiration(String exchange, String routingKey, String messageBody, int expirationTimeInMillis) {
MessageProperties properties = new MessageProperties();
// 设置消息的过期时间
properties.setExpiration(String.valueOf(expirationTimeInMillis));
Message message = MessageBuilder.withBody(messageBody.getBytes())
.andProperties(properties)
.build();
rabbitTemplate.convertAndSend(exchange, routingKey, message);
}
}
import com.ruoyi.quartz.config.RabbitMQConfig;
import org.springframework.amqp.rabbit.annotation.RabbitListener;
import org.springframework.stereotype.Service;
import java.text.SimpleDateFormat;
import java.util.Date;
@Service
public class ReceiverService {
/**
* 监听队列消息 如果要测试死信队列,就不要监听此队列
* @param message 消息
*/
// @RabbitListener(queues = RabbitMQConfig.QUEUE_NAME)
public void receive(String message) {
System.out.println("Received <" + message + ">");
}
/**
* 监听死信队列
* @param message
*/
@RabbitListener(queues = RabbitMQConfig.DEAD_LETTER_QUEUE)
public void processDeadLetter(String message) {
String time = new SimpleDateFormat("yyyy-MM-dd HH:mm:ss").format(new Date());
System.out.println("Received dead letter message: " + message + ",当前时间:" + time);
// 处理死信队列中的消息
}
}
@GetMapping("/sendMessageTtl/{message}")
public void sendMessageTtl(@PathVariable String message){
log.info("当前时间发送:{},发送5条消息给两个TTL队列:{}",new Date().toString(),message);
for (int i = 0; i < 6; i++) {
System.out.println("测试延迟队列======="+DateUtils.parseDateToStr(DateUtils.YYYY_MM_DD_HH_MM_SS,new Date()));
rabbitMQService.sendMessageWithExpiration(RabbitMQConfig.EXCHANGE_NAME,RabbitMQConfig.ROUTING_KEY,"测试5秒延迟==============》",5000);
}
}
Copyright © 2003-2013 www.wpsshop.cn 版权所有,并保留所有权利。