題幹與適用場景
這道資料工程與串流處理題面向資料工程師、串流平台工程師和後端資料基礎設施職位。約束是消費者可能崩潰、發生 rebalance,生產者可能在回應遺失後重試;結果先寫回 Kafka topic。需要區分 Kafka 內部的原子處理與任意外部副作用,不能把「恰好一次」當成整個系統的魔法保證。
面試官考察點
- 能否把 exactly-once 拆成輸出紀錄可見性與輸入 offset 的原子提交。
- 能否正確使用交易式生產者、
transactional.id、read_committed和手動提交 offset。 - 能否解釋 abort、重啟、fencing、rebalance 後為什麼需要回退並重處理。
- 能否指出外部資料庫、搜尋引擎或 HTTP 服務需要自己的交易、冪等鍵或對帳機制。
回答前需要釐清的問題
先確認輸出是否仍在 Kafka、是否使用 Kafka Streams、消費者群組是否會滾動發布,以及外部系統是否必須與 Kafka 使用同一交易。再確認業務要的是每條輸入只產生一個可見輸出,還是外部副作用也必須只發生一次;後者決定是否需要目標系統協作。還要問允許的延遲、交易批次大小、失敗後重試窗口和可接受的積壓。
30 秒回答框架
我會把輸入紀錄、處理後的輸出和消費 offset 放進同一個 Kafka 交易。生產者啟用 transactional.id,消費者關閉自動提交並使用 read_committed,提交成功時三者一起可見,交易中止時輸出不可見且 offset 回到交易前。重啟使用穩定且唯一的交易 ID,舊實例會被 fencing。這個保證只涵蓋 Kafka 內的讀寫;寫外部資料庫時還要把 offset 與結果放進同一儲存交易,或使用冪等鍵、outbox 和對帳來收斂。
分步深入解答
1. 先寫出保證邊界
Kafka 官方把 topic 到 topic 的 exactly-once 定義為:輸出紀錄和消費位置可以在一個交易中原子更新。它不等於函式只執行一次,也不等於任意 HTTP 請求只到達一次。先說清邊界,後面的配置才不會變成口號。
2. 用交易式生產者原子寫輸出與 offset
消費者關閉自動提交,處理一批輸入後呼叫交易式生產者發送輸出,再把該批 offset 作為交易的一部分提交。核心偽代碼如下:
producer.initTransactions();
while (running) {
ConsumerRecords<String, Order> records = consumer.poll(timeout);
producer.beginTransaction();
try {
for (ConsumerRecord<String, Order> r : records) {
producer.send(new ProducerRecord<>("orders-enriched", r.key(), enrich(r.value())));
}
producer.sendOffsetsToTransaction(offsets(records), groupMetadata);
producer.commitTransaction();
} catch (AbortableException e) {
producer.abortTransaction();
consumer.seekToCommitted();
}
}輸出紀錄和 offset 同時提交,避免「輸出已寫但 offset 未前進」造成可見重複,也避免「offset 已前進但輸出未寫」造成漏處理。實際客戶端例外分類需要按版本文件處理,不能把所有例外都盲目重試。
3. 讓消費者只看已提交交易
配置 isolation.level=readcommitted,並保持 enable.auto.commit=false。readuncommitted 會暴露中止交易中的紀錄,消費者可能把本應撤銷的輸出再次傳播。read_committed 會等待或跳過交易標記,因此可見性與提交結果一致。
4. 處理重啟、fencing 與 rebalance
為每個活躍消費者實例使用穩定、叢集範圍唯一的 transactional.id。同一 ID 的新實例註冊時,Kafka 會讓舊實例的未完成交易中止並 fencing,防止兩個實例同時提交。交易中止後,應用程式要重建或明確回退消費者位置,再處理該批輸入;不能繼續使用已前進的本地游標。消費者群組的分區分配確保同一分區同一時間由一個成員處理。
5. 說明外部系統為何需要合作
若結果寫入 PostgreSQL,Kafka 交易不能自動包住資料庫提交。資料庫先提交、程序隨後崩潰會導致 Kafka offset 未提交,重試時必須靠唯一鍵或版本條件更新冪等;offset 先提交則可能漏寫。更強的方案是把結果和 offset 放在同一資料庫交易,或使用 outbox/連接器保存可重播紀錄,再由目標系統做去重與對帳。沒有目標系統協作,只能承諾至少一次執行加可偵測重複,不能宣稱端到端 exactly-once。
6. 用故障注入驗證,而不是只看配置
測試生產者在 sendOffsetsToTransaction 前後崩潰、commit 回應遺失、rebalance、舊實例被 fencing、交易逾時和下游不可用。用 read_committed 消費輸出,核對每個輸入 ID 的可見輸出數、最終 offset、重啟後的積壓和異常交易紀錄;對外部 sink 另外驗證唯一鍵衝突、重播和對帳結果。監控交易 abort 率、consumer lag、processing latency 與 fencing 次數,才能發現「看起來開啟交易、實際上輸出不完整」的配置錯誤。
高品質示範回答
我先限定題目範圍:如果輸入和輸出都在 Kafka,我會採用 Kafka Streams 或等價的交易式 consume-transform-produce。消費者關閉自動提交,生產者用唯一穩定的 transactional.id,一批輸出與該批 offset 透過 sendOffsetsToTransaction 原子提交;下游消費者使用 read_committed,所以中止交易的輸出不可見。重啟時同一交易 ID 會 fencing 舊實例,abort 後必須回退到最近提交的 offset 再處理。若輸出是 PostgreSQL,我不會把 Kafka 的保證擴大為端到端保證,而會把結果和 offset 放進同一資料庫交易,或用 outbox、冪等鍵和對帳。最後用崩潰、回應遺失、rebalance 和 fencing 故障注入,檢查每個輸入的可見輸出、offset 與外部狀態是否一致。
常見錯誤
- 把冪等生產者說成 exactly-once → 它主要防止生產重試在 Kafka 日誌中產生重複,未解決消費輸出與 offset 的原子性 → 補上交易和 offset 提交。
- 只設定
read_committed→ 它只改變可見性,不能替你提交交易或處理崩潰恢復 → 同時關閉自動提交並實作交易生命週期。 - 複用同一個
transactional.id給多個活躍實例 → 實例會互相 fencing,吞吐和可用性失控 → 按活躍實例分配穩定唯一 ID。 - 交易 abort 後繼續使用本地位置 → 本地位置可能落在未提交批次之後 → 重新取得已提交 offset 或明確 seek。
- 承諾 Kafka 能讓資料庫副作用只發生一次 → 兩個系統沒有天然的原子提交 → 使用目標庫交易、冪等寫、outbox 或對帳。
追問及應對
exactly-once 是否代表業務函式只執行一次?
不代表。函式可能在交易中執行後因 commit 失敗而再次執行;保證的是 Kafka 中可見輸出和 offset 的提交結果一致。函式應盡量無外部副作用,或把副作用設計成可重複與可對帳。
為什麼 read_committed 仍可能出現延遲?
消費者要跳過中止紀錄並等待交易完成,開放交易或長交易會推高可見性延遲和 lag。應限制交易批次與逾時,監控 transaction duration,並在低延遲場景權衡至少一次與交易成本。
交易中途發生 rebalance 怎麼辦?
新分配者只能從已提交 offset 繼續。當前交易應 abort,舊實例釋放或被 fencing;新實例重新取得分區並重播未提交批次。提交 offset 前必須確認分區仍由目前成員擁有。
如何把 Kafka 輸出交給 PostgreSQL?
在同一個 PostgreSQL 交易中寫業務結果和消費位置,或先寫帶唯一事件 ID 的 outbox,再由可靠 relay 投遞。若無法把兩者放進同一交易,就明確採用至少一次、唯一約束、版本條件更新和對帳,不使用 exactly-once 作為宣傳語。
哪些指標能證明方案真的生效?
按輸入事件 ID 統計 read_committed 輸出的唯一可見數量,結合已提交 offset、交易 abort/fence 次數、consumer lag、重啟後的重播量和外部 sink 的重複鍵衝突。故障注入前後都應保留樣本,不能只憑配置快照下結論。