1. 題目
一個訂單事件流按 customerId 聚合五分鐘內的金額與訂單數。裝置時鐘可能漂移,網路重試會造成亂序,某些 Kafka 分割區還會暫時沒有新訊息。結果需要盡快展示,同時允許遲到資料在有限時間內更正。請設計事件時間處理方案。
2. 約束與澄清
- 先確認視窗按事件發生時間、寫入時間還是處理時間計算;本題選擇事件時間。
- 設定最大亂序範圍,例如 30 秒,並說明超過範圍後的處理方式。
- 明確結果是追加事件還是可更新快照;下游必須能識別修訂版本。
- 討論閒置分割區、錯誤時間戳、重複事件與重播任務,避免把「沒有新訊息」誤判成輸入結束。
3. 核心概念
事件時間來自記錄本身,水位線表示系統認為事件時間已推進到某個位置。視窗通常在水位線越過視窗末端後觸發;水位線之後到達、但時間戳仍落在視窗內的記錄屬於遲到事件。並行輸入合併時,算子通常要等待各輸入的最小水位線,因此一個不活躍分割區可能拖慢全域進度,必須支援 idle 標記或個別分割區逾時策略。
4. 參考流程
onRecord(event):
ts = extractEventTimestamp(event)
key = canonicalKey(event.customerId)
updateWatermarkGenerator(ts)
window = floorToFiveMinutes(ts)
if ts <= currentWatermark + allowedLateness:
state[window, key] = aggregate(state[window, key], event)
emitUpsert(window, key, state[window, key], revision + 1)
else:
routeToLateData(window, key, event)
onWatermark(wm):
finalizeWindowsBefore(wm)
expireStateAfterRetention()來源端提取事件時間並生成有界亂序水位線;輸入分割區長時間無資料時標記 idle,避免它把合併水位線卡住。視窗狀態先保留到允許遲到期限之後,遲到記錄透過 upsert 或更正事件更新結果;超過期限的資料進入旁路流,供人工核查或離線回補。
5. 準確性與延遲取捨
允許遲到時間越長,結果越接近完整事件時間視圖,但狀態與更正次數也越多。只發送一次最終結果延遲較低,卻無法表達遲到修訂;發送帶視窗鍵、版本號與變更原因的更新事件,要求下游具備冪等合併能力。對於金額等不可遺失指標,可以保留原始事件並安排離線重算;對於即時排行榜,則可接受超過期限後只在下一批刷新。
6. 驗證與觀測
- 建立有序、亂序、遲到 30 秒內與超過期限的事件,逐視窗核對結果與離線基準。
- 注入閒置分割區、時鐘跳變、重複事件與任務重啟,檢查水位線是否停滯或倒退。
- 記錄目前水位線、處理時間與事件時間差、視窗狀態大小、遲到率、旁路數量與更正次數。
- 對每次更新記錄視窗鍵、版本與輸入事件 ID,重播同一日誌時比較最終快照是否一致。
7. 常見誤區
- 用處理時間視窗取代事件時間視窗,卻聲稱結果能抵抗網路亂序。
- 把水位線當成絕對完成標記;實際水位線可能是基於延遲假設的啟發式進度。
- 忽略閒置分割區,導致一個沒有訊息的分割區阻塞所有視窗關閉。
- 遲到資料直接丟棄且沒有旁路、版本或離線回補路徑。
8. 面試評分點
能區分三種時間
應明確事件時間、攝取時間與處理時間的語義,並說明視窗選擇如何影響準確性與延遲。
能生成並合併水位線
應解釋有界亂序、並行分割區最小水位線與 idle 處理,指出水位線是進度估計而非絕對承諾。
能設計遲到與更正路徑
應定義允許遲到期限、旁路流、更新版本與下游冪等合併方式,覆蓋期限內與期限外兩類資料。
能用重播驗證結論
應透過亂序、閒置、重啟與重複事件測試,與離線基準及即時指標對帳,而不是只觀察任務是否存活。