対象:✅ ファブリックデータエンジニアリングおよびデータサイエンス
効率的なスケールダウンは、Spark シャッフル データを Executor の有効期間から切り離す、Microsoft Fabric Spark の機能です。 ローカルの Executor ディスクにシャッフル出力をピン留めする代わりに、Fabric Spark はデータのシャッフルをAzure Blob Storageにルーティング (または必要に応じて移行) し、アダプティブ クエリ実行 (AQE) が書き込み自体を形成できるようにします。 その結果、クラスタのスケールダウンが速くなり、計算コストが低くなり、クエリやノートブック、パイプラインに変更が加わらず、よりレジリエントなジョブが得られます。
Overview
効率的なスケールダウンは、次の 4 つの連携機能から構築されます。
| 能力 | それが何をするか |
|---|---|
| リモート シャッフル マネージャー (RSM) | Executor ローカル ディスクではなく、Azure Blob Storageにシャッフル データを書き込んで読み取ります。 |
| シャッフルによる移行 | 停止する前に、実行プログラムを削除するのではなく、シャッフル ブロックを Executor から移動します。 |
| 意思決定層 | 小さなシャッフルをローカルに保持し、大きなシャッフルをリモート ストレージにオフロードするステージごとのランタイム ルーティング。 |
| AQE シャッフル書き込み | アダプティブ クエリ実行をシャッフル書き込みフェーズに参加させることができるので、パーティション分割が初めて適切になります。 |
Prerequisites
- ネイティブ実行エンジン(NEE)を有効にしてください。
- 自動スケールを有効にしてください(推奨)。 効率的なスケーリングダウンは、この記事の後半で説明するSpark構成を通じてオートスケールなしでも機能します。
- ランタイム 1.3 (Apache Spark 3.5) 以降。
どのように機能するのか
Sparkがクエリを処理する際、ステージ間でデータを再分配することがよくあり、これをシャッフルと呼びます。 通常、各エグゼキューターはローカルディスクにシャッフルデータを保存し、そのデータにエグゼキュータを紐付けます。 すべての消費者が読み終えるまで、遺言執行者は解放できません。 この結合がクラスタが迅速にスケールダウンできず、エグゼキュータを失うと高コストなステージ再試行が発生する最大の理由です。
効率的なスケールダウンにより、この結び付きは解消されます。
- Large shufflesリモート シャッフル マネージャーを介して直接Azure Blob Storageに移動します。
- 小さなシャッフルは、高速化のためローカルディスク上に保持されます。 もし後で実行者を解放する必要がある場合、シャッフル移行によってブロックはバックグラウンドのピアやフォールバックストレージに移動されます。
- 決定層は実行時に各段階で正しい経路を選択します。
- 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
スマートルーティング(意思決定層)
意思決定層は各シャッフル交換を評価し、以下を決定します:
- 大規模なシャッフル → Azure Blob Storage 最大スケールダウンとフォールト トレランスの利点。
- 小規模なシャッフル → ローカル ディスク 小さな転送に対するクラウド I/O オーバーヘッドはありません。 executor が後で停止された場合は、シャッフルマイグレーションが引き継ぎます。
意思決定レイヤーはシャッフルデータを自動でルーティングし、あなたの入力は一切必要ありません。 推奨される粒度はステージ単位です。
主な利点
コストの削減: 使用するコンピューティングに対してのみ支払う
効率的なスケールダウンにより、Executor は作業が完了するとすぐにリリースされます。 下流タスクが後で読み取る可能性のあるシャッフルデータを保持したまま、何もせず遊休状態になることはなくなります。
- より高速なスケールダウン。 自動スケールでは、タスクの完了後すぐにノードが削除されます。
- アイドル状態のコンピューティングを削減。 ローカルシャッフルを提供するためだけに生かされ続ける「ゾンビ」エグゼキューターはありません。
- ディスクオーバープロビジョニングなし。 大規模なシャッフル処理は、大容量のローカル ディスクを必要とせず、Blob Storage に格納されます。
- 有界ストレージ コスト。 フォールバック ストレージは、ブロックが不要になったときに自動的にクリーンアップされます。
回復性の高いジョブ
シャッフル データがローカル ディスクにのみ存在する場合、Executor のクラッシュはデータがなくなったことを意味し、Spark はそれを再計算する必要があります。 効率的なスケールダウンにより、データは既に BLOB ストレージに格納されているか、Executor が終了する前にそこに移行されます。
| Scenario | 効率的なスケールダウンなし | 効率的なスケールダウンで |
|---|---|---|
| Executor がクラッシュする | 失われたデータをシャッフルする。ステージの再実行 | データはストレージ内で安全です。再計算なし |
| ノードのプリエンプション | データがなくなった、コストの高い再試行 | データは保持される。ジョブは通常どおり継続される |
| 正常停止 | シャットダウン時にシャッフルが解除された | ピアまたはフォールバック ストレージに移行されたブロック |
| フェッチ中のネットワークの一時的な不具合 | カスケーディング FetchFailedException |
読み取りはストレージから行われ、影響を受けません |
この設計により、生産中の最も一般的な FetchFailedException 原因が排除されます。
迅速で、本当に柔軟なスケーリング
効率的なスケールダウンがないと、オートスケーラーはノードを再利用できませんが、その上の実行プログラムはシャッフル データまたはキャッシュされたデータを保持します。 効率的なスケールダウンでは、その両方を切り離せます:
- シャッフル データは BLOB ストレージ内にあります (または、シャットダウン時にそこに移行されます)。
- キャッシュはエグゼキューターを固定しなくなりました。 Delta スナップショットキャッシュなどの再生成可能なキャッシュは、スケールダウン保護の対象外です。
自動スケーラーは、ワークロードの変更に応じて、アイドル状態のノードを自由に削除し、クラスターのサイズを変更できます。
偏りのあるシャッフルや大規模なシャッフルでのパフォーマンス向上
AQE Shuffle Writeは、Adaptive Query Executionがシャッフル書き込み自体を形作ることを可能にします。つまり、下流の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 |
効率的なスケールダウンをオンにします。 シャッフル データは、Executor ローカル ディスクではなくAzure Blob Storageに行きます。 |
意思決定層
| Setting | 推奨 | 制御する内容 |
|---|---|---|
spark.sql.rsm.decisionlayer.enabled.level |
stage |
意思決定層でシャッフルがルーティングされる粒度。
stage では、各 Spark ステージが個別に評価されます。 |
AQE シャッフル書き込み
| Setting | 推奨 | 制御する内容 |
|---|---|---|
spark.sql.adaptive.shuffleWrite.enabled |
true |
AQE がシャッフル書き込みフェーズに参加できるようにします。 下流のAQEが再集約することなく利用できるパーティション分割を生成します。 |
Note
AQE 自体 (spark.sql.adaptive.enabled) がオンになっている必要があります。 Fabric Spark では既定でオンになっています。
廃止時のシャッフル移行
| Setting | 推奨 | 制御する内容 |
|---|---|---|
spark.storage.decommission.shuffleBlocks.enabled |
true |
シャッフル ブロックを削除するのではなく、使用停止中の Executor からシャッフル ブロックを移行します。 |
spark.storage.decommission.shuffleBlocks.cleanup |
true |
移行が成功した後、ソース Executor のシャッフル ブロックをクリーンアップします。 |
spark.storage.decommission.shuffleBlocks.migrateToFallbackStorage |
true |
ピア Executor がブロックを受け入れられない場合は、フォールバック ストレージ (Azure Blob Storage) に移行します。 |
spark.storage.decommission.fallbackStorage.cleanUp |
true |
必要なくなったフォールバック ストレージからシャッフル ブロックを削除し、ストレージ コストを制限します。 |
キャッシュ対応の動的割り当て
| Setting | 推奨 | 制御する内容 |
|---|---|---|
spark.dynamicAllocation.preventShutdownExecutorWithCache |
false |
キャッシュされたブロックを保持している場合でも、Executor を解放する動的割り当てを許可します。 |
spark.dynamicAllocation.excludeDeltaSnapshotCache |
true |
エグゼキューターがまだ有効なキャッシュを保持しているかどうかを判断する際に、Delta スナップショット キャッシュを無視します。 差分スナップショットキャッシュは再現性があり、スケールダウンを妨げるべきではありません。 |
高度なチューニング (RSM)
ほとんどのユーザーは、これらの既定値を変更する必要はありません。
書き込み性能
| Setting | デフォルト | 制御する内容 |
|---|---|---|
spark.remote.shuffle.partition.buffersize |
16777216 (16 MB) |
ストレージに書き込む前のパーティションごとのバッファー。 |
spark.remote.shuffle.blocksize |
8388608 (8 MB) |
Blob Storageにアップロードされた個々のブロックのサイズ。 |
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 |
読み取りに使用される最大スレッド数。 |
Reliability
| 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% の実行時間向上を実現し、同じスケールダウン効果が得られます。
制限事項
- NEE が必要です。 効率的なスケールダウンは、ネイティブ実行エンジンによって異なります。
- Azure Blob Storageのみ。 Standard
BlockBlobStorage(HNS は無効)。 Azure Data Lake Gen2/HNS 対応アカウントは、リモート シャッフル ストアとしてサポートされていません。 - Azure Private Linkではサポートされていません。 プライベート リンク ネットワークを使用する環境には、現在互換性がありません。
- 意思決定層の細分 度は現在、各段階ごとに分類されています。 タスクごとまたはパーティションごとのルーティングはスコープ内にありません。
- キャッシュ動作の変更。
preventShutdownExecutorWithCache=falseでは、cache()/persist()データを保持している Executor がスケールダウンされる可能性があります。 ホット データの Executor ローカル キャッシュに大きく依存するワークロードは検証する必要があります。