本文导航
展开本文导航
RabbitMQ和Spring AMQP学习五(延迟队列)
使用场景
对于消息延迟队列的使用场景,网上大多数都是最最常见的场景,即订单超时未支付。实际上使用的场景很多,比如,一个教学系统,老师发作业,老师自定义发作业时间,只有到达定义的发布时间,才真正的进行发作业。这里就可以用到消息队列,即,教师设置发作业时间,将发作业事件推送进延迟消息队列,等待时间到达时,向学生真正的发布作业
使用前准备
选择方案
- 关于rabbitMQ的延迟队列的使用,网上大多数的解决方法是TTL+死信队列的实现。即,消息发布时,设置过期时间,当前生产者发送到的队列是没有任何消费者监听的。当消息到达超时时间,会进入死信队列。而这个死信队列是由对应的消费者监听的,此时会进行消费。这样就间接实现了延迟队列.关于具体实现可自行Google,这里随便给一个实现链接
- 第二种方式,使用RabbitMQ的延迟插件实现
这里选择第二种方式
安装插件
关于RabbitMq的插件相关文档,可见官方整体插件文档。首先查看当前MQ是否已经安装了对应的插件
rabbitmq-plugins list
如果没有,在上述插件文档中搜索所需要的插件,例如这里搜索delay

在下面的download中,点击下载.ez的文件。将下载的插件放到rabbitMQ的安装目录下的plugins下,对于不同系统的plugins文件的位置如下

启动延迟插件
rabbitmq-plugins enable rabbitmq_delayed_message_exchange
编码
rabbitMQ对于上述延迟插件的使用博客见这里
主要分为两步步骤,第一,声明一个延迟exchange;第二,发送消息时设置延迟时间
声明延迟exchange
在声明exchange的的时候,加入x-delayed-type参数,rabbitMQ Java Client使用如下
Map<String, Object> args = new HashMap<>();
// value对应的值为交换机的类型,topic,direct,fanout等
args.put("x-delayed-type", "topic");
channel.exchangeDeclare("delayExchange", true, false, args);
对于Spring AMQP,我们不需要这样声明,只需要在声明具体类型的Exchange之后,将delay属性设置为true即可
@Configuration
public class DelayExchangeConfig {
@Bean
public TopicExchange delayTopicExchange() {
TopicExchange result = new TopicExchange("delay");
// 设置该Exchange为延迟交换机
result.setDelayed(true);
return result;
}
}
那么Spring AMQP内部是如何设置的呢?

在RabbitAdmin中声明Exchange的时候进行了判断,而标红的地方就是上述Java Client的方式,包装了一下而已。至于RabbitAdmin的声明exchange、Queue的时机见这篇文章
查看控制台
声明延迟队列之后,启动生产者端,在rabbitMQ的控制台查看

可以看到,这时exchange的类型变成了x-deleayed-message类型
发送消息设置延迟时间
rabbitMq的Java Client代码如下,在消息头中增加x-delay即可,其值对应的是具体的延迟时间
byte[] messageBodyBytes = "delayed payload".getBytes();
AMQP.BasicProperties.Builder props = new AMQP.BasicProperties.Builder();
headers = new HashMap<String, Object>();
headers.put("x-delay", 5000);
props.headers(headers);
channel.basicPublish("delayExchange", "", props.build(), messageBodyBytes);
Spring AMQP使用如下
设置完延迟exchange之后,还需要,在发送时指定消息队列的延迟时间,部分示例代码如下
public void sendDelayExchange() {
Message delayMessage = MessageUtils.buildStringMessage("this is delay message", MessageProperties.CONTENT_TYPE_TEXT_PLAIN);
// 发送消息时,设置消息为延迟消息
delayMessage.getMessageProperties().setDelay(10000);
log.info("发送延迟消息{},时间为{}", new String(delayMessage.getBody()), LocalDateTime.now().toString());
rabbitTemplate.convertAndSend(delayTopicExchange.getName(), "message.delay", delayMessage);
}
简单验证
消费者端代码不再给出,就是声明一个Queue,这个Queue绑定在延迟exchange上,消费者监听这个Queue即可
启动生产者、消费者之后

可以看到,在使用上述代码之后,发送消息的时候,Spring AMQP自动增加了x-delay头
此时查看控制台如下
此时,延迟exchange中存在一个延迟消息,等待延迟时间(本文中为10秒)之后,在消费者端控制台收到信息如下
可以看到,消息是延迟了10秒之后才被消费者接收