Kafka 後端面試:Share Group 如何提供佇列語意?
題幹與適用場景
訂單通知、圖片轉碼或帳單計算通常是彼此獨立的工作項目。團隊已將事件寫入 Kafka,但傳統 Consumer Group 讓每個 partition 在同一時間只分配給一個消費者;當消費者數量超過 partition 數量時,額外執行個體無法直接增加並行度。請設計一個方案,讓多個消費者協作領取工作,支援逐筆確認、失敗重試與可觀測的投遞次數,同時說明哪些業務仍應保留 partition 順序。
這裡討論的是 Apache Kafka KIP-932 定義的 Share Group。KIP 將它描述為在一般 topic 上提供協作式消費的 group 類型,不等於把 Kafka 變成具有相同語意的 RabbitMQ。回答必須先確認部署版本、用戶端支援與相關 API 是否允許使用,再決定是否上線。
面試官考察點
- 能否解釋 partition 獨佔分配與 Share Group 協作領取的差異。
- 能否把「確認、釋放、拒絕、鎖定逾時」對應到實際處理狀態。
- 是否知道 Share Group 允許消費者數超過 partition 數,但不會自動保留傳統 key 順序。
- 是否把重複投遞、毒性訊息、處理逾時與最大並行度放進同一套失敗模型。
- 是否會核對用戶端、broker、權限、監控與回滾,而不是只背 KIP 名稱。
普通回答會說「Kafka 也能當佇列」;強回答會指出佇列語意的收益、失去的保證與上線前必須驗證的邊界。
回答前需要澄清的問題
- 工作項目是否真的獨立?若同一訂單的事件必須依 key 順序處理,Share Group 可能改變解法。
- 失敗時是短暫重試、轉人工處理,還是永久拒絕?這決定 release、reject 與死信流程。
- 處理時間的 p99、最大並行度與可接受的重複副作用是多少?它們決定 acquisition lock 與冪等設計。
- 業務需要 Kafka 交易的端到端語意嗎?不能假設 Share Group 自動繼承既有 Consumer Group 交易方案。
- 現有 topic 是否仍服務廣播或回放消費者?遷移 group 類型不能改變其他 group 的讀取契約。
30 秒回答框架
「我先確認工作是否允許亂序,以及部署的 Kafka 與用戶端是否支援 KIP-932。Share Group 讓多個消費者協作領取同一 topic 的記錄,消費者數可以超過 partition 數,並能逐筆確認、釋放或拒絕;傳統 Consumer Group 更適合保留 partition 內順序與 offset 語意。我要為每筆工作設計冪等鍵,按處理延遲設定鎖定與重試策略,監控取得、確認、釋放、拒絕與逾時,並為毒性訊息設定隔離路徑。若業務依賴 key 順序、交易邊界或用戶端尚未支援,我保留 Consumer Group,先做小流量獨立 topic 驗證,再決定遷移。」
分步驟深入解答
1. 先畫出兩種分配模型
Consumer Group 通常把 partition 分配給成員;同一 partition 在一個 group 內由一個成員讀取,因此並行度受 partition 數量限制。Share Group 則讓成員從訂閱的 topic 中協作領取記錄,一個 partition 可以同時有多個成員處理不同記錄,成員數可以超過 partition 數。這個差異適合獨立工作項目,卻不能直接推導出全域順序。
Consumer Group: partition-0 -> worker-A
partition-1 -> worker-B
extra workers wait for another partition
Share Group: partition-0 records -> worker-A, worker-B, worker-C
each acquired record is locked for one consumer選擇 Share Group 的理由應是「需要彈性領取與逐筆完成」,而不是「partition 太少所以一定要換」。如果同一 customer 的事件必須依序生效,可以繼續使用 Consumer Group,或在業務層建立序列化佇列。
2. 把記錄生命週期寫成狀態機
KIP-932 描述了有時間限制的 acquisition lock。消費者取得記錄後可以 acknowledge 表示成功、release 讓記錄再次投遞、reject 表示不可處理,或什麼都不做並等待鎖定逾時。KIP 中的預設鎖定時間是 30 秒,但部署應以實際 broker 設定為準,不能把預設值當成 SLA。
available -> acquired -> acknowledged
-> released -> available
-> rejected -> terminal or quarantine
-> lock timeout -> available處理函式必須先用業務冪等鍵登記意圖,再執行外部副作用;否則鎖定逾時或用戶端崩潰會造成重複扣款、重複出貨或重複通知。確認成功只能代表這次領取完成,不能替業務系統回滾已提交的副作用。
3. 設計重試、毒性訊息與並行上限
delivery attempt 計數可用來區分暫時故障與不可處理記錄。對網路抖動可 release 並使用退避;對 schema 不相容、必要欄位缺失等確定性錯誤,應 reject 並寫入隔離 topic 或人工佇列。不要無限 release,否則單一毒性訊息會持續消耗處理預算。
鎖定時間應覆蓋正常處理 p99 加上可解釋的抖動餘量;太短會造成並行重複取得,太長會拖慢復原。還要限制每個 partition 的 acquired record 數,配合 worker semaphore、資料庫連線池與外部 API 額度。監控必須同時呈現 active locks、lock timeout、attempt 分布、reject 數與端到端完成延遲。
4. 重新定義順序與重複語意
傳統 Kafka 敘述常把「partition 內有序」誤當成「業務處理有序」。Share Group 允許多個消費者並行領取,同一 key 的完成順序可能與寫入順序不同;釋放、重新投遞與不同批次處理還會放大這種差異。若業務需要順序,必須把 key 級序列化、版本號檢查或狀態機約束寫進應用,而不是只說「Kafka 有序」。
Exactly-once 也不能由 group 類型自動取得。需要逐項核對訊息讀取、業務寫入與確認邊界;外部資料庫或支付系統仍要依賴冪等鍵、去重表或交易型 outbox。對不支援的組合,明確採用 at-least-once 加冪等,而不是聲稱「Share Group 就是 exactly-once」。
5. 規劃遷移與回滾
先確認 broker 版本、用戶端 API、group 類型設定、ACL、監控與維運命令,再在獨立 topic 或小規模工作負載上壓測。用故障注入驗證:worker 在取得後崩潰、處理超過鎖定時間、連續 reject、broker 重啟與 coordinator 切換。記錄每筆記錄的業務鍵、attempt、狀態與時間線。
若舊消費者仍依賴順序或交易,不能直接把同一個 group 原地改成 Share Group。更安全的方式是複製到專用工作 topic,讓新 group 逐步接管;保持舊 group 可回放,達到錯誤率、重複副作用與延遲門檻後再擴大流量。回滾應停止新 group 的領取並讓舊路徑繼續消費尚未遷移的記錄,避免兩條路徑同時執行同一副作用。
高品質示範回答
「我會先問工作項目能否亂序、處理是否冪等,以及現有部署是否支援 KIP-932。Share Group 適合把 topic 中的獨立記錄當作協作工作項目:多個成員可以從同一 partition 領取不同記錄,成員數不再直接受 partition 數限制,且每筆記錄有 acknowledge、release、reject 與鎖定逾時路徑。它改變了傳統 group 的分配與順序直覺,所以同一業務 key 若要求順序,我會繼續使用 Consumer Group 或增加應用層版本檢查。
我會為每筆記錄設定業務冪等鍵,按 p99 處理時間配置 acquisition lock,限制 active locks 與外部依賴並行;暫時錯誤退避 release,確定性壞資料 reject 到隔離路徑,超過 attempt 閾值停止自動重試。監控取得、確認、釋放、拒絕、鎖定逾時、重複副作用與完成延遲。遷移前做版本、ACL、用戶端與故障注入驗證,先用獨立 topic 小流量執行。除非證明訊息讀取、業務寫入與確認共享同一交易邊界,否則我會明確採用 at-least-once 加冪等,不把 Share Group 宣稱為 exactly-once。」
常見錯誤
- 把 Share Group 說成 RabbitMQ 複製品 → 兩者的儲存、回放與管理模型不同 → 只承諾 KIP 明確的協作領取與確認語意。
- 按 partition 數設定最大 worker 數 → Share Group 的並行模型允許多個成員處理同一 partition → 按鎖定、下游容量與端到端延遲限制並行。
- 預設保留 key 順序 → 多成員領取與重試會改變完成順序 → 對需要順序的 key 使用序列化或版本檢查。
- 處理成功後直接確認但沒有冪等 → 崩潰或鎖定逾時會再次投遞 → 先以業務鍵去重,再確認記錄。
- 無限 release 毒性訊息 → 重試會占滿鎖定與下游預算 → 按錯誤類型、attempt 與隔離策略終止自動重試。
- 把預設 30 秒當成保證 → broker 設定與處理延遲可能不同 → 按實際設定與 p99 測試鎖定時間。
追問與應對
如果同一訂單的事件必須嚴格依序怎麼辦?
不要直接採用 Share Group。保留以訂單 ID 分區的 Consumer Group,或讓同一訂單進入應用層序列化狀態機;若必須共享領取,則需要版本號、前置版本檢查與失敗重排,並承認這會增加複雜度。
一個消費者取得記錄後卡住 2 分鐘怎麼辦?
設定略高於正常 p99 的鎖定時間,監控 lock timeout;逾時後允許再次投遞,但業務處理必須冪等。對長工作可拆成可復原步驟或外置 lease,而不是無限延長鎖定並隱藏故障。
如何處理連續五次 schema 錯誤?
把它視為確定性錯誤,達到門檻後 reject 並寫入隔離 topic,保留原始 payload、schema 版本與錯誤原因。修復消費者後再透過受控重播恢復,不能讓主流程無限重試。
現有系統依賴 Kafka 交易,能否直接切換?
先列出交易涵蓋的讀取、處理與寫入邊界,再驗證 Share Group 用戶端與交易 API 的實際支援。若外部副作用不在同一交易內,就採用 outbox、冪等鍵與補償流程;沒有證據時保留 Consumer Group。
如何證明遷移沒有造成重複扣款?
用唯一業務鍵與去重約束記錄每次執行,注入 worker 崩潰、鎖定逾時、重試與回滾故障,比較執行次數與確認次數。只有副作用次數符合業務不變量、延遲與錯誤率門檻,才擴大流量。