我的 Kafka consumer 一直被踢出群組、rebalance 停不下來,最後發現是我自己一則訊息處理太久
System Design·2026年7月8日·8 分鐘閱讀

我的 Kafka consumer 一直被踢出群組、rebalance 停不下來,最後發現是我自己一則訊息處理太久

L
Leo Wu

Leo Wu,主力 Go 的後端工程師,做過金流與高併發訊息佇列系統,信奉消費者要又快又冪等。

那天早上我盯著一條一路往上的 lag 曲線

那是一個很普通的週二早上,我泡好咖啡打開 Grafana,看到訂單通知服務的 consumer lag 從平常穩定的幾百,一路爬到十幾萬,而且還在往上。與此同時客服群組炸開了:有客人抱怨同一筆訂單的通知簡訊收到了四、五次,有人甚至截圖給我看,同一句「您的訂單已成立」在三分鐘內連發六封。

我第一個反應是:是不是有人手殘把 consumer 部署成無限重啟?我馬上去看 pod 狀態,結果 pod 好好的,RESTARTS 欄位是 0,記憶體、CPU 都不高,健康檢查全綠。日誌裡也沒有 panic、沒有 OOM。程式明明活得好好的,可是 Kafka 那邊卻一直在做 rebalance,consumer group 的成員名單每隔一兩分鐘就重洗一次。訊息重複、lag 暴衝、rebalance 風暴,三個現象疊在一起,我一開始完全兜不起來。

我一開始想錯的方向

老實說,前半個小時我都在錯的路上鬼打牆。

  • 我先懷疑是網路問題,consumer 跟 broker 之間的連線不穩,導致 heartbeat 送不到、被判定離線。於是我去翻 broker log,找 session timeout 相關的訊息,結果沒有。
  • 接著我懷疑是不是有另一個測試環境的 consumer 用了同一個 group.id,兩邊搶 partition 互相把對方踢出來。這種鳥事我以前遇過。查了一輪 client id,也沒有。
  • 我甚至一度想直接把 consumer 數量加倍,想說「多幾個 instance 分擔一下 lag 就會降了吧」。還好我沒真的這樣做——後來才知道,那只會讓 rebalance 更頻繁、火上加油。

真正讓我開竅的,是我去看了單一則訊息的處理時間。我在 consume 的進入點跟結束點各打了一個 timestamp,撈出來一算,嚇一跳:有一批訊息,單則處理竟然要花三、四十秒,甚至偶爾破一分鐘。

Consumer Lag 曲線:單則訊息處理太久超過 max.poll.interval,被踢出群組觸發 rebalance 風暴,lag 一路飆破 15 萬,調小批次量並讓消費者冪等後歸零
Consumer Lag 曲線:單則訊息處理太久超過 max.poll.interval,被踢出群組觸發 rebalance 風暴,lag 一路飆破 15 萬,調小批次量並讓消費者冪等後歸零

真正的根因:我把 poll loop 卡住了

這裡要先講清楚兩個很多人(包括當時的我)會搞混的東西:heartbeatmax.poll.interval.ms 是兩回事。

  • heartbeat 是背景執行緒在做的事,它固定頻率跟 broker 說「我還活著」。這個跟你的業務邏輯多慢無關,只要背景執行緒沒死,heartbeat 就會照送。所以我去查 session timeout 當然查不到東西。
  • max.poll.interval.ms 管的是完全不同的東西:它規定你「兩次呼叫 poll() 之間最多能隔多久」。也就是說,你這一批訊息如果處理太久,久到超過這個間隔還沒回來呼叫下一次 poll(),Kafka 就認定「這個 consumer 雖然心跳還在,但它卡死了、沒在幹活」,於是把它踢出 group,觸發 rebalance

我的 consumer 就是死在這裡。預設一次 poll() 會抓一批(max.poll.records 預設 500 筆)訊息回來,我在一個迴圈裡逐筆處理。而其中某些訊息會去呼叫一個下游的風控服務,那個服務那陣子剛好變慢,一筆要好幾秒;再加上我在迴圈裡還順手做了一堆同步的 DB 寫入。一批 500 筆,只要平均每筆一點點時間,整批加起來就輕鬆突破 max.poll.interval.ms(我當時是預設的五分鐘,但因為單批太肥、下游又慢,還是爆了)。

於是惡性循環就成形了:

  1. 1.我抓了一大批訊息,慢慢處理,處理到一半超時。
  2. 2.Kafka 把我踢出 group,開始 rebalance。rebalance 期間整個 group 都會停止消費,lag 當然只進不出。
  3. 3.partition 重新分配。我這一批訊息因為還沒處理完、offset 根本還沒 commit,所以重新分配後,這些訊息又被原封不動地重放一次。
  4. 4.通知就這樣重複寄出去了。而新接手的 consumer 一樣抓一大批、一樣處理太久、一樣超時被踢……

看起來像 consumer 一直重啟,其實 process 從頭到尾沒死過。它只是一直被 Kafka 判「出局」,一直在重新入隊。lag 不降反升,因為真正在消費的時間,遠少於花在 rebalance 跟重複處理上的時間。

我怎麼把它救回來

盤清楚之後,修法其實不難,難的是要先接受一件事:`rebalance` 不是 bug,它是我處理太慢的症狀。 我一直想「怎麼讓 rebalance 別再發生」,方向就錯了;正確的問題是「怎麼讓每一批訊息快到不會觸發它」。

我做了幾件事,由治標到治本:

  • 先調小 `max.poll.records`。 這是最快能止血的一招。我把一次抓的量從 500 降到 50,等於每一批的總處理時間直接砍到十分之一,馬上就不再超時。這一步幾分鐘內就讓 rebalance 停了、lag 開始往下掉。
  • 把重活移出 poll loop。 這是真正的治本。原本我在消費迴圈裡同步呼叫慢的下游、又同步寫 DB,等於用 Kafka 的 poll loop 當工作佇列在跑。我改成 poll 出來之後,只做「快速的落地」,把真正耗時的風控呼叫丟給後面的 worker pool 非同步處理。poll loop 要保持又輕又快,它的職責是搬運,不是加工。
  • 謹慎地調 `max.poll.interval.ms`。 我有把它從五分鐘往上拉一點,給尖峰一些緩衝。但我很克制,因為這個值調太大有副作用:萬一真的有 consumer 卡死,Kafka 要等更久才會把它踢掉、才會重新分配 partition,故障恢復會變慢。所以這只是保險,不是主力。
  • 確保消費者冪等,能承受重放。 這是最重要、也是我最該早點做的。Kafka 是 at-least-once,只要有 rebalance、有沒 commit 的 offset,訊息就一定會被重放,這是機制,不是意外。我在通知服務裡加了一張去重表,用訂單 ID 加通知類型當唯一鍵,處理前先檢查、送出後才寫入。這樣就算同一則訊息被重放十次,簡訊也只會發一次。

順帶一提,我也重新確認了 offset 的 commit 時機——確保是在訊息真正處理完成之後才 commit,而不是 poll 完就先 commit。commit 太早會變成 at-most-once,訊息可能整批遺失,那是另一種災難。

我學到的教訓

這次事故對我最大的衝擊,是它逼我認清一件事:訊息佇列的 at-least-once 和 rebalance 機制,不是來找你麻煩的,是來逼你把消費者設計成「又快又冪等」的。

  • consumer 慢,代價不是「處理得慢一點」而已,而是會被踢出群組、拖垮整個 group、還把訊息重放給你。慢會被機制放大成災難。
  • poll loop 是命脈,不要在裡面做任何長時間阻塞的事。慢的下游、大量 DB、外部 API,全都該推到別的 worker 去。
  • 冪等不是「有空再加的優化」,是消費者的基本義務。因為重放不是「如果」,是「一定會」。
  • 遇到 rebalance 風暴,先別急著加機器、也別急著調 timeout,先去量單則、單批訊息的處理時間。答案幾乎都藏在那裡。

小結

那天中午 lag 歸零、rebalance 停下來的時候,我看著那條終於掉回底部的曲線,心情很複雜——因為根因說到底就是「我自己寫的消費者太慢、又不冪等」。Kafka 沒做錯任何事,它只是忠實地執行了它的規則。

如果要我把這次的血淚濃縮成一句話:別把 poll loop 當成你的工作區,把它當成一條不能塞住的輸送帶;訊息一定會重來,所以你的消費者從第一天起就該假設自己會被重放。 想清楚這兩件事,rebalance 就從一場惡夢,變回它本來該有的樣子——一個中性、可靠、幫你重新分配工作的機制而已。

#Kafka#訊息佇列#rebalance#消費者#事故覆盤#冪等

留言討論

有想法、有不同經驗、或想糾正我?歡迎在下面留言,免註冊,填個暱稱就能留。

相關文章