文章目錄
- 1. 兩種消費思路
- 2. 確保消費成功兩種思路
- 3. 訊息拒絕
- 4. 訊息確認
- 4.1 自動確認
- 4.2 手動確認
- 4.2.1 推模式手動確認
- 4.2.2 拉模式手動確認
- 5. 冪等性問題
- 6. 小結
前面一篇文章松哥和大家聊了 MQ 高可用之如何確保訊息成功發送,各種配置齊上陣,最終確保了訊息的成功發送,甚至在一些極端情況下還可能發生同一條訊息重復發送的情況,不管怎么樣,訊息總算發送出去了,如果小伙伴們還沒看過上篇文章,建議先看看,再來學習本文:
- 四種策略確保 RabbitMQ 訊息發送可靠性!你用哪種?
今天我們就來聊一聊訊息消費的問題,看看如何確保訊息消費成功,并且確保冪等性,
1. 兩種消費思路
RabbitMQ 的訊息消費,整體上來說有兩種不同的思路:
- 推(push):MQ 主動將訊息推送給消費者,這種方式需要消費者設定一個緩沖區去快取訊息,對于消費者而言,記憶體中總是有一堆需要處理的訊息,所以這種方式的效率比較高,這也是目前大多數應用采用的消費方式,
- 拉(pull):消費者主動從 MQ 拉取訊息,這種方式效率并不是很高,不過有的時候如果服務端需要批量拉取訊息,倒是可以采用這種方式,
兩種方式我都舉個例子看下,
先來看推(push):
這種方式大家比較常見,就是通過 @RabbitListener 注解去標記消費者,如下:
@Component
public class ConsumerDemo {
@RabbitListener(queues = RabbitConfig.JAVABOY_QUEUE_NAME)
public void handle(String msg) {
System.out.println("msg = " + msg);
}
}
當監聽的佇列中有訊息時,就會觸發該方法,
再來看拉(pull):
@Test
public void test01() throws UnsupportedEncodingException {
Object o = rabbitTemplate.receiveAndConvert(RabbitConfig.JAVABOY_QUEUE_NAME);
System.out.println("o = " + new String(((byte[]) o),"UTF-8"));
}
呼叫 receiveAndConvert 方法,方法引數為佇列名稱,方法執行完成后,會從 MQ 上拉取一條訊息下來,如果該方法回傳值為 null,表示該佇列上沒有訊息了,receiveAndConvert 方法有一個多載方法,可以在多載方法中傳入一個等待超時時間,例如 3 秒,此時,假設佇列中沒有訊息了,則 receiveAndConvert 方法會阻塞 3 秒,3 秒內如果佇列中有了新訊息就回傳,3 秒后如果佇列中還是沒有新訊息,就回傳 null,這個等待超時時間要是不設定的話,默認為 0,
這是訊息兩種不同的消費模式,
如果需要從訊息佇列中持續獲得訊息,就可以使用推模式;如果只是單純的消費一條訊息,則使用拉模式即可,切忌將拉模式放到一個死回圈中,變相的訂閱訊息,這會嚴重影響 RabbitMQ 的性能,
2. 確保消費成功兩種思路
在上篇文章中,我們想盡辦法確保訊息能夠發送成功,對于訊息消費成功,其實官方提供了相關的機制,我們一起來看下,
為了保證訊息能夠可靠的到達訊息消費者,RabbitMQ 中提供了訊息消費確認機制,當消費者去消費訊息的時候,可以通過指定 autoAck 引數來表示訊息消費的確認方式,
- 當 autoAck 為 false 的時候,此時即使消費者已經收到訊息了,RabbitMQ 也不會立馬將訊息移除,而是等待消費者顯式的回復確認信號后,才會將訊息打上洗掉標記,然后再洗掉,
- 當 autoAck 為 true 的時候,此時訊息消費者就會自動把發送出去的訊息設定為確認,然后將訊息移除(從記憶體或者磁盤中),即使這些訊息并沒有到達消費者,
我們來看一張圖:

如上圖所示,在 RabbitMQ 的 web 管理頁面:
- Ready 表示待消費的訊息數量,
- Unacked 表示已經發送給消費者但是還沒收到消費者 ack 的訊息數量,
這是我們可以從 UI 層面觀察訊息的消費情況確認情況,
當我們將 autoAck 設定為 false 的時候,對于 RabbitMQ 而言,消費分成了兩個部分:
- 待消費的訊息
- 已經投遞給消費者,但是還沒有被消費者確認的訊息
換句話說,當設定 autoAck 為 false 的時候,消費者就變得非常從容了,它將有足夠的時間去處理這條訊息,當訊息正常處理完成后,再手動 ack,此時 RabbitMQ 才會認為這條訊息消費成功了,如果 RabbitMQ 一直沒有收到客戶端的反饋,并且此時客戶端也已經斷開連接了,那么 RabbitMQ 就會將剛剛的訊息重新放回佇列中,等待下一次被消費,
綜上所述,確保訊息被成功消費,無非就是手動 Ack 或者自動 Ack,無他,當然,無論這兩種中的哪一種,最終都有可能導致訊息被重復消費,所以一般來說我們還需要在處理訊息時,解決冪等性問題,
3. 訊息拒絕
當客戶端收到訊息時,可以選擇消費這條訊息,也可以選擇拒絕這條訊息,我們來看下拒絕的方式:
@Component
public class ConsumerDemo {
@RabbitListener(queues = RabbitConfig.JAVABOY_QUEUE_NAME)
public void handle(Channel channel, Message message) {
//獲取訊息編號
long deliveryTag = message.getMessageProperties().getDeliveryTag();
try {
//拒絕訊息
channel.basicReject(deliveryTag, true);
} catch (IOException e) {
e.printStackTrace();
}
}
}
消費者收到訊息之后,可以選擇拒絕消費該條訊息,拒絕的步驟分兩步:
- 獲取訊息編號 deliveryTag,
- 呼叫 basicReject 方法拒絕訊息,
呼叫 basicReject 方法時,第二個引數是 requeue,即是否重新入隊,如果第二個引數為 true,則這條被拒絕的訊息會重新進入到訊息佇列中,等待下一次被消費;如果第二個引數為 false,則這條被拒絕的訊息就會被丟掉,不會有新的消費者去消費它了,
需要注意的是,basicReject 方法一次只能拒絕一條訊息,
4. 訊息確認
訊息確認分為自動確認和手動確認,我們分別來看,
4.1 自動確認
先來看看自動確認,在 Spring Boot 中,默認情況下,訊息消費就是自動確認的,
我們來看如下一個訊息消費方法:
@Component
public class ConsumerDemo {
@RabbitListener(queues = RabbitConfig.JAVABOY_QUEUE_NAME)
public void handle2(String msg) {
System.out.println("msg = " + msg);
int i = 1 / 0;
}
}
通過 @Componet 注解將當前類注入到 Spring 容器中,然后通過 @RabbitListener 注解來標記一個訊息消費方法,默認情況下,訊息消費方法自帶事務,即如果該方法在執行程序中拋出例外,那么被消費的訊息會重新回到佇列中等待下一次被消費,如果該方法正常執行完沒有拋出例外,則這條訊息就算是被消費了,
4.2 手動確認
手動確認我又把它分為兩種:推模式手動確認與拉模式手動確認,
4.2.1 推模式手動確認
要開啟手動確認,需要我們首先關閉自動確認,關閉方式如下:
spring.rabbitmq.listener.simple.acknowledge-mode=manual
這個配置表示將訊息的確認模式改為手動確認,
接下來我們來看下消費者中的代碼:
@RabbitListener(queues = RabbitConfig.JAVABOY_QUEUE_NAME)
public void handle3(Message message,Channel channel) {
long deliveryTag = message.getMessageProperties().getDeliveryTag();
try {
//訊息消費的代碼寫到這里
String s = new String(message.getBody());
System.out.println("s = " + s);
//消費完成后,手動 ack
channel.basicAck(deliveryTag, false);
} catch (Exception e) {
//手動 nack
try {
channel.basicNack(deliveryTag, false, true);
} catch (IOException ex) {
ex.printStackTrace();
}
}
}
將消費者要做的事情放到一個 try..catch 代碼塊中,
如果訊息正常消費成功,則執行 basicAck 完成確認,
如果訊息消費失敗,則執行 basicNack 方法,告訴 RabbitMQ 訊息消費失敗,
這里涉及到兩個方法:
- basicAck:這個是手動確認訊息已經成功消費,該方法有兩個引數:第一個引數表示訊息的 id;第二個引數 multiple 如果為 false,表示僅確認當前訊息消費成功,如果為 true,則表示當前訊息之前所有未被當前消費者確認的訊息都消費成功,
- basicNack:這個是告訴 RabbitMQ 當前訊息未被成功消費,該方法有三個引數:第一個引數表示訊息的 id;第二個引數 multiple 如果為 false,表示僅拒絕當前訊息的消費,如果為 true,則表示拒絕當前訊息之前所有未被當前消費者確認的訊息;第三個引數 requeue 含義和前面所說的一樣,被拒絕的訊息是否重新入隊,
當 basicNack 中最后一個引數設定為 false 的時候,還涉及到一個死信佇列的問題,這個松哥以后再專門寫文章和大家細聊,
4.2.2 拉模式手動確認
拉模式手動 ack 比較麻煩一些,在 Spring 中封裝的 RabbitTemplate 中并未找到對應的方法,所以我們得用原生的辦法,如下:
public void receive2() {
Channel channel = rabbitTemplate.getConnectionFactory().createConnection().createChannel(true);
long deliveryTag = 0L;
try {
GetResponse getResponse = channel.basicGet(RabbitConfig.JAVABOY_QUEUE_NAME, false);
deliveryTag = getResponse.getEnvelope().getDeliveryTag();
System.out.println("o = " + new String((getResponse.getBody()), "UTF-8"));
channel.basicAck(deliveryTag, false);
} catch (IOException e) {
try {
channel.basicNack(deliveryTag, false, true);
} catch (IOException ex) {
ex.printStackTrace();
}
}
}
這里涉及到的 basicAck 和 basicNack 方法跟前面的一樣,我就不再贅述,
5. 冪等性問題
最后我們再來說說訊息的冪等性問題,
大家設想下面一個場景:
消費者在消費完一條訊息后,向 RabbitMQ 發送一個 ack 確認,此時由于網路斷開或者其他原因導致 RabbitMQ 并沒有收到這個 ack,那么此時 RabbitMQ 并不會將該條訊息洗掉,當重新建立起連接后,消費者還是會再次收到該條訊息,這就造成了訊息的重復消費,同時,由于類似的原因,訊息在發送的時候,同一條訊息也可能會發送兩次(參見四種策略確保 RabbitMQ 訊息發送可靠性!你用哪種?),種種原因導致我們在消費訊息時,一定要處理好冪等性問題,
冪等性問題的處理倒也不難,基本上都是從業務上來處理,我來大概說說思路,
采用 Redis,在消費者消費訊息之前,現將訊息的 id 放到 Redis 中,存盤方式如下:
- id-0(正在執行業務)
- id-1(執行業務成功)
如果 ack 失敗,在 RabbitMQ 將訊息交給其他的消費者時,先執行 setnx,如果 key 已經存在(說明之前有人消費過該訊息),獲取他的值,如果是 0,當前消費者就什么都不做,如果是 1,直接 ack,
極端情況:第一個消費者在執行業務時,出現了死鎖,在 setnx 的基礎上,再給 key 設定一個生存時間,生產者,發送訊息時,指定 messageId,
當然這只是一個簡單思路供大家參考,
松哥在 vhr 專案中也處理了訊息冪等性問題,感興趣的小伙伴可以查看 vhr 原始碼(https://github.com/lenve/vhr),代碼在 mailserver 中,
6. 小結
好啦,今天就和小伙伴們聊了下 RabbitMQ 中和訊息消費相關的幾個話題,感興趣的小伙伴可以實踐下哦~
復制文章標題并在公眾號后臺回復,可以下載本文案例~
轉載請註明出處,本文鏈接:https://www.uj5u.com/qita/298625.html
標籤:其他
下一篇:3.1.4、MySQL__資料庫分組,拼接查詢,日期函式,日期加減,間隔,數值四舍五入,排序,分組,having篩選,分組TopN,流程控制函式,
