具代表性的面試主題

Kafka Streams 如何遷移到 Broker 驅動的 Rebalance Protocol?

後端困難
Offer.cc 編輯團隊發佈 更新

題幹

你的 Kafka Streams 應用要從 classic protocol 遷移到 Kafka 4.2 的 Streams Rebalance Protocol。請說明它解決什麼問題、如何遷移、哪些能力暫時不可用,以及如何驗證不會遺失狀態或 offset。

題幹與適用情境

你負責一個有狀態 Kafka Streams 應用,目前使用 classic group protocol。升級到 Kafka 4.2 後,團隊希望使用 broker 驅動的 Streams Rebalance Protocol,減少執行個體加入、離開或故障時的全域協調停頓。面試官要求你提出遷移計畫、相容邊界與失敗回滾。

假設應用執行在 Kafka Streams 4.2.x,已有 changelog 與 repartition topics,不能接受無計畫的完整重建。官方文件說明新協定由 broker 持續計算任務分配,使用獨立 streams group;Kafka 4.2 新叢集預設啟用,但 client 仍需設定 group.protocol=streams

面試官考察點

  • 能否解釋 broker 驅動協調如何減少 client 全域同步點,而不只是背設定名稱。
  • 能否區分新建 streams group、classic group 線上升級與離線遷移的不同可行性。
  • 是否檢查 broker、client 版本及 4.2.0 的 KAFKA-20254 風險。
  • 是否知道保留的是 committed offsets,其他 group metadata 會重建,不能承諾所有執行狀態無損保留。
  • 能否把不支援的 static membership、topology update、regex 等能力轉成上線門檻。

普通回答是「改設定並滾動重啟」。強回答會先判斷版本與功能缺口,再選擇新 group 或維護窗口遷移,記錄 offset、changelog、repartition topics 與恢復時間目標。

回答前需要釐清的問題

  1. 目前 Kafka 與 Streams client 是否至少 4.2?否則不能直接啟用完整協定。
  2. 現有應用是否依賴 static membership、線上 topology update、regex subscription 或 standby/rack-aware assignment?任一依賴都可能阻止遷移。
  3. 能否安排所有執行個體停止並等待 group 為空?官方 4.2 遷移路徑只支援 offline migration。
  4. 目前版本是 4.2.0 還是 4.2.1 以上?4.2.0 的 offline migration 有已知 broker bug,修復在 4.2.1。
  5. 是否允許建立新的 application.id?新 group 可隔離風險,但會重新建立狀態並改變 offset 管理。

這些答案決定方案:不能停機時不能聲稱支援 classic 到 streams 的線上遷移;依賴未支援功能時應保留 classic protocol 或先重構應用。

30 秒回答框架

「我先確認 broker 與 client 都是 4.2.x,並盤點應用是否使用新協定尚不支援的功能。Streams Rebalance Protocol 把任務協調放到 broker,減少 client 全域 barrier,但遷移不是普通滾動發布:官方路徑要求 group 為空後切換 group.protocol=streams,只有 committed offsets 會保留,changelog 與 repartition topics 繼續存在,其餘 metadata 會重建。我會避開 4.2.0,優先 4.2.1 以上;遷移前記錄 offsets 與狀態檢查點,遷移後驗證恢復、延遲與 rebalance 指標,失敗就切回 classic 或用新的 application.id 重建。」

分步深入解答

1. 說明協定改變了什麼

classic Streams group 在 client 計算成員任務分配,成員變化時容易形成全域協調點。新協定把 streams group 的成員中繼資料與任務分配放到 broker,應用透過專用 heartbeat 與 streams group 協調。官方文件將它描述為 broker-driven,並提供獨立的 streams group 狀態與 Admin API。

2. 先做能力盤點

Kafka 4.2 目前協定有明確限制:static membership 不可用;顯著 topology update 需要新 streams group;只提供 sticky task assignor,warmup tasks 與 rack-aware assignment 不可用;pattern subscription 不支援;classic 與 streams 之間沒有線上遷移。把這些列成發布前檢查表,避免遷移後才發現設定被忽略。

3. 選擇遷移路徑

官方 offline 路徑是:停止所有執行個體,等待 session.timeout.ms 到期或明確 leave,使 group 為空;設定 group.protocol=streams;再啟動執行個體。只有 committed offsets 會在 broker 端保留,changelog 與 repartition topics 仍作為普通內部 topic 存在,其他 group metadata 會重新建立。

text
停止全部執行個體
      ↓
確認 streams group 為空並記錄 committed offsets
      ↓
升級 broker/client 到相容版本
      ↓
設定 group.protocol=streams
      ↓
啟動執行個體,觀察恢復與 rebalance 指標

如果業務不能接受維護窗口,保留 classic protocol,或使用新的 application.id 做平行驗證。不要把 classic consumer 的 rolling upgrade 經驗直接套到 Streams Rebalance Protocol。

4. 處理版本風險

Kafka 官方升級指南指出,4.2.0 的 classic 到 streams offline migration 受到 KAFKA-20254 broker-side bug 影響,建議不要在 4.2.0 執行;修復已包含在 4.2.1。面試回答應把 4.2.1 作為最低遷移版本,而不是只說「Kafka 4.2 已支援」。

5. 設計狀態與 offset 驗證

遷移前記錄每個輸入 topic 的 committed offsets、changelog topic 狀態與處理延遲。遷移後驗證:新 group 能從預期 offset 繼續;state store 能從 changelog 恢復;repartition topic 仍存在且分割區數一致;重複處理與遺失記錄符合既定語義。將結果與遷移前基線比較,而不是只看程序是否啟動。

6. 監控與回滾

使用 streams group 專用狀態、rebalance count/rate、恢復時長、處理延遲與錯誤率作為觀測面。若恢復超時或結果驗證失敗,停止新 group,保留 committed offsets 與日誌,回滾設定到 classic;若 classic group 已清空,必須依據備份或新的 application.id 重新規劃,不能假設 group metadata 會自動回來。

高品質示範回答

「我不會把這次遷移當成普通滾動發布。先確認 broker、client 都是 4.2.x,並檢查應用是否使用 static membership、線上 topology update、regex subscription、warmup 或 rack-aware assignment 依賴。由於官方只支援 offline migration,我會選 4.2.1 以上,在維護窗口停止全部執行個體,確認 group 為空,記錄 committed offsets 與 state store 檢查點,再設定 group.protocol=streams 啟動。

「遷移後我會驗證 offsets 能繼續消費、changelog 能恢復 state store、repartition topics 未遺失,並觀察 streams group 狀態、rebalance 指標、恢復時長與業務延遲。只有 committed offsets 會保留,其他 group metadata 會重建;4.2.0 還有 KAFKA-20254 風險,所以不能把 4.2.0 當成安全版本。若驗證失敗,就停止新 group,回到 classic 或使用新的 application.id 重建,並保留證據以便複盤。」

常見錯誤

  • 錯誤表現:把 group.protocol=streams 當成可滾動切換 → 失敗原因:Streams protocol 不支援線上遷移 → 修正方法:安排 group 為空的維護窗口。
  • 錯誤表現:在 4.2.0 直接遷移 → 失敗原因:官方升級指南記錄 KAFKA-20254 → 修正方法:使用包含修復的 4.2.1 以上。
  • 錯誤表現:承諾所有 group 狀態自動保留 → 失敗原因:只有 committed offsets 保留,其餘 metadata 會重建 → 修正方法:分別記錄 offset、state store 與 topic 驗證項。
  • 錯誤表現:忽略 static membership 或 topology update → 失敗原因:新協定目前不支援這些能力 → 修正方法:遷移前做功能盤點,必要時保持 classic。

追問及應對

業務不能停機,能否兩批執行個體滾動切換?

不能把官方 Streams migration 路徑描述成線上滾動切換。若不能停機,先保留 classic,或建立新的 application.id 做雙寫/影子驗證,再由業務層切流承擔狀態重建成本。

committed offsets 保留了,為什麼還要驗證 state store?

offset 只說明下一筆從哪裡讀,不代表本地狀態已完整恢復。changelog replay、版本不相容或處理失敗都可能讓 state store 與 offset 不一致,必須做狀態與業務結果驗證。

4.2.0 已經 GA,為什麼仍要避開?

GA 表示功能發布,不等於每條遷移路徑都沒有已知缺陷。官方升級指南明確指出 offline migration 的 KAFKA-20254,修復在 4.2.1;遷移版本應依據修復版本而非只看 GA 標籤。

現有應用依賴正則訂閱怎麼辦?

新 streams protocol 目前不支援 pattern-based topic subscription。保持 classic,或把主題發現與訂閱邏輯改成顯式列表後再評估遷移,不能只改 group protocol。

如何判斷是否需要新的 application.id?

若必須平行驗證、舊 group 不能安全清空,或需要隔離狀態恢復風險,可以用新的 application.id 建立獨立 group;代價是重新消費、重建 state store 與額外資源,應先估算恢復時間與儲存成本。

參考資料

  • Apache Kafka Streams Rebalance Protocol developer guide。
  • Apache Kafka 4.2 Streams Upgrade Guide。
  • Apache Kafka 4.2.0 Release Announcement。

公開來源

同類題目