C#实现rabbitmq 延迟队列功能

    最近在研究rabbitmq,项目中有这样一个场景:在用户要支付订单的时候,若是超过30分钟未支付,会把订单关掉。固然咱们能够作一个定时任务,每一个一段时间来扫描未支付的订单,若是该订单超过支付时间就关闭,可是在数据量小的时候并无什么大的问题,可是数据量一大轮训数据库的方式就会变得特别耗资源。当面对千万级、上亿级数据量时,自己写入的IO就比较高,致使长时间查询或者根本就查不出来,更别说分库分表之后了。除此以外,还有优先级队列,基于优先级队列的JDK延迟队列,时间轮等方式。但若是系统的架构中自己就有RabbitMQ的话,那么选择RabbitMQ来实现相似的功能也是一种选择。 咱们项目中用到了rabbitmq,能够作一个延迟队列完美的解决这个问题。数据库

     rabbitmq自己不具备延时消息队列的功能,可是能够经过TTL(Time To Live)、DLX(Dead Letter Exchanges)特性实现。其原理给消息设置过时时间,在消息队列上为过时消息指定转发器,这样消息过时后会转发到与指定转发器匹配的队列上,变向实现延时队列。利用rabbitmq的这种特性,应该有了一个大概的思路。、架构

网上搜了一下  rabbitmq-delayed-message-exchange 这个插件也能够实现延迟队列的功能。今天介绍的是如何用C#来实现。函数

首先了解一下TTL和DLX spa

消息的TTL(Time To Live)

消息的TTL就是消息的存活时间。RabbitMQ能够对队列和消息分别设置TTL。对队列设置就是队列没有消费者连着的保留时间,也能够对每个单独的消息作单独的设置。超过了这个时间,咱们认为这个消息就死了,称之为死信。若是队列设置了,消息也设置了,那么会取小的。因此一个消息若是被路由到不一样的队列中,这个消息死亡的时间有可能不同(不一样的队列设置)。这里单讲单个消息的TTL,由于它才是实现延迟任务的关键。插件

Dead Letter Exchanges

Exchage的概念在这里就不在赘述。一个消息在知足以下条件下,会进死信路由,记住这里是路由而不是队列,一个路由能够对应不少队列。code

1. 一个消息被Consumer拒收了,而且reject方法的参数里requeue是false。也就是说不会被再次放在队列里,被其余消费者使用。blog

2. 上面的消息的TTL到了,消息过时了。rabbitmq

3. 队列的长度限制满了。排在前面的消息会被丢弃或者扔到死信路由上。队列

Dead Letter Exchange其实就是一种普通的exchange,和建立其余exchange没有两样。只是在某一个设置Dead Letter Exchange的队列中有消息过时了,会自动触发消息的转发,发送到Dead Letter Exchange中去。资源

 首先我建了两个控制台项目一个是生产者,一个是消费者。

生产者代码以下 

            var factory = new ConnectionFactory() { HostName = "127.0.0.1", UserName = "test", Password = "test" };
            using (var connection = factory.CreateConnection())
            {
                while (Console.ReadLine() != null)
                {
                    using (var channel = connection.CreateModel())
                    {

                        Dictionary<string, object> dic = new Dictionary<string, object>();
                        dic.Add("x-expires", 30000);
                        dic.Add("x-message-ttl", 12000);//队列上消息过时时间,应小于队列过时时间  
                        dic.Add("x-dead-letter-exchange", "exchange-direct");//过时消息转向路由  
                        dic.Add("x-dead-letter-routing-key", "routing-delay");//过时消息转向路由相匹配routingkey  
                        //建立一个名叫"zzhello"的消息队列
                        channel.QueueDeclare(queue: "zzhello",
                            durable: true,
                            exclusive: false,
                            autoDelete: false,
                            arguments: dic);

                        var message = "Hello World!";
                        var body = Encoding.UTF8.GetBytes(message);

                        //向该消息队列发送消息message
                        channel.BasicPublish(exchange: "",
                            routingKey: "zzhello",
                            basicProperties: null,
                            body: body);
                        Console.WriteLine(" [x] Sent {0}", message);
                    }
                }
            }

            Console.ReadKey();

消费者代码以下:

 

 var factory = new ConnectionFactory() { HostName = "127.0.01", UserName = "test", Password = "test" };

            using (var connection = factory.CreateConnection())
            {
                using (var channel = connection.CreateModel())
                {
                    channel.ExchangeDeclare(exchange: "exchange-direct", type: "direct");
                    string name = channel.QueueDeclare().QueueName;
                    channel.QueueBind(queue: name, exchange: "exchange-direct", routingKey: "routing-delay");

                    //回调,当consumer收到消息后会执行该函数
                    var consumer = new EventingBasicConsumer(channel);
                    consumer.Received += (model, ea) =>
                    {
                        var body = ea.Body;
                        var message = Encoding.UTF8.GetString(body);
                        Console.WriteLine(ea.RoutingKey);
                        Console.WriteLine(" [x] Received {0}", message);
                    };

                    //Console.WriteLine("name:" + name);
                    //消费队列"hello"中的消息
                    channel.BasicConsume(queue: name,
                                         autoAck: true,
                                         consumer: consumer);

                    Console.WriteLine(" Press [enter] to exit.");
                    Console.ReadLine();
                }
            }

            Console.ReadKey();

 

  

效果 :

在等待了12秒后消费者等到了消息。

 

 这样咱们就实现了延迟队列的功能了。

相关文章
相关标签/搜索