本文轉自:https://www.cnblogs.com/haoxinyue/p/6613706.html
場景一:物聯網系統經常會遇到向終端下發命令,如果命令一段時間沒有應答,就需要設置成超時。
場景二:訂單下單之后30分鍾后,如果用戶沒有付錢,則系統自動取消訂單。
上述類似的需求是我們經常會遇見的問題。最常用的方法是定期輪訓數據庫,設置狀態。在數據量小的時候並沒有什么大的問題,但是數據量一大輪訓數據庫的方式就會變得特別耗資源。當面對千萬級、上億級數據量時,本身寫入的IO就比較高,導致長時間查詢或者根本就查不出來,更別說分庫分表以后了。除此之外,還有優先級隊列,基於優先級隊列的JDK延遲隊列,時間輪等方式。但如果系統的架構中本身就有RabbitMQ的話,那么選擇RabbitMQ來實現類似的功能也是一種選擇。
使用RabbitMQ來實現延遲任務必須先了解RabbitMQ的兩個概念:消息的TTL和死信Exchange,通過這兩者的組合來實現上述需求。
消息的TTL(Time To Live)
消息的TTL就是消息的存活時間。RabbitMQ可以對隊列和消息分別設置TTL。對隊列設置就是隊列沒有消費者連着的保留時間,也可以對每一個單獨的消息做單獨的設置。超過了這個時間,我們認為這個消息就死了,稱之為死信。如果隊列設置了,消息也設置了,那么會取小的。所以一個消息如果被路由到不同的隊列中,這個消息死亡的時間有可能不一樣(不同的隊列設置)。這里單講單個消息的TTL,因為它才是實現延遲任務的關鍵。
可以通過設置消息的expiration字段或者x-message-ttl屬性來設置時間,兩者是一樣的效果。只是expiration字段是字符串參數,所以要寫個int類型的字符串:
byte[] messageBodyBytes = "Hello, world!".getBytes(); AMQP.BasicProperties properties = new AMQP.BasicProperties(); properties.setExpiration("60000"); channel.basicPublish("my-exchange", "routing-key", properties, messageBodyBytes);
當上面的消息扔到隊列中后,過了60秒,如果沒有被消費,它就死了。不會被消費者消費到。這個消息后面的,沒有“死掉”的消息對頂上來,被消費者消費。死信在隊列中並不會被刪除和釋放,它會被統計到隊列的消息數中去。單靠死信還不能實現延遲任務,還要靠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個隊列,一個用於發送消息,一個用於消息過期后的轉發目標隊列。
生產者輸出消息到Queue1,並且這個消息是設置有有效時間的,比如60s。消息會在Queue1中等待60s,如果沒有消費者收掉的話,它就是被轉發到Queue2,Queue2有消費者,收到,處理延遲任務。
具體實現步驟如下:
第一步, 首先需要創建2個隊列。Queue1和Queue2。Queue1是一個消息緩沖隊列,在這個隊列里面實現消息的過期轉發。如下圖,設置Dead letter exchange和Dead letter routing key。設置這兩個屬性就是當消息在這個隊列中expire后,采用哪個路由發送。這個dlx的exchange需要事先創建好,就是一個普通的exchange。由於我們還需要向Queue1發送消息,那么還需要創建一個exchange,並且和Queue1綁定。例子中,exchange同樣取名:queue1。
我們還需要建一個Queue2,這個隊列用於消息在Queue1中過期后轉發的目標隊列。所以這個Queue2隊列建好以后,需要綁定Queue1設置的死信路由:dlx。完成Queue2的綁定以后,環境就搭建完成了。
第二步,實現消息的Producer。由於我們的目的是讓進入Queue1的消息過期,然后自動轉送到Queue2中,所以發送的時候,需要設置過期時間。
ConnectionFactory factory = new ConnectionFactory(); factory.setUsername("bsp"); factory.setPassword("123456"); factory.setVirtualHost("/"); factory.setHost("10.23.22.42"); factory.setPort(5672); conn = factory.newConnection(); channel = conn.createChannel(); byte[] messageBodyBytes = "Hello, world!".getBytes(); byte i = 10; while (i-- > 0) { channel.basicPublish("queue1", "queue1", new AMQP.BasicProperties.Builder().expiration(String.valueOf(i * 1000)).build(), new byte[] { i }); }
上面的代碼我模擬了1-10號消息,消息的內容里面是1-10。過期的時間是10-1秒。這里要注意,雖然10是第一個發送,但是它過期的時間最長。
第三步,實現消息的Consumer。Consumer就是延遲任務的具體實施者。由於具體的任務往往是一個比較耗時的任務,所以一般來說,任務一般在異步線程中執行。
ConnectionFactory factory = new ConnectionFactory(); factory.setUsername("bsp"); factory.setPassword("123456"); factory.setVirtualHost("/"); factory.setHost("10.23.22.42"); factory.setPort(5672); conn = factory.newConnection(); channel = conn.createChannel(); channel.basicConsume("queue2", true, "consumer", new DefaultConsumer(channel) { @Override public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties,byte[] body) throws IOException { long deliveryTag = envelope.getDeliveryTag(); //do some work async System.out.println(body[0]); } });
運行后如上面的程序,過了10s以后,消費者開始收到數據,但是它是一次性收到如下結果:
10、9 、8 、7 、6、5 、4 、3 、2 、1
Consumer第一個收到的還是10。雖然10是第一個放進隊列,但是它的過期時間最長。所以由此可見,即使一個消息比在同一隊列中的其他消息提前過期,提前過期的也不會優先進入死信隊列,它們還是按照入庫的順序讓消費者消費。如果第一進去的消息過期時間是1小時,那么死信隊列的消費者也許等1小時才能收到第一個消息。參考官方文檔發現“Only when expired messages reach the head of a queue will they actually be discarded (or dead-lettered).”只有當過期的消息到了隊列的頂端(隊首),才會被真正的丟棄或者進入死信隊列。
所以在考慮使用RabbitMQ來實現延遲任務隊列的時候,需要確保業務上每個任務的延遲時間是一致的。如果遇到不同的任務類型需要不同的延時的話,需要為每一種不同延遲時間的消息建立單獨的消息隊列。