首页技术栈归档照片墙音乐日记随想收藏夹友链留言关于

RabbitMQ延迟队列(基于插件)

写作时间:2026-07-08

RabbitMQ延迟队列(基于插件)

安装插件(Linux安装Docker安装)

Linux

注:版本要一致

  1. 查看RabbitMQ版本号
rabbitmqctl status | grep rabbit
  1. 选择与RabbitMQ版本号对应的插件
# 1. 进入到RabbitMQ插件文件夹所在位置,例如
cd /usr/lib/rabbitmq/lib/rabbitmq_server-3.6.9/plugins

# 2. 下载插件
wget https://dl.bintray.com/rabbitmq/community-plugins/3.3.5/rabbitmq_delayed_message_exchange/rabbitmq_delayed_message_exchange-20171215-3.3.5.zip
  1. 解压缩
unzip -zxvf rabbitmq_delayed_message_exchange-20171215-3.3.5.zip
  1. 启用延时队列插件
# 切换到sbin目录下面
cd /usr/lib/rabbitmq/lib/rabbitmq_server-3.3.5/sbin

# 启用延时队列插件
./rabbitmq-plugins enable rabbitmq_delayed_message_exchange
  1. 重启RabbitMQ服务
systemctl restart rabbitmq-server

或者

在官网上下载 https://www.rabbitmq.com/community-plugins.html,下载

rabbitmq_delayed_message_exchange插件,然后解压放置到RabbitMQ的插件目录

进入RabbitMQ的安装目录下的plgins.目录,执行下面命令让该插件生效,然后重启RabbitMQ

Liunx默认安装目录:/usr/lib/rabbitmq/lib/rabbitmq _server-3.9.13/plugins

安装插件:rabbitmq-plugins enable rabbitmq_delayed_message_exchange

重启服务即可:restart rabbitmq-server

使用Docker进行安装(建议使用)

首先将下载的插件上传到我们的服务器

使用docker ps -a命令查看RabbitMQ容器id

然后进入到容器内部,然后可以看到plugins目录

[root@yuan sbin]# docker exec -it 93ab210462df  /bin/bash
root@93ab210462df:/# ls -a
.  ..  .dockerenv  bin  boot  dev  etc  home  lib  lib32  lib64  libx32  media  mnt  opt  plugins  proc  root  run  sbin  srv  sys  tmp  usr  var
root@93ab210462df:/# 

新开一个窗口,将插件拷贝到容器内plugins目录

[root@yuan rabbitmq]# docker cp rabbitmq_delayed_message_exchange-3.9.0.ez 93ab210462df:/plugins

在容器查看插件是否存在

root@93ab210462df:/# cd plugins
root@93ab210462df:/plugins# ls |grep delay
rabbitmq_delayed_message_exchange-3.9.0.ez

使用命令rabbitmq-plugins enable rabbitmq_delayed_message_exchange启动插件

root@93ab210462df:/plugins# rabbitmq-plugins enable rabbitmq_delayed_message_exchange
Enabling plugins on node rabbit@93ab210462df:
rabbitmq_delayed_message_exchange
The following plugins have been configured:
  rabbitmq_delayed_message_exchange
  rabbitmq_management
  rabbitmq_management_agent
  rabbitmq_prometheus
  rabbitmq_web_dispatch
Applying plugin configuration to rabbit@93ab210462df...
The following plugins have been enabled:
  rabbitmq_delayed_message_exchange

started 1 plugins.

重启容器

root@93ab210462df:/plugins# exit
exit
[root@yuan]# docker restart af99480e815d

基于插件的延迟队列(配置类)

在我们自定义的交换机中,这是一种新的交换类型,该类型消息支持延迟投递机制消息传递后并不会立即投递到目标队列中,而是存储在mnesia(一个分布式数据系统)表中,当达到投递时间时,才投递到目标队列中。

package com.springbootrabbitmq.config;
import org.springframework.amqp.core.Binding;
import org.springframework.amqp.core.BindingBuilder;
import org.springframework.amqp.core.CustomExchange;
import org.springframework.amqp.core.Queue;
import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import java.util.HashMap;
/**
 * 基于插件的延迟队列
 */
@Configuration
public class DelayedQueueConfig {
    //队列
    public static final String DELAYED_QUEUE_NAME = "delayed.queue";
    //交换机
    public static final String DELAYED_EXCHANGE_NAME = "delayed.exchange";
    //routingKey
    public static final String DELAYED_ROUTING_KEY="delayed.routingkey";

    @Bean
    public Queue delayedQue(){
        return new Queue(DELAYED_QUEUE_NAME);
    }
    //声明交换机 基于插件
    @Bean
    public CustomExchange declareExchange(){
        HashMap<String, Object> arguments = new HashMap<>();
        arguments.put("x-delayed-type","direct");
        /**
         * 1.交换机的名称
         * 2.交换机的类型
         * 3.是否需要持久化
         * 4.是否自动删除
         * 5.其他参数
         */
        return new CustomExchange(DELAYED_EXCHANGE_NAME,"x-delayed-message",true,false,arguments);
    }
    //绑定
    @Bean
    public Binding delayedQueueBinding(
            @Qualifier("delayedQue") Queue delayedQue,
            @Qualifier("declareExchange") CustomExchange declareExchange
    ){
        //把delayedQue榜给xeclareExchange直接它的routingKey为DELAYED_ROUTING_KEY,然后进行构建
        return BindingBuilder.bind(delayedQue).to(declareExchange).with(DELAYED_ROUTING_KEY).noargs();
    }
}

基于插件的延迟队列(生产者)

 //开始发送消息 基于插件的消息 及延迟时间
    @GetMapping("/sendMegs/{message}/{delayTime}")
    public void sendMsg(@PathVariable String message,@PathVariable Integer delayTime){
        log.info("当前时间:{},发送一条时长{}毫秒的信息个延迟队列delayed.queue:{}",new Date().toString(),delayTime,message);
        rabbitTemplate.convertAndSend(DelayedQueueConfig.DELAYED_EXCHANGE_NAME,DelayedQueueConfig.DELAYED_ROUTING_KEY,message,mes->{
           // 发送消息的时候 延迟时长 单位:ms
            mes.getMessageProperties().setDelay(delayTime);
            return mes;
        });

基于插件的延迟队列(消费者)

package com.springbootrabbitmq.consumer;
import com.springbootrabbitmq.config.DelayedQueueConfig;
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.util.Date;
/**
 * 消费者基于插件
 */
@Slf4j
@Component
public class DelayQueueConsumer {
    //监听消息
    @RabbitListener(queues = DelayedQueueConfig.DELAYED_QUEUE_NAME)
    public void receiveDelayQueue(Message message){
        String msg=new String(message.getBody());
        log.info("当前时间:{},收到的延迟队列的消息:{}",new Date().toString(),msg);
    }
}

总结

延时队列在需要延时处理的场景下非常有用,使用RabbitMQ来实现延时队列可以很好的利用

RabbitMQ的特性,如:消息可靠发送、消息可靠投递、死信队列来保障消息至少被消费一次以及未被正确处理的消息不会被丢弃。另外,通过RabbitMQ.集群的特性,可以很好的解决单点故障问题,不会因为单个节点挂掉导致延时队列不可用或者消息丢失。

当然,延时队列还有很多其它选择,比如利用Java的DelayQueue,利用Redis.的zset,利用Quartz或者利用kafka.的时间轮,这些方式各有特点,看需要适用的场景

avatar

yuanyourdomain

写代码,做研究,记录生活。

RECOMMENDED

MyBatis 动态 SQL

2026-07-08

Nginx 基础入门

2026-07-08

Maven 多模块与私服

2026-07-08