最近在研究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秒後消費者等到了消息。
這樣咱們就實現了延遲隊列的功能了。