RabbitMQ延迟队列(基于插件)
安装插件(Linux安装Docker安装)
Linux
注:版本要一致
- 查看RabbitMQ版本号
rabbitmqctl status | grep rabbit
- 选择与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
- 解压缩
unzip -zxvf rabbitmq_delayed_message_exchange-20171215-3.3.5.zip
- 启用延时队列插件
# 切换到sbin目录下面
cd /usr/lib/rabbitmq/lib/rabbitmq_server-3.3.5/sbin
# 启用延时队列插件
./rabbitmq-plugins enable rabbitmq_delayed_message_exchange
- 重启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.的时间轮,这些方式各有特点,看需要适用的场景