RabbitMQ延时任务

概念:

消息的TTL(Time To Live)
消息的TTL就是消息的存活时间。RabbitMQ可以对队列和消息分别设置TTL。对队列设置就是队列没有消费者连着的保留时间,也可以对每一个单独的消息做单独的设置。超过了这个时间,我们认为这个消息就死了,称之为死信。
如果队列设置了,消息也设置了,那么会取小的。所以一个消息如果被路由到不同的队列中,这个消息死亡的时间有可能不一样(不同的队列设置)。这里单讲单个消息的TTL,因为它才是实现延迟任务的关键。
可以通过设置消息的expiration字段或者x-message-ttl属性来设置时间,两者是一样的效果。
消息扔到队列中后,过了设置的限定时间,如果没有被消费,它就死了。不会被消费者消费到。这个消息后面的,没有“死掉”的消息对顶上来,被消费者消费。
死信在队列中并不会被删除和释放,它会被统计到队列的消息数中去。单靠死信还不能实现延迟任务,还要靠Dead Letter Exchange。

Dead Letter Exchanges
Exchage的概念在这里就不在赘述,可以从这里进行了解。一个消息在满足如下条件下,会进死信路由,记住这里是路由而不是队列,一个路由可以对应很多队列。
1. 一个消息被Consumer拒收了,并且reject方法的参数里requeue是false。也就是说不会被再次放在队列里,被其他消费者使用。
2. 上面的消息的TTL到了,消息过期了。
3. 队列的长度限制满了。排在前面的消息会被丢弃或者扔到死信路由上。
Dead Letter Exchange其实就是一种普通的exchange,和创建其他exchange没有两样。只是在某一个设置Dead Letter Exchange的队列中有消息过期了,会自动触发消息的转发,发送到Dead Letter Exchange中去。

实现延迟队列:
延迟任务通过消息的TTL和Dead Letter Exchange来实现。我们需要建立2个队列,一个用于发送消息,一个用于消息过期后的转发目标队列。

发布者Code:

        public void PublishDelayMessage<T>(T message, int expireMinutes, bool durable = true) where T : class
        {
            if (expireMinutes <= 0)
            {
                throw new ArgumentException("expireMinutes 必须大于0");
            }

            using (var channel = connection.CreateModel())
            {
                var body = Encoding.UTF8.GetBytes(JsonConvert.SerializeObject(message));

                //创建默认的死信交换机
                string receiveExchangeName = $"{typeof(T).FullName}.DelayReceive";
                string receiveQueueName = $"{receiveExchangeName}.{expireMinutes}";
                channel.ExchangeDeclare(exchange: receiveExchangeName, type: "direct", durable: durable);

                string bufferExchange = $"{typeof(T).FullName}.DelayBuffer";
                string bufferQueueName = $"{bufferExchange}.{expireMinutes}";
                channel.ExchangeDeclare(exchange: bufferExchange, type: "direct", durable: durable);

                //创建消息缓冲队列,在这个队列里面实现消息的过期转发
                var properties = channel.CreateBasicProperties();
                //properties.Expiration = (expireMinutes * 60000).ToString();
                properties.Expiration= (expireMinutes * 6000).ToString();
                Dictionary<string, object> arguments = new Dictionary<string, object>();
                arguments.Add("x-dead-letter-exchange", receiveExchangeName);
                arguments.Add("x-dead-letter-routing-key", receiveQueueName);
                channel.QueueDeclare(queue: bufferQueueName, durable: durable, exclusive: false, autoDelete: false, arguments: arguments);
                channel.QueueBind(queue: bufferQueueName, exchange: bufferExchange, routingKey: bufferQueueName);

                //这个队列用于消息在缓冲队列中过期后转发的目标队列
                channel.QueueDeclare(queue: receiveQueueName, durable: durable, exclusive: false, autoDelete: false, arguments: null);
                channel.QueueBind(queue: receiveQueueName, exchange: receiveExchangeName, routingKey: receiveQueueName);

                channel.BasicPublish(exchange: bufferExchange, routingKey: bufferQueueName, basicProperties: properties, body: body);
            }
        }

虽然没贴出全部的代码,但是最核心的已经有了

1,设置消息的过期时间

2.设置缓冲队列,并且在消息过期以后转发到真实的路由中

看Wireshark抓包分析:

1.过期时间

可以看到发布消息的properties里面设置了expiration

2.过期转发

可以看到缓冲队列在声明的时候,设置了arguments

里面配置了x-dead-letter-exchange,x-dead-letter-routing-key

这样缓冲队列里面的消息过期以后,就将消息转发给配置的对应配置的交换机路由。

时间: 2025-01-12 17:18:40

RabbitMQ延时任务的相关文章

spring boot Rabbitmq集成,延时消息队列实现

本篇主要记录Spring boot 集成Rabbitmq,分为两部分, 第一部分为创建普通消息队列, 第二部分为延时消息队列实现: spring boot提供对mq消息队列支持amqp相关包,引入即可: [html] view plain copy <!-- rabbit mq --> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-

基于rabbitMQ 消息延时队列方案 模拟电商超时未支付订单处理场景

前言 传统处理超时订单 采取定时任务轮训数据库订单,并且批量处理.其弊端也是显而易见的:对服务器.数据库性会有很大的要求,并且当处理大量订单起来会很力不从心,而且实时性也不是特别好 当然传统的手法还可以再优化一下,即存入订单的时候就算出订单的过期时间插入数据库,设置定时任务查询数据库的时候就只需要查询过期了的订单,然后再做其他的业务操作 jdk延迟队列 DelayQueue 采取jdk自带的延迟队列能很好的优化传统的处理方案,但是该方案的弊.端也是非常致命的,所有的消息数据都是存于内存之中,一旦

【RabbitMQ】一文带你搞定RabbitMQ延迟队列

本文口味:鱼香肉丝? ?预计阅读:10分钟 一.说明 在上一篇中,介绍了RabbitMQ中的死信队列是什么,何时使用以及如何使用RabbitMQ的死信队列.相信通过上一篇的学习,对于死信队列已经有了更多的了解,这一篇的内容也跟死信队列息息相关,如果你还不了解死信队列,那么建议你先进行上一篇文章的阅读. 这一篇里,我们将继续介绍RabbitMQ的高级特性,通过本篇的学习,你将收获: 什么是延时队列 延时队列使用场景 RabbitMQ中的TTL 如何利用RabbitMQ来实现延时队列 二.本文大纲

RabbitMQ队列延迟

RabbitMQ队列延迟 1. 场景: “订单下单成功后,15分钟未支付自动取消”   1.传统处理超时订单     采取定时任务轮训数据库订单,并且批量处理.其弊端也是显而易见的:对服务器.数据库性会有很大的要求,     并且当处理大量订单起来会很力不从心,而且实时性也不是特别好.当然传统的手法还可以再优化一下,     即存入订单的时候就算出订单的过期时间插入数据库,设置定时任务查询数据库的时候就只需要查询过期了的订单,     然后再做其他的业务操作 2.rabbitMQ延时队列方案  

CenterOS - CenterOS下安装RabbitMQ

CenterOS下安装RabbitMQ 下载erlang wget https://bintray.com/rabbitmq-erlang/rpm/download_file?file_path=erlang%2F21%2Fel%2F6%2Fx86_64%2Ferlang-21.3.8.14-1.el6.x86_64.rpm  安装erlang rpm -ivh download_file\?file_path\=erlang%2F21%2Fel%2F6%2Fx86_64%2Ferlang-21

RabbitMQ:伪延时队列

目录 一.什么是延时队列 二.RabbitMQ实现 三. 延时队列的问题 四.解决RabbitMQ的伪延时方案 ps:伪延时队列先卖个关子,我们先了解下延时队列. 一.什么是延时队列 所谓延时队列是指消息push到队列后,监听的消费者不能第一时间获取消息,需要等到指定时间才能消费. 一般在业务里面需要对某些消息做定时发送,不想走定时任务或者是用户下单之后多长时间自动失效类似的场景可以考虑通过延时队列实现. 二.RabbitMQ实现 MQ本身并不支持直接的延时队列实现,但是我们可以通过Rabbit

springboot使用RabbitMQ实现延时任务

延时队列顾名思义,即放置在该队列里面的消息是不需要立即消费的,而是等待一段时间之后取出消费.那么,为什么需要延迟消费呢?我们来看以下的场景 订单业务: 在电商/点餐中,都有下单后 30 分钟内没有付款,就自动取消订单.短信通知: 下单成功后 60s 之后给用户发送短信通知.失败重试: 业务操作失败后,间隔一定的时间进行失败重试. 本文基于springboot,使用rabbitmq_delayed_message_exchange插件实现延时队列(RabbitMQ及其插件环境安装点此),具体实践如

RabbitMq 实现延时队列-Springboot版本

rabbitmq本身没有实现延时队列,但是可以通过死信队列机制,自己实现延时队列: 原理:当队列中的消息超时成为死信后,会把消息死信重新发送到配置好的交换机中,然后分发到真实的消费队列: 步骤: 1.创建带有时限的队列 dealLineQueue; 2.创建死信Faout交换机dealLineExchange; 3.创建消费队列realQueue,并和dealLineExchange绑定 4.配置dealLineQueue 的过期时间,消息过期后的死信交换机,重发的routing-key: 以下

RabbitMQ实现延时队列(死信队列)

基于队列和基于消息的TTL TTL是time to live 的简称,顾名思义指的是消息的存活时间.rabbitMq可以从两种维度设置消息过期时间,分别是队列和消息本身. 队列消息过期时间-Per-Queue Message TTL: 通过设置队列的x-message-ttl参数来设置指定队列上消息的存活时间,其值是一个非负整数,单位为微秒.不同队列的过期时间互相之间没有影响,即使是对于同一条消息.队列中的消息存在队列中的时间超过过期时间则成为死信. 死信交换机DLX 队列中的消息在以下三种情况