一、消费端限流

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 死信处理

  1. 丢弃,如果不是很重要,可以选择丢弃
  2. 记录死信入库,然后做后续的业务分析或处理
  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分钟之内没有支付成功,那么这个订单将进行取消处理。这时就可以使用延时队列将订单信息发送到延时队列。

  • 需求:

  1. 下单后,30分钟未支付,取消订单,回滚库存

  2. 新用户注册成功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:指定延迟交换器的基础路由行为,支持 directtopicfanoutheaders 等 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;
            });

  1. 延迟设置

    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()方法执行

消费者端:

更多推荐