可能是雲訊息佇列 Kafka 版用戶端版本過低或者Consumer沒有獨立線程維持心跳。
問題現象
使用ApsaraMQ for Kafka時,消費用戶端頻繁出現Rebalance。
可能原因
可能導致故障的原因包括:
v0.10.2之前版本的用戶端:Consumer沒有獨立線程維持心跳,而是把心跳維持與poll介面耦合在一起。其結果就是,如果使用者消費出現卡頓,就會導致Consumer心跳逾時,引發Rebalance。
v0.10.2及之後版本的用戶端:如果消費時間過慢,超過一定時間(
max.poll.interval.ms設定的值,預設5分鐘)未進行poll拉取訊息,則會導致用戶端主動離開隊列,而引發Rebalance。NOT_COORDINATOR / join group failed:消費者組協調器(Group Coordinator)發生變更,或用戶端心跳逾時後重新尋找 Coordinator 時,會觸發 Rebalance。用戶端日誌中可見
NOT_COORDINATOR或join group failed關鍵詞。MemberIdRequiredException / consumer poll timeout:v0.10.2 之前版本的用戶端,心跳與 poll 介面耦合,消費卡頓導致心跳逾時時可能出現此報錯;v0.10.2 及之後版本中,若消費速度過慢,超過
max.poll.interval.ms配置的時間未進行 poll,用戶端會主動離開消費組,用戶端日誌中可見consumer poll timeout相關關鍵詞。CommitFailedException:消費者處理訊息耗時過長,或未能及時發送心跳,導致被消費者組踢出(觸發 Rebalance)後仍嘗試提交 Offset,引發該異常。
解決方案
首先您需要瞭解以下幾點資訊:
session.timeout.ms:心跳逾時時間(可以由用戶端自行設定)。max.poll.records:每次poll返回的最大訊息數量。v0.10.2之前版本的用戶端:心跳是通過poll介面來實現的,沒有內建的獨立線程。
v0.10.2及之後版本的用戶端:為了防止用戶端長時間不進行消費,Kafka用戶端在v0.10.2及之後的版本中引入了
max.poll.interval.ms配置參數。heartbeat.interval.ms:心跳發送間隔(由用戶端自行設定)。
參考以下說明調整參數值:
session.timeout.ms:v0.10.2之前的版本可適當提高該參數值,需要大於消費一批資料的時間,但不要超過30s,建議設定為25s;而v0.10.2及其之後的版本,保持預設值10s即可。
max.poll.records:降低該參數值,建議遠遠小於
<單個線程每秒消費的條數> * <消費線程的個數> *<max.poll.interval.ms>的積。max.poll.interval.ms:該值要大於
<max.poll.records> / (<單個線程每秒消費的條數> * <消費線程的個數>)的值。heartbeat.interval.ms:建議設定為不超過 10000 ms(10 s),且需小於
session.timeout.ms的三分之一,以確保心跳在會話逾時前能發送至服務端。
如果您使用 Spring Kafka 架構,可以在專案的
application.yml檔案中通過以下屬性配置上述參數:spring: kafka: consumer: properties: session.timeout.ms: 25000 max.poll.interval.ms: 300000 max.poll.records: 50 heartbeat.interval.ms: 8000盡量提高用戶端的消費速度,消費邏輯另起線程進行處理。
減少Group訂閱Topic的數量,一個Group訂閱的Topic最好不要超過5個,建議一個Group只訂閱一個Topic。
將用戶端升級至0.10.2以上版本。
常見問題
如何臨時減少 Rebalance 以加快堆積訊息消費?
在訊息出現堆積的緊急情境下,可參考以下措施快速恢複消費:
確保用戶端版本 ≥ 0.10.2,以使用獨立心跳線程,減少因心跳與 poll 耦合導致的逾時 Rebalance。
適當調大
max.poll.interval.ms,使其大於單次完整消費邏輯的最大耗時,避免用戶端主動離組。降低
max.poll.records值,減少單次拉取的訊息量,確保每次 poll 拉取的訊息能在逾時時間內處理完畢。減少 Group 訂閱的 Topic 數量(建議不超過 5 個,最好一個 Group 只訂閱一個 Topic)。
避免消費線程阻塞,將消費邏輯非同步處理(另起獨立線程執行商務邏輯,主線程只負責 poll 和提交 Offset)。