跳到主要内容
JackSparrow414
返回

RabbitMQ和Spring AMQP学习五(延迟队列)

本文导航

展开本文导航

RabbitMQ和Spring AMQP学习五(延迟队列)

使用场景

对于消息延迟队列的使用场景,网上大多数都是最最常见的场景,即订单超时未支付。实际上使用的场景很多,比如,一个教学系统,老师发作业,老师自定义发作业时间,只有到达定义的发布时间,才真正的进行发作业。这里就可以用到消息队列,即,教师设置发作业时间,将发作业事件推送进延迟消息队列,等待时间到达时,向学生真正的发布作业

使用前准备

选择方案

这里选择第二种方式

安装插件

关于RabbitMq的插件相关文档,可见官方整体插件文档。首先查看当前MQ是否已经安装了对应的插件

rabbitmq-plugins list

如果没有,在上述插件文档中搜索所需要的插件,例如这里搜索delay RabbitMQ 插件页面中的 rabbitmq_delayed_message_exchange 下载入口

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

启动延迟插件

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 声明为 x-delayed-message 的源码

在RabbitAdmin中声明Exchange的时候进行了判断,而标红的地方就是上述Java Client的方式,包装了一下而已。至于RabbitAdmin的声明exchange、Queue的时机见这篇文章

查看控制台

声明延迟队列之后,启动生产者端,在rabbitMQ的控制台查看 RabbitMQ Exchanges 列表显示 delay exchange 的类型为 x-delayed-message

可以看到,这时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 启动日志显示延迟交换机与队列的声明 Spring AMQP 发送日志显示消息头中的 x-delay 参数

可以看到,在使用上述代码之后,发送消息的时候,Spring AMQP自动增加了x-delay头

此时查看控制台如下 RabbitMQ 延迟交换机详情显示一条等待投递的消息 此时,延迟exchange中存在一个延迟消息,等待延迟时间(本文中为10秒)之后,在消费者端控制台收到信息如下 消费者日志显示延迟十秒后收到消息 可以看到,消息是延迟了10秒之后才被消费者接收

示例代码

生产者端代码

参考

  1. Spring AMQP的使用示例文档

  2. RabbitMQ的相关博客


分享本文:

继续阅读本系列

RabbitMQ 与 Spring AMQP

  1. RabbitMQ和Spring AMQP学习一
  2. RabbitMQ和Spring AMQP学习二(消息模型详解、Spring AMQP启动流程分析-超详细)
  3. Rabbit MQ 和Spring AMQP学习三(消息的可靠性)
  4. Rabbit MQ和Spring AMQP学习四(JSON消息体)
  5. RabbitMQ和Spring AMQP学习五(延迟队列)正在阅读

评论

欢迎提问、纠错和分享经验。发表评论需要登录 GitHub;中英文版本共享评论。

评论仅在正式网站开放。