Dataflow 服務的動態工作重新平衡功能可讓服務根據執行階段條件,動態重新劃分工作。這些條件可能包括:
- 工作指派不平衡
- 工作站完成工作的時間比預期更慢
- 工作站完成工作的時間比預期更快
Dataflow 服務會自動偵測這些情況,並動態將工作指派給未使用的工作站或使用率偏低的工作站,以縮短整體工作處理時間。如要進一步瞭解 Dataflow 如何分配及擴充執行作業,請參閱「瞭解 Dataflow 中的平行處理」。
限制
只有在 Dataflow 服務平行處理某些輸入資料時,才會發生動態工作重新平衡:從外部輸入來源讀取資料時、使用具體化的中繼 PCollection 時,或使用匯總結果 (例如 GroupByKey) 時。如果作業中有大量步驟融合,作業的中間 PCollection 數量就會較少,動態工作重新平衡會受限於來源具體化 PCollection 中的元素數量。如要確保動態工作重新平衡機制可套用至管道中的特定 PCollection,可以透過幾種不同方式防止融合,確保動態平行處理。
動態重新平衡工作無法比單一記錄更精細地重新平行化資料。 如果資料含有導致處理時間大幅延遲的個別記錄,工作仍可能會延遲。Dataflow 無法將個別「熱門」記錄細分並重新分配給多個 worker。
Java
如果您為管道的最終輸出設定固定數量的分片 (例如使用 TextIO.Write.withNumShards 寫入資料),Dataflow 會根據您選擇的分片數量限制平行化。
Python
如果您為管道的最終輸出設定固定數量的分片 (例如使用 beam.io.WriteToText(..., num_shards=...) 寫入資料),Dataflow 會根據您選擇的分片數量限制平行化。
Go
如果為管道的最終輸出設定固定數量的分片,Dataflow 會根據您選擇的分片數量限制平行化。
使用自訂資料來源
Java
如果管道使用您提供的自訂資料來源,您必須實作 splitAtFraction 方法,讓來源與動態工作重新平衡功能搭配運作。
如果 splitAtFraction 實作有誤,來源記錄可能會重複或遭到捨棄。如需實作 splitAtFraction 的說明和提示,請參閱「RangeTracker 的 API 參考資料」。
Python
如果管道使用您提供的自訂資料來源,RangeTracker必須實作 try_claim、try_split、position_at_fraction 和 fraction_consumed,來源才能使用動態工作重新平衡功能。
詳情請參閱「RangeTracker 的 API 參考資料」。
Go
如果管道使用您提供的自訂資料來源,您必須實作有效的 RTracker,才能讓來源搭配動態工作重新平衡功能運作。
詳情請參閱 RTracker API 參考資訊。
動態工作重新平衡功能會使用自訂來源的 getProgress() 方法傳回值來啟動。getProgress() 的預設實作會傳回 null。如要確保自動調度資源功能會啟動,請確認自訂來源會覆寫 getProgress(),並傳回適當的值。