適用於:✅ 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 基準)
| 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()資料的執行器可能會遭到縮減。 嚴重依賴執行者本地快取來處理熱資料的工作負載應該進行驗證。