題幹與適用情境
一個電商平台把訂單事件寫入可重播的持久日誌,同時保留不可變原始資料。穩態吞吐未提供,尖峰為每秒 2 萬筆,促銷時可能短暫升至 10 倍。庫存異常要在事件產生後 10 秒內告警;營運儀表板要在 5 分鐘內更新;財務在隔日關帳,要求同一營業日可重跑、可對帳。約 2% 的事件會延遲,最長 6 小時。
以上吞吐、突發倍數、延遲比例與時限都是題目假設,不代表任何引擎的效能承諾。假設每筆事件有穩定的 eventid、orderid、event_time 和 schema 版本;財務口徑使用明確時區的營業日。團隊有四名資料工程師,現有資料倉儲與排程器能可靠執行批次工作,但還沒有成熟的有狀態串流工作輪值體系。
問題不要求三個消費端統一選擇同一種模式。目標是從行動時限、結果完整性、復原方式與維護成本推導最小足夠方案。本題歸入 data,因為核心能力是資料處理語意、時效與資料品質取捨,而非跨業務領域的平台架構。
面試官考察點
第一個訊號是候選人是否從業務決策的最晚時限倒推,而非看到訊息佇列就預設串流處理。來源是無界事件流,只表示資料持續抵達;它不決定每個下游都必須逐筆處理。財務隔日關帳可以讀取同一日誌的有界快照,以批次處理得到更容易重算的結果。
第二個訊號是能否區分處理延遲與行動延遲。引擎在 500 毫秒內完成計算,但接收端每 5 分鐘更新一次,儀表板仍不可能達到秒級。強回答會把預算拆成擷取、排隊、計算、寫入、快取與告警投遞,並為每一段設指標。
第三個訊號是正確性邊界。串流引擎的「正好一次」處理不能保證延遲資料已經完整,也不能自動保護外部副作用。批次工作可重跑也不代表安全;涵蓋範圍、業務鍵、快照版本與原子發布不清楚時,重跑同樣會重複或混合口徑。
最後看維運判斷。連續串流處理需要長期容量、積壓復原、狀態、檢查點、發布與輪值能力。若 5 分鐘時限能由微批次穩定滿足,引入逐筆有狀態計算只會增加故障面。反過來,10 秒告警若真的會觸發補貨或停售,按小時跑批就失去業務價值。
回答前需要釐清的問題
- 時限從哪裡開始,到哪裡結束? 本題從事件在來源系統提交成功開始,到告警送達或查詢結果可見為止。若只承諾「進入數倉」,架構會低估接收端與快取延遲。
- 10 秒告警會觸發自動動作嗎? 自動扣減、停售或通知需要去重、冪等與稽核;僅供觀察的告警可以接受較高誤報或重複率。
- 5 分鐘儀表板需要最新估算還是完整數字? 若允許附「資料截至時間」的可修訂估算,微批次足夠;若每次顯示都必須包含最長 6 小時的延遲事件,5 分鐘新鮮度與完整性要求互相衝突,必須修改產品契約。
- 財務關帳何時凍結,之後能否調整? 隔日產生初版、延遲後出調整分錄,與等待 6 小時再凍結是兩種口徑。它們決定批次工作的截點與修訂協定。
- 三個用途能否共用原始層與轉換定義? 可以共用正規化事件、業務鍵與測試樣本;若分別複製過濾、時區與金額規則,批次與串流結果終究會漂移。
- 團隊能否全天維護狀態串流工作? 沒有積壓告警、檢查點復原演練與安全發布能力時,應把連續串流範圍限制在確實需要秒級動作的最小路徑。
30 秒回答框架
「我先依業務行動時限拆分消費端,不依資料來源一刀切。10 秒庫存異常使用連續串流處理;5 分鐘營運儀表板用一到兩分鐘微批次,保留寫入與快取預算;隔日財務關帳使用營業日快照的批次處理,並作為對帳真源。三條路徑共用不可變原始事件、正規化規則與業務鍵。
串流結果都是可修訂檢視,不宣稱延遲資料已經完整;外部動作依事件與規則版本冪等。上線前用同一歷史區間重播批次、微批次與串流結果,影子執行比較筆數、金額與尾端延遲,再只為秒級路徑啟用動作。選擇標準是滿足行動 SLO 的最簡單模式,而且能證明復原後結果收斂。」
分步深入解答
第一步:把三個需求寫成結果契約
先為每個輸出寫清楚鍵、時限、完整性與修訂方式,再選 Spark、Flink 或其他產品。
| 用途 | 結果鍵 | 可見時限 | 完整性與修訂 | 建議模式 |
|---|---|---|---|---|
| 庫存異常 | 商品、規則版本、時間窗 | 10 秒 | 可快速修訂;動作必須冪等 | 連續串流 |
| 營運儀表板 | 指標、維度、視窗 | 5 分鐘 | 顯示截至時間;延遲後覆蓋新版本 | 1–2 分鐘微批次 |
| 財務關帳 | 營業日、帳戶、幣別 | 隔日 | 依凍結快照重算;差異走調整流程 | 批次 |
同一事件可服務三份不同契約。訊息佇列只是輸入傳輸層,不能把表中三列強迫變成同一種處理方式。
第二步:從行動時限分配端到端預算
10 秒告警可以先分配一個可測預算:來源提交與傳輸 2 秒,排隊 2 秒,計算 2 秒,接收端與規則動作 2 秒,告警投遞 2 秒。實際數字要透過壓測修正,但總和不能超過業務時限。每一段都記錄 p95、p99 與最大積壓年齡;只看運算子耗時會漏掉慢接收端。
5 分鐘儀表板可每 1 或 2 分鐘啟動微批次。假設使用 2 分鐘切片,排程等待最多約 2 分鐘,計算與寫入各保留 1 分鐘,快取更新保留 1 分鐘。若 10 倍突發讓計算超出預算,先增加平行度、縮短切片或延後非關鍵維度;不能用平均耗時宣稱達標。
財務的主要約束是可重算與口徑凍結,不是最低延遲。批次工作讀取明確的輸入快照或 offset 範圍,把結果寫入帶有 run_id 的暫存區,驗證通過後原子發布。重跑相同輸入與規則版本應得到相同結果。
第三步:讓三條路徑共用事實,避免複製實作債務
原始事件先進入可重播日誌或物件儲存,保留原始 payload、schema 版本、擷取時間與來源位置。正規化層統一處理 schema 演進、時區、金額單位、取消狀態和業務鍵。之後再由三種模式讀取。
共用事實不要求三個引擎逐行程式碼相同。更實用的邊界是共用資料契約、黃金樣本與確定性規則;若批次與串流必須分別實作視窗邏輯,就用同一輸入樣本做等價測試。財務還可以保留更嚴格的獨立驗證,不能因為重用程式碼而失去制衡。
第四步:分別設計重複、延遲與外部副作用
庫存串流依 event_id 去重,依事件時間計算短視窗,輸出包含規則版本與結果版本。傳送停售、補貨或通知時使用穩定的動作冪等鍵;檢查點成功不代表外部 HTTP 呼叫只發生一次。失敗重試必須能識別已執行動作。
營運微批次依固定半開區間讀取來源 offset,每次對目標視窗做版本化覆蓋,而非追加一個新的總數。延遲事件進入後續批次並提高結果版本。儀表板顯示「資料截至某時」,讓使用者知道五分鐘新鮮度不等於六小時尾端已完整。
財務批次處理等待約定截點後讀取凍結輸入。超過截點的事件進入下一次調整執行,保留原執行、輸入範圍、規則版本與差異核准。串流正好一次可保證某些持久化處理結果不重複,卻不能證明延遲資料完整,也不能替外部副作用自動提供相同語意。
第五步:把容量與成本比較成可重算問題
題設尖峰每秒 2 萬筆;10 倍突發就是每秒 20 萬筆。若每筆正規化事件以 1 KiB 粗估,入口尖峰約為每秒 195 MiB。這個數字只用於容量推導,真實設計必須用壓縮前後的大小分布重算。
連續串流要依尖峰與積壓復原速率配置長期容量。若故障 10 分鐘期間持續每秒 2 萬筆,會積壓 1,200 萬筆;復原後只提供等於目前流量的能力,積壓永遠清不掉。微批次則要證明單一切片能在下一切片到來前完成。批次處理可在低價時段集中資源,但小批次過密會讓啟動、提交與小檔案開銷占比上升。
成本比較至少包含常駐運算、狀態與檢查點儲存、接收端寫入、資料掃描、輪值與發布複雜度。不能只比較雲端帳單;四人團隊維護三套相近邏輯的時間也是成本。建議方案把連續串流限制在庫存告警,避免財務與儀表板承擔不必要的全天狀態成本。
第六步:用影子重播完成遷移,而非一次切換
先固定一段包含正常流量、10 倍突發、重複、亂序與 6 小時延遲的歷史輸入。舊批次結果作為基線,新微批次與串流工作以影子模式執行,不觸發真實庫存動作。比較每個業務鍵的筆數、金額、版本與延遲修訂軌跡,並解釋每個差異。
接著逐級放量:先只寫影子表,再開放內部儀表板,最後才讓串流告警觸發可回復動作。故障注入涵蓋檢查點前後當機、接收端逾時、分割區停滯、schema 變更與積壓追趕。驗收指標包括端到端 p99、最大積壓年齡、微批次完成時間、延遲修訂率、重複動作數、批次與串流差異金額及復原時間。
退出條件也要明確。若微批次持續超出 5 分鐘,檢查排程、傾斜、接收端與突發容量後再考慮連續串流;若 10 秒告警不再驅動行動,可以降為微批次以減少輪值成本。處理模式是可驗證的業務選擇,不是永久身分。
高品質示範回答
「我會把無界輸入與處理模式分開看。三個消費端的行動時限不同,所以不會為了技術統一把它們全做成串流。庫存異常確實要在 10 秒內觸發動作,我會用連續串流,逐段分配來源、排隊、計算、寫入與通知預算,並依 eventid 與規則版本產生冪等動作。營運儀表板只要求 5 分鐘,我會先用一到兩分鐘微批次,結果依視窗與版本覆蓋,頁面顯示資料截至時間。財務關帳則讀取凍結的營業日快照做批次處理,用 runid 暫存、核對後原子發布,延遲資料走調整執行。
三條路徑共用原始可重播事件、正規化資料契約與黃金樣本,但財務保留獨立驗證。題設每秒 2 萬筆、短暫 10 倍,我會驗證每秒 20 萬筆入口,以及故障期間積壓的復原倍率,而不只壓測穩態。串流工作的檢查點不能取代外部動作冪等,正好一次也不能證明 6 小時延遲尾端已經完整。
遷移時先重播同一段歷史,在影子表比較批次、微批次與串流的鍵級結果及修訂軌跡。連續串流先只發觀察告警,達到 10 秒 p99、沒有重複動作、積壓可在目標時間清空後再啟用自動動作。最終規則是:選擇滿足行動時限、完整性與復原要求的最簡單模式;更低延遲只有在改變業務決策時才值得承擔額外狀態與輪值成本。」
常見錯誤
- 看到 Kafka 就全部改成串流處理 → 輸入連續不代表每個結果都需逐筆更新,財務與儀表板承擔了無收益的狀態和維運成本 → 依每個消費端的行動時限分別選擇。
- 只說「串流處理延遲低」 → 接收端、快取與通知可能吃掉全部預算 → 拆出端到端路徑並監控各段尾端延遲。
- 把 5 分鐘新鮮度稱作最終正確 → 最長 6 小時的延遲事件仍會改變結果 → 明確定義估算、修訂、凍結與調整協定。
- 認為檢查點等於外部動作正好一次 → 重試可能重複呼叫庫存或通知服務 → 使用動作冪等鍵、稽核紀錄與可重播測試。
- 微批次不斷追加總數 → 同一視窗每批都被重複累計 → 依視窗與版本原子覆蓋,或定義清楚可撤回增量協定。
- 批次與串流各自複製業務規則 → 時區、取消單與金額口徑逐漸漂移 → 共用契約與黃金樣本,並做鍵級等價比較。
- 只依穩態吞吐擴充 → 10 倍突發與故障積壓無法在時限內清空 → 同時驗證尖峰、積壓復原倍率與接收端容量。
- 上線就觸發真實動作 → 語意差異直接影響庫存 → 先影子寫入與觀察告警,再逐級開放可回復動作。
追問及應對
追問一:5 分鐘儀表板為什麼不用連續串流?
若一到兩分鐘微批次在尖峰和復原情境下穩定完成,而且最終寫入與快取仍在 5 分鐘內,連續串流不會改變營運決策。它只增加長期狀態、檢查點、積壓復原與發布成本。若之後時限縮短到 30 秒,或微批次排程與啟動開銷已占掉大部分預算,再用同一歷史重播比較連續串流;升級條件應是 SLO 證據,不是「更即時」這個標籤。
追問二:能否用一條串流同時產出告警、儀表板和財務結果?
技術上可以,但要避免一個檢查點、schema 變更或錯誤視窗同時阻斷三種用途。告警可讀取正規化串流並維護短狀態,儀表板從版本化彙總讀取,財務仍從不可變歷史依凍結快照重算。共用輸入與定義,隔離發布與失敗域。只有團隊證明統一工作的故障影響、回填與稽核都可接受時,才值得合併執行單元。
追問三:正好一次處理是否能讓財務直接採用串流結果?
不能只憑這個標籤決定。處理語意可以防止已提交輸出因重試重複,卻不保證延遲紀錄已經到齊,也不會自動涵蓋外部副作用。財務還需要凍結輸入範圍、規則版本、可重複執行、總帳約束與差異核准。串流結果可作為早期估算;只有這些稽核條件同樣成立,並經過長期對帳證明後,才可能升級為關帳輸入。
追問四:10 倍突發時如何證明系統能復原?
記錄突發時間與故障時間,計算積壓筆數,再測實際淨清空速率。若目前輸入為每秒 2 萬筆、消費端處理每秒 3 萬筆,淨清空只有每秒 1 萬筆;積壓 1,200 萬筆需要約 20 分鐘。驗收同時看端到端告警延遲、接收端限額、狀態成長與自動擴縮延遲。只展示尖峰吞吐而不計算淨復原能力,會掩蓋長尾違約。
追問五:何時應該把連續串流降回微批次?
當業務不再依秒級結果行動、微批次已能滿足新時限,或串流工作的輪值與狀態成本持續高於它避免的損失時,可以降級。先影子執行微批次並比較結果與延遲,再停用真實動作、保留可重播輸入與回復時間窗。降級後仍要監控延遲修訂和尖峰完成時間,避免用成本最佳化換來不可見的資料陳舊。