1. 程式人生 > >C#實現rabbitmq 延遲佇列功能

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 

訊息的TTL(Time To Live)

訊息的TTL就是訊息的存活時間。RabbitMQ可以對佇列和訊息分別設定TTL。對佇列設定就是佇列沒有消費者連著的保留時間,也可以對每一個單獨的訊息做單獨的設定。超過了這個時間,我們認為這個訊息就死了,稱之為死信。如果佇列設定了,訊息也設定了,那麼會取小的。所以一個訊息如果被路由到不同的佇列中,這個訊息死亡的時間有可能不一樣(不同的佇列設定)。這裡單講單個訊息的TTL,因為它才是實現延遲任務的關鍵。

Dead Letter Exchanges

Exchage的概念在這裡就不在贅述。一個訊息在滿足如下條件下,會進死信路由,記住這裡是路由而不是佇列,一個路由可以對應很多佇列。

1. 一個訊息被Consumer拒收了,並且reject方法的引數裡requeue是false。也就是說不會被再次放在佇列裡,被其他消費者使用。

2. 上面的訊息的TTL到了,訊息過期了。

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秒後消費者等到了訊息。

 這樣我們就實現了延遲佇列的功能了。