Dataflow 的設計宗旨是透過在代管的運算執行個體集區中分配工作,執行大型資料處理管道。瞭解 Dataflow 如何平行處理作業,有助於設計有效率的管道、避免效能瓶頸,以及盡可能降低資源成本。
本頁面說明 Dataflow 如何平行處理資料、管理及擴充執行作業、限制平行處理的常見因素,以及可用於最佳化管道處理量的技術。
平行處理模型:水平與垂直
Dataflow 會使用兩種互補策略來實現平行處理:
水平平行處理:管道資料會分割並同時在多個 worker 執行個體 (虛擬機器) 中處理。Dataflow 可透過自動水平調度資源功能,根據工作負載需求自動調整工作站集區大小。根據預設,Dataflow 會為每項工作設定 4,000 個工作站的資源限制,您可以使用配額要求調整這項限制。
垂直平行處理:單一 worker 例項中的多個 CPU 核心和執行緒會同時處理管道資料。每個 worker VM 都會執行 worker 程序和 harness 執行緒,以使用可用的運算資源。透過動態執行緒調度,Dataflow 可根據 CPU 使用率和記憶體空間,調整批次管道中每個工作站的現用執行緒數量。在 Dataflow Prime 中,自動垂直調度資源功能會動態調整分配給工作站的記憶體和運算資源。
工作單元和執行階層
為在工作站和執行緒之間分配處理作業,Dataflow 會將 Apache Beam 管道分成個別的工作單元:
- PCollection 和分區:
PCollection代表分散式資料集。如果是有界資料 (批次管道),Dataflow 會將資料集分割成多個分割區或分片。如果是無界限資料 (串流管道),資料會持續傳入,並以訊息或串流分割區的形式擷取。 - 套件:Dataflow 會將元素分組為任意套件,供
DoFn處理。套裝組合是失敗和重試的單位:如果處理元素時引發未處理的例外狀況,系統會重試整個套裝組合。記憶體消耗量高的作業可能會增加工作人員的記憶體壓力,並導致記憶體不足錯誤。 - 階段和步驟融合:在圖表最佳化期間,Dataflow 會將相鄰的轉換合併為融合的執行階段,以消除中繼資料具體化的負擔。在融合階段中,元素會在單一執行緒的緊密執行迴圈中處理,然後傳遞至下一個階段或隨機播放邊界。
如要進一步瞭解管道轉換和圖表產生作業,請參閱「管道生命週期」。
代管式平行處理和自動調度資源
根據預設,Dataflow 會自動管理管道平行處理作業,不需手動調整分區,方法如下:
- 水平自動調度:
- 批次管道:評估預估剩餘工作總量、來源待處理工作量和 CPU 使用率,藉此擴大或縮減工作站集區,以快速且經濟實惠的方式完成工作。
- 串流管道:分析系統延遲、待處理工作數量和 CPU 使用率,在處理量尖峰期間增加工作站數量,並在流量較少的期間減少工作站數量。詳情請參閱「調整串流自動水平調度資源」。
- 動態工作重新平衡 (DWR):在批次管道中,Dataflow 會監控個別工作站工作的進度。如果 worker 提早完成工作,或另一個 worker 因資料偏斜 (落後者) 而落後,Dataflow 會動態分割 slow worker 剩餘未處理的工作,並重新指派給閒置的 worker。詳情請參閱「動態重新平衡工作」。
- 動態執行緒資源配置:在採用可攜式 Runner 的批次管道中,系統會根據 CPU 使用率和記憶體空間,自動調整每個 worker 的並行處理執行緒數量。詳情請參閱「動態執行緒調整」。
- 自動垂直調度資源:在 Dataflow Prime 中,Dataflow 會動態調度工作站記憶體和運算資源,避免發生記憶體不足錯誤,並盡可能提升資源使用率。詳情請參閱「垂直自動調度資源」。
限制平行處理的因素
由於下列資料特徵或管道圖設計,管道可能無法達到預期的平行處理:
無法分割的輸入來源
如果輸入來源無法分割成獨立範圍,Dataflow 就會強制使用單一工作執行緒循序讀取來源:
- 無法分割的檔案壓縮:系統無法從任意位元組偏移量平行讀取
.gz(gzip) 或.bzip2(沒有索引) 等格式。讀取單一大型壓縮檔時,資料解壓縮並重新分配前,擷取階段會限制為單一執行緒。 - 解決方法:以可分割的檔案格式 (例如 Parquet、Avro 或 Snappy 壓縮格式) 儲存資料,或將輸入資料分割成多個較小的檔案,存放在 Cloud Storage 中。
步驟融合和高擴散傳遞功能
步驟融合可減少序列化負擔,進而提升效能,但如果平行處理量低的步驟產生大量輸出元素 (「高擴散傳遞功能」作業),可能會無意間限制平行處理量並增加記憶體壓力:
- 範例:來源會讀取五個檔案,並與產生 1,000,000 個輸出元素的
FlatMap轉換融合。如果FlatMap轉換與下游轉換融合,所有 1,000,000 個元素最多會繼續在五個背景工作執行緒上執行,嚴重限制下游處理量。此外,如果中繼轉換在提交前於記憶體中大幅擴充,大型套件可能會耗盡可用的 worker 記憶體。 - 解決方法:在高擴散傳遞功能步驟和下游轉換之間插入
Redistribute(或傳統的Reshuffle) 轉換,即可中斷融合,並在工作人員集區中重新分配工作。如要偵錯記憶體相關問題,請參閱「疑難排解記憶體不足錯誤」。
主要傾斜和熱鍵
匯總作業 (GroupByKey、CoGroupByKey、Combine.PerKey) 會依相關聯的鍵將元素分組。
- 熱鍵瓶頸:Dataflow 會將具有相同鍵的所有元素,路由至單一背景工作執行緒進行匯總。如果單一鍵包含的資料量占整個資料集的大部分,該 worker 就會成為落後者,上游 worker 可能會遇到背壓。例如預設
null鍵或極受歡迎的類別鍵。 - 解決方法:
- 盡可能使用 Combiner (
CombineFn或Combine.PerKey) 取代GroupByKey,讓 Dataflow 在重組前執行部分本機組合。 - 在熱鍵中加入隨機整數前置字串或後置字串 (鍵鹽化),將鍵空間分配給各個 worker,然後進行第二階段的彙整,合併鹽化結果。
- 盡可能使用 Combiner (
下游接收器節流
將管道輸出內容寫入資料庫或第三方 API 等外部服務時,高平行處理可能會使目的地系統飽和:
- 節流:數百個工作執行緒同時發出寫入呼叫,可能會導致速率限制錯誤、連線逾時或資料庫效能降低。
- 解決方法:
- 使用
GroupByKey分組元素,或使用具有受控平行處理的批次處理接收器,限制寫入平行處理。 - 在接收器
DoFn實作中,導入用戶端指數輪詢和重試邏輯。
- 使用
最佳化策略
如要最佳化 Dataflow 工作的平行處理,請考慮下列方法:
- 防止與
Redistribute發生不當融合:Redistribute.arbitrarily():中斷步驟融合,並在所有可用的 worker 之間平均重新分配元素。Redistribute.byKey():在工作執行緒之間重新平衡鍵/值組,同時保留鍵的區域性。- 如需實作範例,請參閱「防止融合」。
- 監控落後者和瓶頸: 使用 Google Cloud 控制台執行詳細資料,找出落後者數量偏高或進度停滯的階段:
後續步驟
- 瞭解管道生命週期。
- 瞭解水平自動調度資源。
- 瞭解動態重新平衡工作。
- 請參閱 Dataflow 管道最佳做法。
- 瞭解如何排解記憶體不足錯誤。