Kafka Consumer Lag 排查與 Offset 復原實戰

🌏 Read this article in English

本文提供在 Apache Kafka 集群中,當 Consumer Lag 累積導致系統延遲時,透過 CLI 工具安全重置 Offset 並建立停損機制的具體操作步驟。目標是調整消費邊界以恢復服務可用性。

復現環境條件:

  • Kafka Broker: 2.8+ (CLI 語法基於 kafka-consumer-groups.sh)
  • CLI 可用性: 作業主機須能執行 kafka-consumer-groups.sh,並由套件版本或既有部署清單核對其版本與目標 Broker 相容。若指令不存在、無法執行或版本不符,停止操作。
  • 目標集群與身分驗證: --bootstrap-server 須指向本次核准操作的集群,SSL、SASL 或其他既有驗證設定須由同一操作身分載入。使用相同連線參數執行 --describe --state,只有在成功取得指定 Group 狀態、且未出現連線逾時或驗證錯誤時才能繼續。
  • Offset 操作權限: 操作身分須具備查詢 Consumer Group、讀取目標 Topic 中繼資料及變更 committed offset 的權限。先由既有權限清單或核准紀錄確認,再觀察 --describe --state--describe 與 dry-run 是否出現授權錯誤;任一檢查失敗即停止,不以其他身分繞過。
  • 復原檔目錄: 匯出目錄須存在,且同一操作身分能建立新的、不覆寫既有內容的檔案。若無法建立檔案、空間不足或目標檔名已存在,停止重置並保留既有紀錄。
  • Consumer Group: 處於非活躍狀態 (EmptyDead#MEMBERS 為 0)
  • 目標操作: 先以 --export 匯出保存原始邊界,再透過 --from-file 搭配 --dry-run--execute 完成指定 Topic 的 Offset 重置或回退,並通過四大成功判準。

以上任一項連線、身分驗證、授權、CLI 或檔案寫入檢查未通過,都不進入 Offset 寫入流程。

Lag 的本質是 Partition 中最新訊息的 Offset 與 Consumer Group 已提交的 Offset 之間的差距。這個數值的變化揭示了系統的動態:若 Lag 持續增加,代表生產者速度大於消費者速度,或是消費者發生了阻塞;若 Lag 減少但速度慢,則代表消費者效能瓶頸仍在;若 Lag 出現跳變或歸零,通常與 Rebalance、Offset 重置或量測口徑更新有關。

本文旨在提供一套結構化的排查框架,協助在確保資料邊界清晰的前提下,安全地調整 Offset 以恢復系統平衡。透過標準化工具與決策矩陣,我們可以在緊急情況下快速降低風險,同時明確記錄每一次操作的邊界條件與復原路徑。

排查前置作業與停損機制

在執行任何重置 Offset 的操作之前,首要任務是釐清 Lag 的成因與當前狀態。未經驗證的 Offset 調整會直接改變消費邊界,引發資料重複或遺失的風險,因此必須建立清晰的停損點(Stop Condition)與觀察指標。

Definition: Consumer Lag refers to the difference between the latest offset in a partition and the current committed offset of a consumer group.

量化業務影響:如何判斷緊急程度?

在決定是否需要進行 Offset Reset 之前,必須先量化 Lag 對業務造成的實際延遲時間。這對於 PM 與決策者評估「資料即時性」至關重要。

計算公式:

$$\text{Estimated Delay (Time)} = \frac{\text{Current Lag (Messages)}}{\text{Consumption Rate (Messages/Second)}}$$

上述公式僅適用於「無新增流量或生產速率已知且穩定」的情境,代表的是清除現有 backlog 所需的時間。若生產者持續輸入訊息,實際清空時間應使用 Lag ÷ (Consumption Rate - Production Rate) 計算。

例如,若當前 Lag 為 1,000,000 則訊息,而消費者處理速率為 5,000 msg/s,且假設生產速率相對穩定,則代表系統延遲已達 200 秒。如果該業務場景(如支付確認)要求延遲必須小於 5 秒,則這是一個需要立即介入的高優先級事件;若僅是日誌收集,則可能只需觀察趨勢即可。

注意:在群組非活躍時 Lag 不增加,僅表示 Log End Offset 未前進或監控未更新,不能推論生產與消費已達平衡,因為此時消費速率為零。真正的延遲應由訊息時間戳(Timestamp)量測,而非單純依賴 Lag 數值。

確認 Consumer Group 狀態與邊界

在執行寫入或重置操作前,確認 Consumer Group 處於非活躍狀態是降低風險的關鍵實踐。這不僅是為了避免資料競爭,也是 Kafka 官方建議的安全實踐。如果 Consumer 仍在運行,調整 Offset 可能會導致訊息重複消費或丟失,具體取決於 auto.offset.reset 設定與 Consumer 的內部狀態機。

操作步驟:

  1. 檢查 Group 狀態與成員數量:使用 --describe --state 參數查詢 Consumer Group 的精確狀態。若 STATEEmptyDead,且 #MEMBERS 為 0(無任何 active members),方可進行後續作業。
  2. 停損條件:若查詢顯示狀態為 StablePreparingRebalance,或者 #MEMBERS 大於 0,必須立即停止操作。此時應先停止 Consumer 實例或等待 Rebalance 結束,切勿強行重置。

注意--to-offset 在 Kafka CLI 中僅接受單一 Long 整數,並會套用到所選 Topic 的所有 Partition;若需要針對各 Partition 指定不同的 Offset(例如 Partition 0 重置到 1500、Partition 1 重置到 450),--to-offset 1500,450 並非合法語法。針對逐 Partition 差異化重置或精確復原,必須先使用 --export 匯出備份 CSV,再透過 --from-file 載入設定。

決策矩陣:擴容、修復還是重置?

面對 Lag,通常有三種主要的處置方向:水平擴展程式碼或配置優化、以及 Offset 重置。這三者並非互斥,但在緊急情況下需要優先級判斷。

決策的關鍵不在於「消除 Lag」,而在於「在業務容忍度與工程成本之間尋找平衡點」。若業務允許短暫的資料延遲,且重新部署 Consumer 耗時過長,Offset Reset 往往是恢復系統可用性的快速途徑。然而,這必須伴隨嚴格的驗證與記錄機制,且需明確:重置僅是事件緩解手段,系統的根本平衡仍取決於 Consumer 根因是否已排除。

策略適用情境風險備援方案
水平擴展Lag 穩定增加,硬體資源充足,無程式碼瓶頸低(僅涉及資源分配)恢復原部署規模或調整 Rebalance 策略
程式碼優化Consumer 處理邏輯複雜、I/O 阻塞或配置不當中(需重新部署與測試)回滾至上一個穩定版本
Offset 重置需要快速恢復服務,且允許部分歷史資料被略過或重播高(可能導致資料不一致或重複處理)保留原始 Offset 記錄,隨時可手動復原

Offset Reset 是一種以犧牲資料完整性為代價換取系統可用性的 trade-off。

  • To-Latest: 略過未處理的訊息,應用層需容忍遺失。
  • To-Earliest/Specific: 重播歷史訊息,應用層需具備冪等性。

Kafka Log 中的資料不會因此被刪除。

重要限制:Kafka 支援增加 Topic 的 partition 數,但不支援直接縮減既有 Topic 的 partition 數。在緊急復原情境下,建立新 Topic 並遷移資料通常比線上擴增 Partition 更具可預測性與較低的 Rebalance 風險。

安全重置 Offset 的實戰步驟

本節將詳細說明如何使用 Kafka CLI 工具進行 Offset 的安全重置與回退。我們將遵循「先檢查狀態、預覽模擬、匯出備份、執行寫入、事後驗證」的原則,確保每一步都在可控範圍內。

步驟一:檢查 Group 狀態與定位目標 Partition

在重置之前,必須先以 --describe --state 驗證 Consumer Group 的活躍狀態(確認 STATE 欄位與 #MEMBERS 數量),並以 --describe 明確知道哪些 Partition 發生了 Lag 以及當前的 CURRENT-OFFSET

操作與驗證:

使用 kafka-consumer-groups.sh 查詢 Consumer Group my-consumer-group 的狀態與細節。

kafka-consumer-groups.sh --bootstrap-server localhost:9092 \
  --group my-consumer-group \
  --describe \
  --state

# 2. 查詢特定 Consumer Group 的 Lag 與 Offset 詳情
kafka-consumer-groups.sh --bootstrap-server localhost:9092 \
  --group my-consumer-group \
  --describe

預期狀態輸出範例(–describe –state):

GROUP              COORDINATOR (ID)         ASSIGNMENT-STRATEGY  STATE           #MEMBERS
my-consumer-group  10.0.1.15:9092 (1001)    -                    Empty           0

停損與驗證:若 --describe --state 輸出為 Stable#MEMBERS 大於 0,表示仍有成員在線上運作,必須停止操作。只有當 STATEEmptyDead#MEMBERS 為 0 時,方可由 --describe 確認各 Partition 的 Offset 詳情:

GROUP              TOPIC      PARTITION  CURRENT-OFFSET  LOG-END-OFFSET  LAG             CONSUMER-ID     HOST            CLIENT-ID
my-consumer-group  my-topic   0          1500            2500            1000            -               -               -
my-consumer-group  my-topic   1          450             3000            2550            -               -               -

核對 CURRENT-OFFSET 與預期目標。若 LAG 異常巨大,請再次確認 Partition 分配狀況與生產者寫入速率。

步驟二:使用 Dry Run 模擬重置結果

Kafka CLI 提供了 --dry-run 參數,這是安全操作的關鍵環節。它不會實際修改任何 Offset,但會顯示如果執行該命令,哪些 Partition 的 Offset 將被更改為指定的值。

Dry Run 範例邏輯:

假設我們決定將 my-consumer-group 的 Offset 重置到最新(Latest),以跳過所有未處理的舊訊息。

# 使用 --dry-run 模擬重置至 Latest
kafka-consumer-groups.sh --bootstrap-server localhost:9092 \
  --group my-consumer-group \
  --reset-offsets \
  --to-latest \
  --topic my-topic \
  --dry-run

預期輸出範例:

The following offsets will be set.
Topic: my-topic Partition: 0 Offset current: 1500, new: 2500 (latest)
Topic: my-topic Partition: 1 Offset current: 450, new: 3000 (latest)

驗證:對照輸出中的 new offset 是否符合預期目標值(例如應為最新 Log End Offset),並確認受影響的 Partition 數量是否與步驟一觀察到的一致。

步驟三:匯出原始 Offset 復原檔(–export)

在執行實際重置之前,必須將每個 Partition 的當前 Offset 匯出為可供 --from-file 讀取的 CSV 復原檔。這份檔案同時是審計紀錄與回退入口。由於 --to-offset 只接受單一數值,無法用於 multi-partition 異值復原,因此匯出 CSV 檔是達成逐 Partition 精確控制的唯一途徑。

注意:對單一 Consumer Group 執行的 --export 匯出檔內容為 topic,partition,offset 三個欄位,不包含 Group ID 欄位;目標 Consumer Group 名稱係由命令列 --group 參數鎖定與確認。

復原檔名稱需包含本次操作的唯一識別值,且目標檔案必須不存在;既有復原紀錄一律保留。<operation-id> 可替換為變更單號或時間戳。set -C 會拒絕覆寫已存在的檔案。

# 匯出 my-consumer-group 目前所有 Topic/Partition 的 committed offset
backup_file="my-consumer-group-offsets-before-reset-<operation-id>.csv"

if [ -e "$backup_file" ]; then
  echo "復原檔已存在,保留既有紀錄並停止此次操作:$backup_file"
  exit 1
fi

set -C
kafka-consumer-groups.sh --bootstrap-server localhost:9092 \
  --group my-consumer-group \
  --reset-offsets \
  --all-topics \
  --to-current \
  --dry-run \
  --export > "$backup_file"
export_status=$?
set +C

if [ "$export_status" -ne 0 ]; then
  echo "匯出失敗;保留未完成檔案供稽核,但不得作為回退來源:$backup_file"
  exit 1
fi

驗證:逐列核對匯出檔中的 topic,partition,offset 數值,確認它們與步驟一 --describe 查得的 CURRENT-OFFSET 完全一致。若有 Partition 缺漏或 Offset 不符,停止重置;保留該檔案,待 Group 狀態再次確認後重新匯出。若寫入失敗或未完成核對,不進入下一步。

步驟四:執行實際重置操作

實際重置的前提是:步驟二的 dry-run 符合預期、步驟三的復原檔已完成逐 Partition 核對,而且 Consumer Group 仍無活躍成員。

執行前再次進行狀態檢查:

kafka-consumer-groups.sh --bootstrap-server localhost:9092 \
  --group my-consumer-group \
  --describe \
  --state

停損條件:若狀態為 StablePreparingRebalance,或 #MEMBERS 大於 0,不執行重置。先停止 Consumer 或等待 Rebalance 完成,再從步驟一重新確認。

只有狀態為 EmptyDead 且無活躍成員時,才執行已完成 dry-run 的重置命令:

# 執行實際重置至 Latest
kafka-consumer-groups.sh --bootstrap-server localhost:9092 \
  --group my-consumer-group \
  --reset-offsets \
  --to-latest \
  --topic my-topic \
  --execute

在啟動 Consumer 前,須以步驟一的 --describe 指令完成逐 Partition 檢查,確認 CURRENT-OFFSET 等於 dry-run 的 new offset。若任一 Partition 不符、Group 狀態改變,或出現 IllegalGenerationException,停止啟動 Consumer,並使用下方的匯出檔回退流程。

完整回退流程(–from-file)

當重置結果不符預期或業務需要撤銷變更時,必須透過先前匯出的 CSV 檔案執行回退。回退過程同樣嚴格遵循五階段防護:

1. 狀態檢查:先確認 Consumer 已停止,由 --describe --state 確認 Group 為 EmptyDead#MEMBERS 為 0。若為 Stable 或仍在 Rebalance,停止回退。

kafka-consumer-groups.sh --bootstrap-server localhost:9092 \
  --group my-consumer-group \
  --describe \
  --state

2. Dry Run 預覽:搭配 --group my-consumer-group 使用 --from-file 預覽將載入的每個 Partition Offset。

kafka-consumer-groups.sh --bootstrap-server localhost:9092 \
  --group my-consumer-group \
  --reset-offsets \
  --from-file "$backup_file" \
  --dry-run

3. 逐 Partition 核對:逐一比對 --dry-run 輸出中的 new offset 與 $backup_file 內的 CSV 數據(topic,partition,offset),確認 Topic、Partition 與 Offset 完全相符,並確認命令列指定的 Group 為目標 my-consumer-group

4. 二次狀態檢查與執行回退:核對無誤後,再次檢查 --describe --state 確認 Group 仍為 EmptyDead,接著執行回退寫入:

kafka-consumer-groups.sh --bootstrap-server localhost:9092 \
  --group my-consumer-group \
  --reset-offsets \
  --from-file "$backup_file" \
  --execute

5. 事後驗證:回退執行完成後,立刻執行 --describe 進行逐 Partition 查驗,確保每個 Partition 的 CURRENT-OFFSET 與備份 CSV 記錄完全一致。此流程僅復原 Kafka 中的 committed offset,重置期間下游已處理或遺失的資料仍須透過應用層冪等機制或補償路徑處置。

常見邊界條件與復原路徑

即使遵循了嚴格的操作步驟,在分散式環境中仍可能面臨非預期的邊界狀態。

常見情境與應對

  1. Consumer Group 並非 Empty(出現 active members):若在執行重置時,有 Consumer 實例正在運行或嘗試 Rebalance,會導致 Offset 更新失敗或引發 IllegalGenerationException

    • 停損條件--describe --state 顯示 StablePreparingRebalance#MEMBERS 大於 0。
    • 復原路徑:立即關閉所有 Consumer 實例,等待 Coordinator 確認 Group 變更為 EmptyDead 後,重新執行狀態驗證。
  2. 目標 Offset 設定錯誤或格式不符:例如將 --to-offset 誤填為多個 Partition 的列表(如 1500,450 導致語法解析錯誤),或重置到超過 Log End Offset 的位置。

    • 停損條件:CLI 工具回傳語法錯誤、重置後 Lag 出現負值,或 Consumer 拋出 OffsetOutOfRangeException
    • 復原路徑:單一指定數值時使用修正後的 --to-offset <Long>;涉及多 Partition 差異化值時,改用步驟三匯出的 CSV 檔,循 --from-file 搭配 --dry-run--execute 重新進行復原。
  3. 資料一致性邊界與冪等設計:重置到 Latest 會略過訊息,重置到 Earliest 或特定舊 Offset 則會重播訊息。

    • 停損條件:下游系統監測到重複訂單或資料缺漏等業務層異常。
    • 復原路徑:重置 CLI 無法復原業務資料狀態,需依據應用層的冪等性(Idempotency)設計處理重複訊息,並以補償任務補充跳過的歷史數據。

驗證與下一步安全擴展

重置 Offset 並非操作的終點。只有同時通過四大成功判準,才能將這次復原視為可稽核的成功。

重新啟動 Consumer 會觸發實際資料處理。啟動前判準是:受影響 Partition 的 CURRENT-OFFSET 已與核准的 dry-run 結果一致、可用的 CSV 復原檔已完成核對、Lag 與錯誤日誌可持續觀察,且 Consumer 停止或隔離路徑與應用層補償流程均可執行。任一條件未成立,就不啟動 Consumer。

四大成功判準

  1. Offset 逐 Partition 一致:執行 --describe 查驗,每個受影響 Partition 的 CURRENT-OFFSET 必須精確等於已核准的 dry-run 目標值(或回退時等於 CSV 復原檔的記錄值)。任何一列 Partition 不符即判定為未完成。
  2. Lag 趨勢符合處置目標:重新啟動 Consumer 後,持續監測 Consumer Group my-consumer-group 的 Group 狀態、各 Partition Offset 與 Lag 變化。若目標是清除積壓以恢復即時消費,Lag 應呈穩定下降或維持在低位。若 Lag 止跌回升、Consumer CPU 滿載,或業務驗證失敗,停止或隔離 Consumer,保存當下的 --describe 輸出、Consumer 日誌與業務指標,不再繼續處理資料。若核准的消費邊界本身有誤,待 Group 回到 EmptyDead 後進入 --from-file 回退流程;若下游業務狀態已被改變,Offset 回退不足以恢復一致性,則維持隔離並進入應用層冪等或資料補償流程。
  3. 錯誤日誌未持續累積:檢查 Consumer 實例日誌,確認沒有持續出現 OffsetOutOfRangeExceptionRebalanceInProgressExceptionIllegalGenerationException 等異常。若錯誤日誌不斷增加,應立即暫停 Consumer。
  4. 業務指標符合既有判準:比對上游生產端與下游消費端的業務指標(如每秒訂單處理數、交易總額對齊度),確認無超出業務容忍界限的資料斷層。若指標異常,應觸發資料補償流程。

下一步安全擴展:建立自動化防護機制

四大判準完全通過後,可將本次經驗轉化為預防機制:

  • 建立基於速率的 Lag 告警:在 Prometheus / Grafana 中針對 kafka_consumergroup_lag 設定導數告警(Rate of Change)。相較於固定閾值,Lag 的長尾暴增能更早預警 Consumer 阻塞。
  • 檢視 auto.offset.reset 配置:確認各微服務在無 Committed Offset 或 Offset 過期時的預設行為(earliest vs latest),避免非預期的自動跳頁。
  • 維運 SOP 知識庫化:將包含備份 CSV 檔、--describe --state 檢查截圖、--dry-run 比對紀錄與業務驗證結果在內的完整變更歷程歸檔,作為未來事件處理的標準作業程序。

當 Partition Offset、Lag 趨勢、錯誤日誌與業務指標四者皆達標時, Offset 重置處置才算正式閉環。下一個安全擴展點是把這四大判準寫入自動化監控與維運 SOP,逐步降低手動維護的 blast radius。

Sources