RabbitMQ消费端限流、消息超时、死信队列、延迟队列
一、消费端限流
1.1 概述

- 生产者发送消息过多,消费端并发达到上限。设定最多从队列取回请求的数量
1.2 代码实现
- 生产者端代码
@Test
public void testSendMessage() {
for (int i = 0; i < 100; i++) {
rabbitTemplate.convertAndSend(
EXCHANGE_DIRECT,
ROUTING_KEY,
"Hello atguigu" + i);
}
}
- 消费端代码
// 正常业务操作
log.info("消费端接收到消息内容:" + dataString);// System.out.println(10 / 0);
TimeUnit.SECONDS.sleep(1);// 给 RabbitMQ 服务器返回 ACK 确认信息
channel.basicAck(deliveryTag, false);
1.3 测试
1.3.1 未使用prefetch(未开启限流)
- 启动消费者前,启动生产者,队列消息情况
Ready表示已经发送到队列的消息数量
Unacked表示已经发送到消费端但是消费端尚未返回ACK信息的消息数量
Total未被删除的消息总数
启动消费者,队列消息情况,所有消息全被消费端取走进行逐个处理

1.3.2 设定prefetch
- YAML配置
spring:
rabbitmq:
host: 虚拟机ip
port: 5672
username: guest
password: 123456
virtual-host: /
listener:
simple:
acknowledge-mode: manual
prefetch: 1 # 设置每次最多从消息队列服务器取回多少消息



消息不是一次性全部取走的,而是有个过程
二、消息超时
2.1 概述
TTL:Time To Live(存活时间/过期时间)
消息到达存活时间后,还没有被消费,会被自动清除。
RabbitMQ可以对消息设置过期时间,也可以对整个队列设置过期时间

2.2 实现
2.2.1 队列层面设置

绑定:

测试:
-
不启动消费端程序
-
向设置了过期时间的队列中发送100条消息
-
等10秒后,看是否全部被过期删除

2.2.2 消息层面设置
发送消息时代码:
import org.springframework.amqp.core.Message;
import org.springframework.amqp.core.MessagePostProcessor;@Test
public void testSendMessageTTL() {
// 1、创建消息后置处理器对象
MessagePostProcessor messagePostProcessor = (Message message) -> {
// 设定 TTL 时间,以毫秒为单位
message.getMessageProperties().setExpiration("5000");
return message;
};
// 2、发送消息
rabbitTemplate.convertAndSend(
EXCHANGE_DIRECT,
ROUTING_KEY,
"Hello", messagePostProcessor);
}
效果:

三、死信队列
3.1 概述
3.1.1 什么是死信队列(DLX)
- 死信队列,英文缩写:DLX 。Dead Letter Exchange(死信交换机),当消息成为Dead message后,可以被重新发送到另一个交换机,这个交换机就是DLX。

- 死信,无法被消费的消息。
3.1.2 消息成为死信的条件(三种情况)
-
拒绝:消费者拒接消息,basicNack()/basicReject(),并且不把消息重新放入原目标队列,requeue=false
-
溢出:队列中消息数量到达限制。比如队列最大只能存储10条消息,且现在已经存储了10条,此时如果再发送一条消息进来,根据先进先出原则,队列中最早的消息会变成死信
-
超时:消息到达超时时间未被消费
3.1.3 死信处理
- 丢弃,如果不是很重要,可以选择丢弃
- 记录死信入库,然后做后续的业务分析或处理
- 通过死信队列,由负责监听死信的应用程序进行处理
- 通过死信队列,将产生的死信通过程序的配置路由到指定的死信队列,然后应用监听死信队列,对接收到的死信做后续的处理
3.2 消费端拒收消息
3.2.1 死信资源准备
a.创建死信交换机和死信队列
-
死信交换机:exchange.dead.letter.video
-
死信队列:queue.dead.letter.video
-
死信路由键:routing.key.dead.letter.video
正常创建即可
b.创建正常交换机和正常队列
-
正常交换机:exchange.normal.video
-
正常队列:queue.normal.video
-
正常路由键:routing.key.normal.video

Arguments(参数)配置
- x - dead - letter - exchange:指定死信交换器,当队列中的消息成为死信(如过期、被否定确认等情况 )时,会被发送到该交换器,图中值为 “exchange.dead.letter.video” 。
- x - dead - letter - routing - key:指定死信路由键,配合死信交换器,决定死信消息的路由,值为 “routing.key.dead.letter.video” 。
- x - max - length:限制队列中消息的最大数量,图中设为 10 ,达到此数量后,新消息入队会按溢出行为处理 。
- x - message - ttl:消息存活时间(Time To Live ),单位毫秒,图中为 10000 ,即消息在队列中最长存活 10 秒,超时未被消费将成为死信 。
全部设置完成后:

3.2.2 代码实现
a.常量声明代码
public static final String EXCHANGE_NORMAL = "exchange.normal.video";
public static final String EXCHANGE_DEAD_LETTER = "exchange.dead.letter.video";
public static final String ROUTING_KEY_NORMAL = "routing.key.normal.video";
public static final String ROUTING_KEY_DEAD_LETTER = "routing.key.dead.letter.video";
public static final String QUEUE_NORMAL = "queue.normal.video";
public static final String QUEUE_DEAD_LETTER = "queue.dead.letter.video";
b.发送消息代码
@Test
public void testSendMessageButReject() {
rabbitTemplate
.convertAndSend(
EXCHANGE_NORMAL,
ROUTING_KEY_NORMAL,
"测试死信情况1:消息被拒绝");
}
c.接收消息代码
- 消费端正常队列监听
@RabbitListener(queues = {QUEUE_NORMAL})
public void processMessageNormal(Message message, Channel channel) throws IOException {
// 监听正常队列,但是拒绝消息
log.info("★[normal]消息接收到,但我拒绝。");
channel.basicReject(message.getMessageProperties().getDeliveryTag(), false);
}
- 消费端死信队列监听
@RabbitListener(queues = {QUEUE_DEAD_LETTER})
public void processMessageDead(String dataString, Message message, Channel channel) throws IOException {
// 监听死信队列
log.info("★[dead letter]dataString = " + dataString);
log.info("★[dead letter]我是死信监听方法,我接收到了死信消息");
channel.basicAck(message.getMessageProperties().getDeliveryTag(), false);
}
3.2.3 输出结果

3.3 消息数量超过队列容纳极限
3.3.1 死信资源准备上同
3.3.2 代码实现
a.生产者代码
@Test
public void testSendMultiMessage() {
for (int i = 0; i < 20; i++) {
rabbitTemplate.convertAndSend(
EXCHANGE_NORMAL,
ROUTING_KEY_NORMAL,
"测试死信情况2:消息数量超过队列的最大容量" + i);
}
}
b.消费端接收代码
@RabbitListener(queues = {QUEUE_NORMAL})
public void processMessageNormal(Message message, Channel channel) throws IOException {
// 监听正常队列
log.info("★[normal]消息接收到。");
channel.basicAck(message.getMessageProperties().getDeliveryTag(), false);
}
重启服务
3.3.3 输出结果
消费端死信队列接收到前10条消息:

3.4 消息超时未消费
3.4.1 死信资源准备上同
3.4.2 代码实现
a.生产者代码
@Test
public void testSendMessageTimeout() {
rabbitTemplate
.convertAndSend(
EXCHANGE_NORMAL,
ROUTING_KEY_NORMAL,
"测试死信情况3:消息超时");
}
b.不写消费代码,让队列中消息超时
重启服务
3.3.3 输出结果
没有消费端监听程序,所以消息未超时前滞留在队列中:

消息超时后,进入死信队列:

四、延迟队列
4.1 概述
4.1.1 什么是延迟队列
-
延迟队列存储的对象是对应的延时消息,所谓”延时消息”是指当消息被发送以后,并不想让消费者立即拿到消息,而是等待指定时间后,消费者才拿到这个消息进行消费。
-
场景:在订单系统中,一个用户下单之后通常有30分钟的时间进行支付,如果30分钟之内没有支付成功,那么这个订单将进行取消处理。这时就可以使用延时队列将订单信息发送到延时队列。
-
需求:
-
下单后,30分钟未支付,取消订单,回滚库存
-
新用户注册成功30分钟后,发送短信问候
4.1.2 延迟队列实现
- 延迟队列实现需求

RabbitMQ中没有提供延时队列功能
方案1:借助消息超时时间+死信队列
//转发到死信队列时,大多是信息会出现乱序
方案2:给RabbitMQ安装插件

4.2 延迟插件
插件官网地址:https://github.com/rabbitmq/rabbitmq-delayed-message-exchange
延迟期限:最多两天
4.2.1 安装插件
- 卷映射目录

和容器内/plugins目录对应的宿主机目录是:/var/lib/docker/volumes/rabbitmq-plugin/_data
- 下载延迟插件
官方文档说明地址:Community Plugins | RabbitMQ

下载安装插件:
wget https://github.com/rabbitmq/rabbitmq-delayed-message-exchange/releases/download/v3.12.0/rabbitmq_delayed_message_exchange-3.12.0.ez
mv rabbitmq_delayed_message_exchange-3.12.0.ez /var/lib/docker/volumes/rabbitmq-plugin/_data
也可直接将文件放到上方卷目录中。
- 启用插件
# 登录进入容器内部
docker exec -it rabbitmq /bin/bash# rabbitmq-plugins命令所在目录已经配置到$PATH环境变量中了,可以直接调用
rabbitmq-plugins enable rabbitmq_delayed_message_exchange# 退出Docker容器
exit# 重启Docker容器
docker restart rabbitmq
- 确认
确认点1:在 RabbitMQ 管理界面,依次进入 Overview(概览)→ Nodes(节点)→ Advanced(高级)→ Plugins(插件) ,通过该路径可查看当前节点已启用的插件列表 。

确认点2:如果创建新交换机时可以在type中看到x-delayed-message选项,那就说明插件安装好了

4.2.2 创建交换机
rabbitmq_delayed_message_exchange插件在工作时要求交换机是x-delayed-message类型才可以,创建方式如下:

-
x-delayed-type:指定延迟交换器的基础路由行为,支持direct、topic、fanout、headers等 RabbitMQ 原生交换器类型。 -
这里填
direct,表示延迟消息的路由逻辑将遵循direct交换器的规则(按Routing Key精准匹配队列)。
4.2.3 代码测试
1.生产者端
@Test
public void testSendDelayMessage() {
rabbitTemplate.convertAndSend(
EXCHANGE_DELAY,//延迟消息交换机名称(对应管理界面的
exchange.delay.happy)。
ROUTING_KEY_DELAY,//
ROUTING_KEY_DELAY:路由键,用于绑定队列。
"测试基于插件的延迟消息 [" + new SimpleDateFormat("hh:mm:ss").format(new Date()) + "]",//消息内容:包含当前时间戳的测试消息,格式为
测试基于插件的延迟消息 [HH:mm:ss]。
messageProcessor -> {// 设置延迟时间:以毫秒为单位
messageProcessor.getMessageProperties().setHeader("x-delay", "10000");return messageProcessor;
});
延迟设置:
messageProcessor.getMessageProperties().setHeader("x-delay", "10000"):设置x-delay头部为10000毫秒(即 10 秒),表示消息将在发送后延迟 10 秒才会被投递到队列。}
2、消费者端
可分两种情况
- 资源已经创建
直接进行消费即可
import com.rabbitmq.client.Channel;
import lombok.extern.slf4j.Slf4j;
import org.springframework.amqp.core.Message;
import org.springframework.amqp.rabbit.annotation.RabbitListener;
import org.springframework.stereotype.Component;import java.io.IOException;
import java.text.SimpleDateFormat;
import java.util.Date;
@Component
@Slf4j
public class MyDelayMessageListener {
public static final String QUEUE_DELAY = "queue.delay.video";
@RabbitListener(queues = {QUEUE_DELAY})
public void process(String dataString,Message message,Channel channel) throws IOException {
log.info("[生产者]" + dataString);
log.info("[消费者]" + new SimpleDateFormat("hh:mm:ss").format(new Date()));
channel.basicAck(message.getMessageProperties().getDeliveryTag(), false);
}}
- 资源未创建,消费者开启手动
import com.rabbitmq.client.Channel;
import lombok.extern.slf4j.Slf4j;
import org.springframework.amqp.core.Message;
import org.springframework.amqp.rabbit.annotation.*;
import org.springframework.stereotype.Component;import java.io.IOException;
import java.text.SimpleDateFormat;
import java.util.Date;
@Component
@Slf4j
public class MyDelayMessageListener {
public static final String EXCHANGE_DELAY = "exchange.delay.video";
public static final String ROUTING_KEY_DELAY = "routing.key.delay.video";
public static final String QUEUE_DELAY = "queue.delay.video";
@RabbitListener(bindings = @QueueBinding(
value = @Queue(value = QUEUE_DELAY, durable = "true", autoDelete = "false"),
exchange = @Exchange(
value = EXCHANGE_DELAY,
durable = "true",
autoDelete = "false",
type = "x-delayed-message",
arguments = @Argument(name = "x-delayed-type", value = "direct")),
key = {ROUTING_KEY_DELAY}
))
public void process(String dataString, Message message, Channel channel) throws IOException {
log.info("[生产者]" + dataString);
log.info("[消费者]" + new SimpleDateFormat("hh:mm:ss").format(new Date()));
channel.basicAck(message.getMessageProperties().getDeliveryTag(), false);
}}
3.执行效果
交换机类型:
![]()
生产者端:使用插件后,消息发送成功也会触发returnedMessage()方法执行

消费者端:

更多推荐


所有评论(0)