高效的縮減規模與遠端重組管理器

適用於:✅ Fabric 資料工程與資料科學

高效率縮減是 Microsoft Fabric Spark 的一項功能,可將 Spark 重組資料與執行程式的生命週期解耦。 Fabric Spark 不會將 shuffle 輸出固定在執行器的本機磁碟上,而是將 shuffle 資料導向 Azure Blob 儲存體(或在需要時移轉至該處),並讓自適應查詢執行(AQE)決定實際的寫入方式。 結果是更快的叢集縮減、更低的運算成本,以及更具韌性的作業——且不需更改查詢、筆記本或管線。

Overview

高效縮減機制由四項相互協作的能力構成:

能力 其功能是什麼
Remote Shuffle Manager(RSM) 寫入與讀取作業會將混洗資料寫入及讀取自 Azure Blob 儲存體,而非執行器的本機磁碟。
重組遷移 在執行器退役之前,會先將 shuffle 區塊從該執行器移出,而不是直接丟棄。
決策層 按階段的執行時路由,將小型 shuffle 保留在本機,並將大型 shuffle 轉移至遠端儲存體。
AQE Shuffle 寫入 讓 Adaptive Query Execution 參與洗牌寫入階段,確保第一次分割是正確的。

先決條件

  • 啟用原生執行引擎(NEE)。
  • 啟用自動調整(建議)。 透過本文後述的 Spark 配置,有效縮放也能在無需自動縮放的情況下運作。
  • Runtime 1.3(Apache Spark 3.5) 或更新版本。

運作原理

當 Spark 處理查詢時,通常會在不同階段之間重新分散資料——也就是shuffle。 通常,每個執行器都會將 Shuffle 資料儲存在其本機磁碟上,因此執行器會受制於這些資料。 執行器必須等到所有消費者都完成讀取後才能釋放。 這種耦合是叢集無法快速縮減規模的最大原因,以及失去執行者會導致昂貴的階段重試。

有效縮放會打破這種耦合:

  • 大型重組作業會透過 Remote Shuffle Manager 直接寫入 Azure Blob 儲存體。
  • 小規模的 shuffle 作業 會保留在本機磁碟上,以提升速度。 若其執行器之後需要釋出,Shuffle 遷移會在背景中將區塊移動到對等節點或後備儲存體。
  • 決策層會在執行時為每個階段選擇正確的路徑。
  • AQE 混洗寫入可確保寫入器產生的分區方式能由下游 AQE 直接使用,而無需重新合併分區,避免浪費 I/O。
                ┌───────────────────────────┐
   Query  ───►  │   AQE + decision layer    │   per-stage choice
                └─────────────┬─────────────┘
                              │
                ┌─────────────▼─────────────┐
                │   AQE Shuffle Write       │   partition-aware writer
                └─────┬─────────────────┬───┘
                      │                 │
              local   ▼                 ▼   remote
        ┌────────────────────┐   ┌──────────────────┐
        │  Local disk +      │   │  RSM → Azure     │
        │  shuffle migration │   │  Blob Storage    │
        └─────────┬──────────┘   └─────────┬────────┘
                  │ on decommission        │
                  ▼                        ▼
        fallback storage   Remote shuffle store

智慧路由(決策層)

決策層會評估每個 Shuffle 交換,並決定:

  • 大型重分割作業 → Azure Blob 儲存體。 最大化縮減規模與容錯能力的效益。
  • 小型 shuffle → 本機磁碟。 小量傳輸無需承受雲端 I/O 的額外負擔。 如果執行器之後除役,則會由 Shuffle 遷移接手。

決策層會自動路由洗牌資料,且不需要你輸入。 建議的細緻度是按階段計算的。

主要優點

降低成本:只支付你所使用的運算費用

透過高效的縮減,執行人在工作完成後立即被釋放。 它們不再閒置不動,持有供下游任務最終讀取的 shuffle 資料。

  • 更快速縮減。 自動縮放會在任務完成後立即移除節點。
  • 減少閒置運算。 沒有任何「殭屍」執行器會僅為了提供其本機 shuffle 服務而持續存活。
  • 沒有磁碟過度配置。 大規模混洗作業會改為寫入 Blob 儲存體,而不需要大型本機磁碟空間。
  • 有界儲存成本。 當不再需要區塊時,備援儲存會自動清理。

更具韌性的職業

當洗牌資料只存在於本地磁碟時,執行器當機表示資料已遺失,Spark 必須重新計算。 透過有效率的縮容,資料不是已經位於 Blob 儲存體中,就是會在執行器終止前先遷移到該處。

劇本 沒有有效縮減 具備高效縮減能力
執行者當機 Shuffle 資料遺失;階段已重新執行 資料安全存放;無重計算
節點搶占 資料消失,重試昂貴 資料存活下來;工作照常進行
平順除役 Shuffle 在關機時被刪除 區塊遷移至對等或備用儲存
擷取期間的網路異常 層疊式 FetchFailedException 讀取來自儲存空間,沒有受到影響

此設計消除了在生產環境中造成 FetchFailedException 的最常見原因。

更快、真正具彈性的擴充

如果沒有高效率的縮容,只要節點上的任何執行器仍持有洗牌資料或快取資料,自動擴展器就無法回收該節點。 有效的縮編可將兩者解耦:

  • 重組資料會儲存在 Blob 儲存體中(或在關閉時遷移至該處)。
  • 快取不再將執行器固定住。 可重現的快取,如 Delta 快照快取,則不在縮減保護範圍內。

自動縮放器可以自由移除閒置節點,並根據工作負載變化調整叢集大小。

在傾斜且大洗牌時表現較佳

AQE Shuffle Write 讓自適應查詢執行(Adaptive Query Execution)直接決定 shuffle 寫出的方式——選擇可供下游 AQE 直接使用、無須再次合併分區的分區方式,並為遠端儲存產生數量更少且大小更合適的區塊。 結合決策層,你可以在大型或偏斜查詢上獲得更快的壁掛時鐘時間,而小型查詢則延遲不變。

開始

套用此配置以實現完整且高效的縮減堆疊:

# Remote Shuffle Manager
spark.conf.set("spark.remote.shuffle.enabled", "true")

# Decision layer — per-stage routing of local vs. remote shuffle
spark.conf.set("spark.sql.rsm.decisionlayer.enabled.level", "stage")

# AQE participates in shuffle write
spark.conf.set("spark.sql.adaptive.shuffleWrite.enabled", "true")

# Shuffle migration on executor decommission
spark.conf.set("spark.storage.decommission.shuffleBlocks.enabled", "true")
spark.conf.set("spark.storage.decommission.shuffleBlocks.cleanup", "true")
spark.conf.set("spark.storage.decommission.shuffleBlocks.migrateToFallbackStorage", "true")
spark.conf.set("spark.storage.decommission.fallbackStorage.cleanUp", "true")

不需要變更程式碼。 你也可以在你的環境 Spark 屬性中設定這些。

設定參考

遠端重組管理器(RSM)

Setting 建議 它控制的是什麼
spark.remote.shuffle.enabled true 開啟高效縮減。 Shuffle 資料會送到 Azure Blob 儲存體,而不是執行器本地磁碟。

決策層

Setting 建議 它控制的是什麼
spark.sql.rsm.decisionlayer.enabled.level stage 決策層路由隨機重排的粒度。 stage 會對每個 Spark 階段分別進行評估。

AQE 洗牌寫入

Setting 建議 它控制的是什麼
spark.sql.adaptive.shuffleWrite.enabled true 使 AQE 參與 shuffle 寫入階段。 產生可供下游 AQE 使用而無須重新合併分區的分割方式。

Note

AQE 本身(spark.sql.adaptive.enabled)必須啟用。 在 Fabric Spark 裡預設是開啟的。

除役時的重整遷移

Setting 建議 它控制的是什麼
spark.storage.decommission.shuffleBlocks.enabled true 將 shuffle 區塊從正在除役的 executor 遷移出去,而不是直接將其丟棄。
spark.storage.decommission.shuffleBlocks.cleanup true 在成功遷移後,清理原始執行器上的洗牌區塊。
spark.storage.decommission.shuffleBlocks.migrateToFallbackStorage true 如果沒有任何對等執行器無法接受這些區塊,則會將它們遷移到備用儲存(Azure Blob 儲存體)。
spark.storage.decommission.fallbackStorage.cleanUp true 一旦不再需要 shuffle 區塊,便會從備用儲存體中將其移除,藉此控制儲存成本。

快取感知型動態配置

Setting 建議 它控制的是什麼
spark.dynamicAllocation.preventShutdownExecutorWithCache false 允許動態配置釋放執行器,即使其仍持有已快取的區塊。
spark.dynamicAllocation.excludeDeltaSnapshotCache true 在判斷執行器是否仍持有有用快取時,忽略 Delta 快照快取。 Delta 快照快取可重建,不應阻止縮容。

進階調音(RSM)

大多數使用者不需要更改這些預設值。

寫入效能

Setting 預設值 它控制的是什麼
spark.remote.shuffle.partition.buffersize 16777216 (16 MB) 在寫入儲存體之前,先為每個分割區進行緩衝。
spark.remote.shuffle.blocksize 8388608 (8 MB) 上傳到 Blob 儲存體 的各個區塊大小。
spark.remote.shuffle.write.maxthreads cores × 16 寫入洗牌資料所用的最大執行緒數。
spark.remote.shuffle.write.maxtasks 16384 可同時進行的最大寫入作業數。

讀取效能

Setting 預設值 它控制的是什麼
spark.remote.shuffle.read.parallel.enabled true 平行下載串流用於隨機讀取。
spark.remote.shuffle.read.parallelism 4 每個任務的平行下載數。
spark.remote.shuffle.read.prefetchqueuesize 250 讀取時預取佇列深度。
spark.remote.shuffle.read.maxthreads cores × 4 閱讀時使用的最大執行緒數。

可靠性

Setting 預設值 它控制的是什麼
spark.remote.shuffle.retries 5 在發生暫時性儲存錯誤時重試。
spark.remote.shuffle.retrydelayms 800 重試之間的初始退避時間。
spark.remote.shuffle.retrymaxdelayms 60000 退避上限。

Compression

Setting 預設值 它控制的是什麼
spark.remote.shuffle.compression 用途 spark.io.compression.codec 遠端混洗資料的壓縮格式(例如,lz4、zstd)。

表現成績

圖表顯示在 TPC-DS 基準測試中,啟用與停用高效率縮減時的計算成本節省情形,顯示成本降低 54%。

計算成本節省(TPC-DS 基準)

Metric 沒有有效縮減 具備高效縮減能力
總運算量(VM-Minutes) 14,952 6,880
成本降低 — 54%

總工作執行時間可能更長(自動縮放使用較少同時執行者),但計費計算量減少超過一半。

決策層效能(TPC-DS,RSM 開啟)

將小洗牌路由到本地磁碟,且只將大洗牌送至遠端儲存,相較於遠端路由每個洗牌,執行時提升高達 57%,且享有相同的縮減效益。

Limitations

  • 需要 NEE。 高效的縮減取決於原生執行引擎。
  • 僅限 Azure Blob 儲存體。 標準 BlockBlobStorage,且已停用 HNS。 Azure Data Lake Gen2 / HNS 啟用的帳號不支援作為遠端洗牌儲存庫。
  • Azure Private Link 不支援。 使用私有連結網路的環境目前並不相容。
  • 決策層的細緻 度目前是按階段進行的。 每個任務或每個分區的路由不屬於涵蓋範圍。
  • 快取行為的改變。 使用 preventShutdownExecutorWithCache=false 時,持有 cache()/persist() 資料的執行器可能會遭到縮減。 嚴重依賴執行者本地快取來處理熱資料的工作負載應該進行驗證。