Effektiv skalning och fjärrblandningshanterare

Gäller för:✅ Fabric Data Engineering och Data Science

Effektiv nedskalning är en funktion i Microsoft Fabric Spark som frikopplar Spark shuffle-data från exekverarens livslängd. I stället för att fästa shuffle-utdata på lokala kördiskar dirigerar Fabric Spark shuffle-data till Azure Blob Storage (eller migrerar dem dit på begäran) och låter Adaptive Query Execution (AQE) forma själva skrivningen. Resultatet är snabbare klusternedskalning, lägre beräkningskostnader och mer motståndskraftiga jobb – utan förändringar i dina frågor, anteckningsböcker eller pipelines.

Overview

Effektiv nedskalning bygger på fyra samarbetsfunktioner:

Capability Vad det gör
Fjärrhanterare för blandning (RSM) Skriver och läser shuffle-data till Azure Blob Storage i stället för körbara lokala diskar.
Blandad migrering Flyttar shuffle-block från en exekutor innan den avvecklas, i stället för att ta bort dem.
Beslutslagret Körningsroutning per fas som håller små omfördelningar lokala och flyttar över stora omfördelningar till fjärrlagring.
AQE Shuffle-skriv Gör att adaptiv frågeexekvering deltar i skrivfasen för shuffle så att partitioneringen blir korrekt direkt.

Förutsättningar

  • Aktivera Native Execution Engine (NEE).
  • Aktivera autoskalering (rekommenderas). Effektiv nedskalning fungerar också utan autoskalering via de Spark-konfigurationer som beskrivs senare i denna artikel.
  • Runtime 1.3 (Apache Spark 3.5) eller senare.

Så här fungerar det

När Spark bearbetar en förfrågan omfördelar det ofta data mellan olika steg – en shuffle. Normalt lagrar varje exekutor shuffle-data på sin lokala disk, vilket binder exekutörerna till den datan. Exekverarna kan inte frisläppas förrän varje konsument har läst klart. Denna koppling är den enskilt största anledningen till att kluster inte kan skalas ner snabbt och varför förlusten av en exekutor leder till dyra stage-retrys.

Effektiv skalning bryter den här kopplingen:

  • Large shuffles gå direkt till Azure Blob Storage via Remote Shuffle Manager.
  • Små omblandningar lagras på den lokala disken för att ge högre hastighet. Om deras exekutor senare behöver släppas, flyttar shuffle-migreringen blocken till peers eller till reservlagring i bakgrunden.
  • Beslutslagret väljer rätt vägval per steg vid körning.
  • AQE Shuffle Write säkerställer att skrivprocessen skapar en partitionering som efterföljande AQE använder utan att slå ihop partitionerna på nytt, och därmed undviker onödiga I/O-operationer.
                ┌───────────────────────────┐
   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

Smart routing (beslutslager)

Beslutslagret utvärderar varje shuffle-utbyte och avgör:

  • Stora shuffle-operationer → Azure Blob Storage. Maximal skalnings- och feltoleransförmån.
  • Små omfördelningar → lokal disk. Inga molnbaserade I/O-omkostnader för små överföringar. Om exekutorn senare avvecklas tar shuffle migration över.

Beslutslagret leder automatiskt shuffle-data och kräver ingen input från dig. Den rekommenderade kornigheten är per steg.

Viktiga fördelar

Lägre kostnader: Betala endast för den beräkning du använder

Med effektiv nedskalning släpps utförare så snart deras arbete är klart. De står inte längre sysslolösa och lagrar shuffle-data som efterföljande uppgifter så småningom kan läsa.

  • Snabbare nedskalning. Automatisk skalning tar bort noder omedelbart efter att uppgiften har slutförts.
  • Mindre outnyttjad beräkningskapacitet. Inga "zombie" exekutorer hålls vid liv bara för att tjäna sin lokala blandning.
  • Ingen överallokering av disk. Stora dataomfördelningar lagras i bloblagring i stället för att kräva stora lokala diskar.
  • Begränsad lagringskostnad. Reservlagring rensas automatiskt när blocken inte längre behövs.

Mer motståndskraftiga jobb

När shuffle-data bara finns på en lokal disk innebär en körningskrasch att data är borta och Spark måste räkna om dem. Med effektiv nedskalning finns datan antingen redan i Blob Storage eller migreras dit innan exekveraren försvinner.

Scenario Utan effektiv skalning Med effektiv nedskalning
Executor kraschar Shuffle-data förlorade; steg kördes igen Data är säkra i lagringen. ingen omkomputation
Nodpreemption Data borta, dyra återförsök Data överlever; jobbet fortsätter normalt
Graciös avveckling Shuffle togs bort vid avstängning Block migrerade till peer- eller reservlagring
Nätverksblips under hämtning Kaskadvis FetchFailedException Läsningar hämtas från lagring, utan påverkan

Denna design eliminerar den vanligaste orsaken till FetchFailedException i produktion.

Snabbare, verkligt elastisk skalning

Utan effektiv nedskalning kan autoskalaren inte återta en nod medan någon exekverare på den fortfarande har kvar shuffle-data eller cachelagrade data. Effektiv skalning frikopplar båda:

  • Shuffle-data finns i Blob Storage (eller migrerar dit vid avstängning).
  • Cacheminnet låser inte längre fast körningsprocesser. Reproducerbara cachar, till exempel Delta snapshot cache, undantas från skydd mot nedskalning.

Autoskalaren kan fritt ta bort inaktiva noder och ändra klustrets storlek när arbetsbelastningen förändras.

Bättre prestanda för skeva och stora blandningar

AQE Shuffle Write låter Adaptive Query Execution forma själva shuffle-skrivningen – genom att välja en partitionering som AQE längre ned i flödet använder utan att slå samman den igen, och genom att producera färre block i bättre storlek för fjärrlagring. I kombination med beslutslagret får du snabbare väggklocktid på stora/skeva frågor och oförändrad latens för små.

Get started

Använd den här konfigurationen för att aktivera den fullständiga effektiva skalningsstacken:

# 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")

Inga kodändringar krävs. Du kan också ange dessa i din miljö Spark-egenskaper.

Konfigurationsreferens

Fjärrhanterare för shuffle (RSM)

Setting Rekommenderad Vad den styr
spark.remote.shuffle.enabled true Slår på effektiv nedskalning. Shuffle-data går till Azure Blob Storage i stället för körbara lokala diskar.

Beslutslagret

Setting Rekommenderad Vad den styr
spark.sql.rsm.decisionlayer.enabled.level stage Granularitet med vilken beslutslagrets rutter blandas. stage utvärderar varje Spark-fas oberoende av varandra.

AQE Shuffle-skriv

Setting Rekommenderad Vad den styr
spark.sql.adaptive.shuffleWrite.enabled true Gör det möjligt för AQE att delta i shuffleskrivfasen. Skapar partitionering som efterföljande AQE använder utan att slås samman igen.

Note

själva AQE (spark.sql.adaptive.enabled) måste vara aktiverat. Den är aktiverad som standard i Fabric Spark.

Flyttning av shuffle vid avveckling

Setting Rekommenderad Vad den styr
spark.storage.decommission.shuffleBlocks.enabled true Migrerar shuffle-block från en exekutor som avvecklas, i stället för att ta bort dem.
spark.storage.decommission.shuffleBlocks.cleanup true Rensar shuffle-block på källexekutorn efter en lyckad migrering.
spark.storage.decommission.shuffleBlocks.migrateToFallbackStorage true Om ingen peer-exekverare kan ta emot blocken migreras de till reservlagring (Azure Blob Storage).
spark.storage.decommission.fallbackStorage.cleanUp true Tar bort shuffleblock från reservlagringen när de inte längre behövs och håller därmed lagringskostnaderna nere.

Cachemedveten dynamisk allokering

Setting Rekommenderad Vad den styr
spark.dynamicAllocation.preventShutdownExecutorWithCache false Tillåter dynamisk allokering att frigöra exekverare även när de har cachelagrade block.
spark.dynamicAllocation.excludeDeltaSnapshotCache true Ignorerar cachen för Delta-ögonblicksbilder när det avgörs om en exekutor fortfarande har kvar användbar cache. Cachen för delta-snapshots är reproducerbar och bör inte blockera nedskalning.

Avancerad finjustering (RSM)

De flesta användare behöver inte ändra dessa standardvärden.

Skrivprestanda

Setting Standardinställning Vad den styr
spark.remote.shuffle.partition.buffersize 16777216 (16 MB) Buffert per partition innan du skriver till lagring.
spark.remote.shuffle.blocksize 8388608 (8 MB) Storleken på enskilda block som laddats upp till Blob Storage.
spark.remote.shuffle.write.maxthreads cores × 16 Maximalt antal trådar som används för att skriva shuffle-data.
spark.remote.shuffle.write.maxtasks 16384 Maximalt antal samtidiga skrivåtgärder.

Läsprestanda

Setting Standardinställning Vad den styr
spark.remote.shuffle.read.parallel.enabled true Parallella nedladdningsströmmar vid shuffle-läsningar.
spark.remote.shuffle.read.parallelism 4 Parallella nedladdningsströmmar per aktivitet.
spark.remote.shuffle.read.prefetchqueuesize 250 Prefetch-ködjup vid läsning.
spark.remote.shuffle.read.maxthreads cores × 4 Maximalt antal trådar som används för läsning.

Reliability

Setting Standardinställning Vad den styr
spark.remote.shuffle.retries 5 Försök igen vid tillfälliga lagringsfel.
spark.remote.shuffle.retrydelayms 800 Initial väntetid mellan återförsök.
spark.remote.shuffle.retrymaxdelayms 60000 Övre gräns för backoff.

Compression

Setting Standardinställning Vad den styr
spark.remote.shuffle.compression Använder spark.io.compression.codec Komprimeringsformat för fjärrblandade data (till exempel lz4, zstd).

Prestandaresultat

Diagram som visar minskade beräkningskostnader med effektiv nedskalning aktiverad jämfört med inaktiverad i ett TPC-DS-benchmark, med en kostnadsminskning på 54 procent.

Kostnadsbesparingar för beräkningsresurser (TPC-DS-benchmark)

Metric Utan effektiv skalning Med effektiv nedskalning
Total beräkningskapacitet (VM-minuter) 14,952 6,880
Kostnadsminskning 54%

Jobbets totala körtid kan bli längre (autoskalning använder färre samtidiga exekverare), men den debiterade beräkningskapaciteten minskar med mer än hälften.

Prestanda för beslutslagret (TPC-DS, RSM on)

Att routa små shuffles till lokal disk och endast stora shuffles till fjärrlagring ger upp till 57% förbättring i körtid jämfört med att routa varje shuffle på distans, med samma nedskalningsfördel.

Limitations

  • NEE krävs. Effektiv nedskalning är beroende av Native Execution Engine.
  • Azure Blob Storage bara. Standard BlockBlobStorage med HNS inaktiverat. Azure Data Lake Gen2/HNS-aktiverade konton stöds inte som fjärrlagringsplats.
  • Stöds inte med Azure Private Link. Miljöer som använder privata länknätverk är för närvarande inte kompatibla.
  • Beslutslagretsgranulariteten är för närvarande per steg. Routning per aktivitet eller per partition finns inte i omfånget.
  • Ändring av cachebeteende. Med preventShutdownExecutorWithCache=falsekan exekutorer som innehåller cache()/persist() data skalas ned. Arbetsbelastningar som är starkt beroende av executor-local cache för frekventa data bör verifieras.