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 模型與時間語意
Pipeline、轉換、PCollection、執行器:
- 一個 Beam pipeline 會將一個由 PTransform 組成的有向無環圖應用於 PCollection(有界或無界)。
- 執行器(如 Dataflow、Spark、Flink、Direct)負責執行 pipeline;Dataflow 提供託管的自動擴展、檢查點以及營運可視性。
- 轉換包含逐元素處理(ParDo)、分組與合併(GroupByKey、Combine)、聯結(CoGroupByKey)以及 IO 操作(PubSubIO、BigQueryIO、FileIO)。
視窗:
- 固定視窗 (Fixed windows):非重疊的時間區段(例如,1 分鐘的滾動視窗),用於週期性匯總。
- 滑動視窗 (Sliding windows):重疊的視窗,用於平滑的滾動指標(例如,每 1 分鐘滑動一次的 5 分鐘視窗)。
- 會話視窗 (Session windows):動態視窗,在一段無活動的間隔後關閉,非常適合用於使用者會話或設備的突發活動。
- 全域視窗 (Global window):整個無界串流的預設無視窗視圖;通常與觸發器搭配使用,以進行週期性的具體化。
事件時間 vs. 處理時間:
- 事件時間 (Event time):事件在來源端發生的時間;即使傳輸延遲不同,也能實現邏輯上一致的匯總。
- 處理時間 (Processing time):pipeline 觀測到事件的時間;對於營運觸發器很有用,但對於語意上的正確性則不然。
浮水印 (Watermarks):
- 浮水印估計事件時間的完整性(即執行器猜測它已看到 T 時間點之前的所有事件)。
- 在背壓或來源延遲的情況下,浮水印可能不規律地推進或停滯;任何時間戳小於浮水印的到達資料都屬於延遲資料。
觸發器與延遲性:
- 預設:AfterWatermark 觸發器,當浮水印超過視窗結尾時觸發;若允許的延遲 (allowed lateness) = 0,則延遲資料會被丟棄。
- 提早觸發 (Early firings)(基於處理時間或計數)可提供低延遲的初步結果。
- 延遲觸發 (Late firings) 允許在延遲資料到達時進行修正;累積模式 (accumulation mode) 決定窗格 (pane) 是要累積結果還是丟棄先前的輸出。
- 根據業務容忍度以及儲存/計算的權衡來選擇允許的延遲;越高的延遲會增加狀態保留時間和成本。
有狀態處理、計時器、會話化、重複資料刪除:
- 有狀態的 DoFn (Stateful DoFns) 會持有每個鍵的狀態(例如,最後看到的事件、運行中的匯總值),並設定計時器來發出或清除狀態。
- 會話化 (Sessionization) 可以很自然地透過 SessionWindows 來表達;對於自訂邏輯,則使用基於鍵的狀態 (keyed state) 和處理時間/事件時間計時器。
- 重複資料刪除 (Deduplication):為每個事件使用一個穩定的 ID,並在每個視窗中使用 Distinct/Combine,或使用每個鍵的狀態(例如,帶有 TTL 的布隆過濾器或集合)。在記憶體和偽陽性與嚴格的準確性之間進行權衡。
失敗模式與權衡:
- 使用處理時間視窗來計算業務指標,會在流量尖峰或重試時導致數據漂移;應優先使用事件時間視窗。
- 過小的視窗搭配頻繁的提早觸發,會導致過多的窗格發出和接收器的寫入放大。
- 無限制的允許延遲會導致狀態膨脹;務必為狀態設定 TTL 邊界,並設定計時器以清除休眠的鍵。
操作串流工作負載的 Dataflow
Worker 的大小設定與自動擴展:
- 水平自動擴展會根據待辦項目 (backlog)、浮水印延遲 (watermark lag)、CPU 和吞吐量 (throughput) 來新增/移除 worker;設定合理的
maxWorkers來吸收流量高峰。 - 針對瓶頸選擇機器類型:CPU 密集型(更多 vCPU)、記憶體密集型(高記憶體類型)、網路密集型(較大的 VM 可減少 shuffle 的額外開銷)。
- 對於繁重的 shuffle 或基於檔案的 sink,增加開機磁碟的大小。監控系統延遲 (system lag) 和待辦項目秒數 (backlog seconds)。
- 水平自動擴展會根據待辦項目 (backlog)、浮水印延遲 (watermark lag)、CPU 和吞吐量 (throughput) 來新增/移除 worker;設定合理的
Streaming Engine 與 shuffle:
- Streaming Engine 將狀態 (state) 和 shuffle 外部化到服務後端,改善彈性、減少 worker 的記憶體壓力,並實現更快的更新。
- 對於批次處理繁重的階段或大規模的鍵值分組,使用 Dataflow Shuffle 將 shuffle I/O 從 worker 卸載。兩者都能減少 hot-worker 故障和磁碟抖動 (disk thrash)。
背壓 (Backpressure)、熱鍵 (hot keys) 與傾斜 (skew):
- Dataflow 透過動態工作重新平衡來管理背壓;儘管如此,在適用情況下,仍需調整來源的流量控制(例如,Pub/Sub 未處理的訊息/位元組數量)。
- 熱鍵(例如,熱門的 id)會產生落後者 (stragglers)。可透過鍵值分片 (key sharding) (key#N)、部分預先彙總後再重新分鍵 (re-keying),或基於 sketch 的近似計算來緩解。
- 來自離群記錄(巨大的 payload)或突發性發布者的傾斜,可能需要每個發布者使用獨立分區、批次處理或壓縮。
Pub/Sub 整合:
- 使用 Pub/Sub 主題進行資料擷取;啟用訊息屬性以傳遞元資料(例如,deviceId、事件時間戳)。
- 使用
PubSubIO進行擷取;從屬性或 payload 中提取事件時間戳,否則退回使用發布時間。 - 排序鍵 (Ordering keys) 提供每個鍵的順序性;但由於至少一次 (at-least-once) 的交付機制,Dataflow 的下游行為仍需具備冪等性 (idempotent)。
串流至 BigQuery 的模式:
- 偏好使用
BigQueryIO搭配 Storage Write API,以透過串流偏移量 (stream offsets) 和自動重試,在單一串流中實現高吞吐量、低延遲的「僅一次」(exactly-once) 語意。 - 對於低速率的簡單 pipeline,串流插入 (streaming inserts) 是可接受的;設定
insertId來對客戶端的重試進行重複資料刪除。 - 對串流緩衝區 (streaming buffers) 的查詢是最終一致的;對於時間敏感的分析,可在緩衝區延遲後再查詢(例如,等待約 2 倍的觀測可用性延遲),或透過微批次視窗 (micro-batch windows) 和 Storage Write API 的 committed 模式來實現 (materialize)。
- 偏好使用
僅一次 (Exactly-once) 的效果、冪等性 (idempotency)、重播 (replay) 與 sink:
- Beam 保證至少一次 (at-least-once) 的處理;「僅一次」(exactly-once) 必須在 sink 端透過冪等寫入、交易或重複資料刪除鍵來實現。
- BigQuery:使用 Storage Write API 的預設串流 (default streams) 或已提交串流 (committed streams) 來達成單一串流內的 exactly-once;若使用 streaming inserts,則設定一個穩定的
insertId。 - 檔案:寫入帶有唯一名稱的暫存檔,在視窗完成時進行最終確認 (finalize),並確保原子性的重新命名;避免覆寫以防止部分重複。
- 外部資料庫:使用由穩定 id 作為鍵的 upsert 操作,或實作重複資料刪除視窗。
- 為重播而設計:維持具確定性的轉換;確保 sink 在重試時能刪除重複資料。
無法處理的信件 (Dead-letter) 處理、錯誤路由、可觀測性 (observability):
- 在
ParDo中用try/catch包裝有風險的解析/擴充,並透過TupleTag將失敗的項目發送到一個無法處理信件的PCollection;其中應包含 payload、錯誤碼和上下文。 - 將 DLQ(Dead-Letter Queues)路由到 BigQuery 或 Cloud Storage 進行分析;考慮使用一個獨立的 Pub/Sub 主題來進行重新處理。
- 可觀測性:使用 Dataflow 工作指標(浮水印延遲、系統延遲、吞吐量)、自訂計數器、分佈指標,以及在 Cloud Logging 中的各步驟日誌。在 Cloud Monitoring 中針對延遲和錯誤率建立警報。使用 Error Reporting 來彙總例外錯誤。
- 在
效能調校模式:
- 高效率讀取:對於 BigQuery 來源,偏好使用 Storage Read API 或基於查詢的讀取,只選擇需要的欄位和過濾條件。
- Combiner 優化:在
GroupByKey之前使用 combiners 來減少 shuffle 的資料量。 - 側輸入 (Side inputs):將小的參考資料快取在記憶體中;注意扇出 (fanout) 和更新的頻率。
- 序列化:使用緊湊的 schema(例如 Avro/Proto),並避免在熱路徑 (hot paths) 上進行過多的 JSON 解析。
部署、範本與升級策略
Flex Templates:
- 將 pipeline 打包成容器化、參數化的範本,以實現可重現的部署。Flex Templates 支援自訂相依性、GPU 映像檔和環境隔離。
- 將執行階段參數(例如:輸入訂閱、輸出資料表、無效信件接收器、maxWorkers)外部化,以實現針對特定環境的部署。
Pipeline 更新與相容性:
- 如果轉換名稱、狀態規格和輸出類型保持相容,Dataflow 支援對許多串流 pipeline 進行就地更新。請使用穩定的 PTransform 名稱。
- 對於不相容的圖形或狀態變更,執行受控的切換:啟動新的工作,然後清空(drain)舊的工作以完成處理中的工作並停止讀取新元素。
清空(Draining)與快照:
- 清空(Drain)會優雅地完成處理、寫入剩餘的輸出並終止;與 Pub/Sub 的保留或快照協調以避免資料間隙。
- 為確保連續性,您可以建立一個 Pub/Sub 快照,啟動新的 pipeline 並從該快照或適當的時間戳記開始讀取,驗證輸出後,再清空舊的工作。
設定範例:
- 具有早期/延遲觸發器和累計功能的視窗化範例:
undefined
- 使用 Storage Write API 的 BigQueryIO 範例:
undefined
- 常見陷阱:
- 在串流模式下寫入基於檔案的接收器時,若未使用視窗化寫入,可能會導致最終化(finalization)停滯;請啟用視窗化寫入和觸發器。
- 無限制增長:忘記限制狀態或允許的延遲時間,可能導致記憶體洩漏和擴展失敗。
- 缺少時間戳記:未指派事件時間戳記會導致 pipeline 預設使用處理時間,並在變動延遲下失去正確性。
實務問題情境
NovaTrack 公司接收來自全球 50,000 個溫度感測器的物聯網遙測資料,必須提供分鐘級的匯總資料、保存原始資料,並提供一個即時儀表板。預期會偶爾出現格式錯誤的訊息和亂序傳送。解決方案必須能自動擴展、呈現不良記錄以供檢查,並支援零停機時間的升級。
解決方法:
接收與時間語意
- 建立一個區域性的 Pub/Sub 主題和每個區域的發布者,並帶有
deviceId和eventTs(RFC3339) 屬性。在可行時,按deviceId啟用排序鍵。 - 理由:Pub/Sub 提供持久、有彈性的傳入機制,並具備至少一次的傳送保證。在邊緣附加事件時間戳記可以保留真實的事件時間;按裝置排序可減少裝置內的重新排序,而不會造成中央瓶頸。
- 建立一個區域性的 Pub/Sub 主題和每個區域的發布者,並帶有
使用事件時間視窗的 Dataflow 串流 pipeline
- 透過 PubSubIO 從專用訂閱中讀取,將
eventTs提取為 Beam 時間戳記,如果缺少則退回使用publishTime。 - 應用 1 分鐘的 FixedWindows,並在 30 秒時設定早期觸發器,對每個延遲元素設定延遲觸發;將允許的延遲時間設為 10 分鐘並累計窗格。
- 理由:事件時間視窗確保了分鐘級匯總的準確性;早期觸發以低於一分鐘的更新頻率為儀表板提供資料;延遲觸發則在延遲資料到達時修正匯總。延遲時間的限制可以控制狀態大小和成本。
- 透過 PubSubIO 從專用訂閱中讀取,將
驗證、豐富化與無效信件路由
- 實作一個 ParDo 來解析 JSON、驗證結構和範圍,並透過在工作啟動時從 BigQuery 載入的旁路輸入(side input)來豐富小型靜態參考資料。
- 使用 TupleTags 將有效記錄發送到主輸出,並將失敗的記錄發送到包含 payload、錯誤、
deviceId和解析時間戳記的無效信件 PCollection;將 DLQ 寫入一個分區的 BigQuery 資料表。 - 理由:旁路輸入將參考資料保存在記憶體中以實現低延遲。無效信件的捕獲允許檢查和有針對性地重新處理不良資料列,而不會阻塞主流程。
匯總與熱點鍵(hot-key)緩解
- 按
deviceId分組,並使用 CombineFns 計算每分鐘的 avg/min/max。對於前 N 名的區域性指標,按region#N進行分片以避免熱點鍵,然後再重新匯總。 - 理由:Combiners 能最小化 shuffle 的資料量和成本;鍵分片可防止在區域性扇入(fan-in)期間出現單一鍵的瓶頸。
- 按
接收器與恰好一次(exactly-once)效果
- 使用帶有 Storage Write API 的 BigQueryIO 將已驗證的原始事件和分鐘級匯總寫入 BigQuery。根據
deviceId+eventTs設定一個穩定的insert id,以在任何自訂重試中實現冪等性。 - 理由:Storage Write API 提供高吞吐量、低延遲的接收,並在串流中具有恰好一次的語意。穩定的 ID 確保在發生重播時下游可以進行重複資料刪除。
- 使用帶有 Storage Write API 的 BigQueryIO 將已驗證的原始事件和分鐘級匯總寫入 BigQuery。根據
儀表板一致性策略
- 儀表板查詢分區的匯總資料表,並相對於浮水印(watermark)回溯 2 分鐘,或對串流資料設定一個 2 倍於觀察到的可用性延遲的固定延遲。
- 理由:BigQuery 的串流可見性是最終一致的;稍微延遲讀取可以避免遺漏處理中的資料列,同時保持近乎即時的行為。
操作:自動擴展與串流引擎
- 啟用 Streaming Engine;根據預期峰值設定
maxWorkers(例如,平均值的 3 倍),選擇一個適合 CPU 密集型解析和加密的機器類型,並增加啟動磁碟以容納暫時的 shuffle 資料。 - 監控浮水印延遲、待辦積壓秒數、CPU 和每個步驟的吞吐量;對持續的延遲和 DLQ 率飆升設定警報。
- 理由:Streaming Engine 將狀態/shuffle 外部化,以實現彈性和更簡單的升級;適當的規模設定和監控可防止無聲的 SLO 違規。
- 啟用 Streaming Engine;根據預期峰值設定
使用 Flex Templates 進行部署與升級
- 將 pipeline 打包為一個帶有參數的 Flex Template:輸入訂閱、輸出資料表、DLQ 資料表、
maxWorkers和區域。對於不相容的變更,啟動一個新的 pipeline,目標是同一個主題但使用新的訂閱,驗證輸出後,再清空舊的工作。可選擇性地建立一個 Pub/Sub 快照,並讓新的訂閱從該快照開始讀取,以保證沒有資料間隙。 - 理由:Flex Templates 實現了可重複、參數化的部署。一個經過驗證的藍/綠切換與清空(drain)操作可以實現零資料遺失和最小的停機時間。
- 將 pipeline 打包為一個帶有參數的 Flex Template:輸入訂閱、輸出資料表、DLQ 資料表、
重新處理與批次補回
- 透過旁路輸出將壓縮後的原始事件 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.
通過考試 →