Google PDE: 使用 Dataflow 與 Apache Beam 進行串流處理 — 學習指南

屬於 Google Professional Data Engineer — 學習指南. 使用經過驗證的解答練習: Google 考試中心, 或參加限時模擬考試: ExamRoll.io.

總覽

Google Cloud 上的串流處理,核心是圍繞著由 Dataflow 執行器所執行的 Apache Beam 統一程式設計模型。Beam 提供了一個邏輯上的抽象層——在 PCollection 上進行轉換的 pipeline——它將您的程式碼與執行細節(如平行處理、自動擴展和容錯能力)解耦。在串流處理中,正確性取決於時間語意(事件時間 vs. 處理時間)、視窗化(固定、滑動、會話、全域)、浮水印、觸發器以及延遲資料的處理。要在 Dataflow 上達到卓越的營運,需要正確的 worker 規模設定、自動擴展策略、串流引擎、shuffle 選擇、冪等接收器設計、無效信件處理以及穩健的可觀測性。

Apache Beam 模型與時間語意

失敗模式與權衡:

操作串流工作負載的 Dataflow

部署、範本與升級策略

undefined

undefined

實務問題情境

NovaTrack 公司接收來自全球 50,000 個溫度感測器的物聯網遙測資料,必須提供分鐘級的匯總資料、保存原始資料,並提供一個即時儀表板。預期會偶爾出現格式錯誤的訊息和亂序傳送。解決方案必須能自動擴展、呈現不良記錄以供檢查,並支援零停機時間的升級。

解決方法:

  1. 接收與時間語意

    • 建立一個區域性的 Pub/Sub 主題和每個區域的發布者,並帶有 deviceIdeventTs (RFC3339) 屬性。在可行時,按 deviceId 啟用排序鍵。
    • 理由:Pub/Sub 提供持久、有彈性的傳入機制,並具備至少一次的傳送保證。在邊緣附加事件時間戳記可以保留真實的事件時間;按裝置排序可減少裝置內的重新排序,而不會造成中央瓶頸。
  2. 使用事件時間視窗的 Dataflow 串流 pipeline

    • 透過 PubSubIO 從專用訂閱中讀取,將 eventTs 提取為 Beam 時間戳記,如果缺少則退回使用 publishTime
    • 應用 1 分鐘的 FixedWindows,並在 30 秒時設定早期觸發器,對每個延遲元素設定延遲觸發;將允許的延遲時間設為 10 分鐘並累計窗格。
    • 理由:事件時間視窗確保了分鐘級匯總的準確性;早期觸發以低於一分鐘的更新頻率為儀表板提供資料;延遲觸發則在延遲資料到達時修正匯總。延遲時間的限制可以控制狀態大小和成本。
  3. 驗證、豐富化與無效信件路由

    • 實作一個 ParDo 來解析 JSON、驗證結構和範圍,並透過在工作啟動時從 BigQuery 載入的旁路輸入(side input)來豐富小型靜態參考資料。
    • 使用 TupleTags 將有效記錄發送到主輸出,並將失敗的記錄發送到包含 payload、錯誤、deviceId 和解析時間戳記的無效信件 PCollection;將 DLQ 寫入一個分區的 BigQuery 資料表。
    • 理由:旁路輸入將參考資料保存在記憶體中以實現低延遲。無效信件的捕獲允許檢查和有針對性地重新處理不良資料列,而不會阻塞主流程。
  4. 匯總與熱點鍵(hot-key)緩解

    • deviceId 分組,並使用 CombineFns 計算每分鐘的 avg/min/max。對於前 N 名的區域性指標,按 region#N 進行分片以避免熱點鍵,然後再重新匯總。
    • 理由:Combiners 能最小化 shuffle 的資料量和成本;鍵分片可防止在區域性扇入(fan-in)期間出現單一鍵的瓶頸。
  5. 接收器與恰好一次(exactly-once)效果

    • 使用帶有 Storage Write API 的 BigQueryIO 將已驗證的原始事件和分鐘級匯總寫入 BigQuery。根據 deviceId + eventTs 設定一個穩定的 insert id,以在任何自訂重試中實現冪等性。
    • 理由:Storage Write API 提供高吞吐量、低延遲的接收,並在串流中具有恰好一次的語意。穩定的 ID 確保在發生重播時下游可以進行重複資料刪除。
  6. 儀表板一致性策略

    • 儀表板查詢分區的匯總資料表,並相對於浮水印(watermark)回溯 2 分鐘,或對串流資料設定一個 2 倍於觀察到的可用性延遲的固定延遲。
    • 理由:BigQuery 的串流可見性是最終一致的;稍微延遲讀取可以避免遺漏處理中的資料列,同時保持近乎即時的行為。
  7. 操作:自動擴展與串流引擎

    • 啟用 Streaming Engine;根據預期峰值設定 maxWorkers(例如,平均值的 3 倍),選擇一個適合 CPU 密集型解析和加密的機器類型,並增加啟動磁碟以容納暫時的 shuffle 資料。
    • 監控浮水印延遲、待辦積壓秒數、CPU 和每個步驟的吞吐量;對持續的延遲和 DLQ 率飆升設定警報。
    • 理由:Streaming Engine 將狀態/shuffle 外部化,以實現彈性和更簡單的升級;適當的規模設定和監控可防止無聲的 SLO 違規。
  8. 使用 Flex Templates 進行部署與升級

    • 將 pipeline 打包為一個帶有參數的 Flex Template:輸入訂閱、輸出資料表、DLQ 資料表、maxWorkers 和區域。對於不相容的變更,啟動一個新的 pipeline,目標是同一個主題但使用新的訂閱,驗證輸出後,再清空舊的工作。可選擇性地建立一個 Pub/Sub 快照,並讓新的訂閱從該快照開始讀取,以保證沒有資料間隙。
    • 理由:Flex Templates 實現了可重複、參數化的部署。一個經過驗證的藍/綠切換與清空(drain)操作可以實現零資料遺失和最小的停機時間。
  9. 重新處理與批次補回

    • 透過旁路輸出將壓縮後的原始事件 Avro 檔案儲存在 Cloud Storage;當模型或結構變更時,執行一個批次的 Dataflow pipeline 來補回或重新處理資料到 BigQuery。
    • 理由:持久的原始封存支援可重現性和結構演進,而不會影響熱路徑(hot path)。

這個設計能夠在處理全球規模的亂序和延遲資料的同時,產生正確、低延遲的匯總資料,且成本受控、錯誤隔離清晰、可觀測性強,並提供安全的升級路徑。


BigQuery 分析與倉儲工程 · 所有領域 · 訊息傳遞、事件擷取與即時服務

練習這些題目 → · 在 ExamRoll.io 上限時練習 →

Pass the whole exam — not just this question

You found this answer. Get every verified question and explanation in one place, and save hours of prep. Free to start.

通過考試 →

瀏覽 Google →

Related guides

一站式存取

一份訂閱。所有考試。

每個方案都可無限存取答案搜尋、練習測驗、AI 解釋和完整的資源庫 — 支援 20 多種語言。

每月
24.87
Just €0.83/day
包含所有內容:
  • 無限答案搜尋
  • 無限練習測驗
  • AI 驅動的解釋
  • 完整資源庫
  • 20 多種語言
  • 每週內容更新
  • 獎勵與推薦
  • 優先支援
開始免費試用

無需信用卡*

最佳價值
12 個月
179.87
Just €0.49/daySave 40%
包含所有內容:
  • 無限答案搜尋
  • 無限練習測驗
  • AI 驅動的解釋
  • 完整資源庫
  • 20 多種語言
  • 每週內容更新
  • 獎勵與推薦
  • 優先支援
開始免費試用

無需信用卡*

✓ 包含免費方案 · ✓ 隨時取消 · ✓ 所有方案解鎖完整產品