Effektiv nedskalering og fjernstyring av stokking

Gjelder for:✅ Fabric Data Engineering og Data Science

Effektiv nedskalering er en funksjon i Microsoft Fabric Spark som frikobler Spark shuffle-data fra eksekutørens levetid. I stedet for å feste shuffle-utgangen til lokale executor-disker, ruter Fabric Spark shuffle-data til Azure Blob Storage (eller migrerer dem dit på forespørsel) og lar Adaptive Query Execution (AQE) forme selve skrivingen. Resultatet er raskere nedskalering av klynger, lavere beregningskostnader og mer robuste jobber – uten endringer i spørringer, notatbøker eller pipelines.

Oversikt

Effektiv nedskalering er bygget opp fra fire samarbeidende egenskaper:

Kapabilitet Hva det gjør
Remote Shuffle Manager (RSM) Skriver og leser og stokker data til Azure Blob Storage i stedet for til executor lokale disker.
Shuffle-migrasjon Bevegelser stokker blokker en eksekutør før den blir avviklet, i stedet for å droppe dem.
Beslutningslag Per-trinns runtime-ruting som holder små stokkinger lokalt og overfører store stokkinger til fjernlagring.
AQE Shuffle Skriv La adaptiv spørringsutførelse delta i shuffle-skrivefasen slik at partisjoneringen er riktig første gang.

Forutsetninger

  • Aktiver Native Execution Engine (NEE).
  • Aktiver autoskalering (anbefalt). Effektiv nedskalering fungerer også uten autoskalering via Spark-konfigurasjonene som beskrives senere i denne artikkelen.
  • Kjøretid 1.3 (Apache Spark 3.5) eller nyere.

Slik fungerer det

Når Spark behandler en spørring, omfordeler den ofte data mellom trinnene – en stokking. Vanligvis lagrer hver eksekutor stokkingsdata på sin lokale disk, som knytter eksekutørene til disse dataene. Bobestyrerne kan ikke frigjøres før alle kunder er ferdige med å lese. Denne koblingen er den største grunnen til at klynger ikke kan skaleres raskt ned, og hvorfor tap av en executor fører til dyre stage-retries (stage-repres).

Effektiv nedskalering bryter denne koblingen:

  • Store stokkinger gå direkte til Azure Blob Storage via Remote Shuffle Manager.
  • Små stokkinger blir værende på lokal disk for hastighet. Hvis eksekutoren deres senere må frigjøres, flytter shuffle-migreringen blokkene til peers eller til fallback-lagring i bakgrunnen.
  • Beslutningslaget velger riktig vei per trinn ved kjøring.
  • AQE Shuffle Write sikrer at skriveren produserer partisjonering som AQE nedstrøms bruker uten å koaleres, og unngår bortkastet 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

Smart ruting (beslutningslag)

Beslutningslaget evaluerer hver stokkingsutveksling og avgjør:

  • Store bevegelser → Azure Blob Storage. Maksimal nedskalering og feiltoleransefordel.
  • Små stokkinger → lokal disk. Ingen sky-I/O-overhead for små overføringer. Hvis eksekutoren senere avvikler, tar shuffle migration over.

Beslutningslaget ruter data automatisk og krever ingen input fra deg. Den anbefalte granulariteten er per trinn.

Hovedfordeler

Lavere kostnader: Betal kun for beregningen du bruker

Med effektiv nedskalering blir eksekutorer frigitt så snart arbeidet er gjort. De står ikke lenger ubrukt og holder stokkingsdata som oppgaver nedstrøms kanskje til slutt leser.

  • Raskere nedskalering. Autoscale fjerner noder umiddelbart etter at oppgaven er fullført.
  • Mindre inaktiv beregning. Ingen «zombie»-eksekutorer holdt i live kun for å servere sin lokale omrokering.
  • Ingen disk-overprovisionering. Store stokkinger går til blob-lagring i stedet for å kreve store lokale disker.
  • Begrenset lagringskostnad. Fallback-lagring ryddes automatisk opp når blokker ikke lenger trengs.

Mer robuste jobber

Når shuffle-data kun ligger på lokal disk, betyr et executor-krasj at dataene er borte og Spark må beregne dem på nytt. Med effektiv nedskalering er data enten allerede i blob-lagring eller migrert dit før eksekutoren forsvinner.

Scenario Uten effektiv nedskalering Med effektiv nedskalering
Eksekutorkrasj Shuffle data mistet; Stadier utført på nytt Data er trygge i lagring; ingen omberegning
Nodefortrinn Data borte, dyre forsøk Data overlever; Jobben fortsetter som vanlig
Elegant avvikling Shuffle droppet ved nedstengning Blokker migrert til peer- eller fallback-lagring
Nettverksblips under henting Kaskadering FetchFailedException Avlesninger kommer fra lagring, upåvirket

Dette designet eliminerer den vanligste årsaken til FetchFailedException produksjon.

Raskere, virkelig elastisk skalering

Uten effektiv nedskalering kan ikke autoscaleren gjenerobre en node mens enhver utøver på den fortsatt holder shuffle-data eller bufret data. Effektiv nedskalering kobler fra begge:

  • Shuffle-data er i blob-lagring (eller migrerer dit ved nedstengning).
  • Cache fester ikke lenger eksekutorer. Reproduserbare cacher som Delta snapshot-cache er unntatt fra nedskaleringsbeskyttelse.

Autoscaleren kan fritt fjerne inaktive noder og endre størrelsen på klyngen som respons på endringer i arbeidsmengden.

Bedre ytelse på skjeve og store stokkinger

AQE Shuffle Write lar adaptiv spørringsutførelse forme selve shuffle-skrivingen – ved å velge partisjonering som nedstrøms AQE bruker uten å samle seg på nytt, og produsere færre, bedre blokker for fjernlagring. Kombinert med beslutningslaget får du raskere veggklokketid på store/skjeve forespørsler og uendret forsinkelse for små.

Get started

Bruk denne konfigurasjonen for å aktivere den fulle, effektive nedskaleringsstakken:

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

Ingen kodeendringer er nødvendige. Du kan også sette disse i Spark-egenskapene i miljøet ditt.

Konfigurasjonsreferanse

Remote Shuffle Manager (RSM)

Setting Anbefalt Hva den kontrollerer
spark.remote.shuffle.enabled true Slår på effektiv nedskalering. Shuffle-data går til Azure Blob Storage i stedet for til executor lokale disker.

Beslutningslag

Setting Anbefalt Hva den kontrollerer
spark.sql.rsm.decisionlayer.enabled.level stage Granulariteten som beslutningslagets ruter omrokerer med. stage vurderer hvert Spark-stadium uavhengig.

AQE Shuffle Skriv

Setting Anbefalt Hva den kontrollerer
spark.sql.adaptive.shuffleWrite.enabled true La AQE delta i shuffle write-fasen. Fører til partisjonering som nedstrøms AQE forbruker uten å re-koaleschere.

Note

AQE selv (spark.sql.adaptive.enabled) må være på. Den er på som standard i Fabric Spark.

Shuffle migrasjon ved avvikling

Setting Anbefalt Hva den kontrollerer
spark.storage.decommission.shuffleBlocks.enabled true Flytter blokkerer en eksekutor som avvikler, i stedet for å droppe dem.
spark.storage.decommission.shuffleBlocks.cleanup true Rydder opp i shuffle-blokkeringer på kildeutøveren etter en vellykket migrering.
spark.storage.decommission.shuffleBlocks.migrateToFallbackStorage true Hvis ingen peer-utøver kan akseptere blokkene, migrerer den dem til fallback-lagring (Azure Blob Storage).
spark.storage.decommission.fallbackStorage.cleanUp true Fjerner stokkeblokker fra fallback-lagring når de ikke lenger trengs, og begrenser lagringskostnadene.

Cache-bevisst dynamisk allokering

Setting Anbefalt Hva den kontrollerer
spark.dynamicAllocation.preventShutdownExecutorWithCache false Tillater dynamisk allokering for å frigi eksekutorer selv når de holder bufrede blokker.
spark.dynamicAllocation.excludeDeltaSnapshotCache true Ignorerer Delta-snapshot-cache når den skal avgjøre om en eksekutor fortsatt har nyttig cache. Delta snapshot-cache er reproduserbar og bør ikke blokkere nedskalering.

Avansert stemming (RSM)

De fleste brukere trenger ikke å endre disse standardinnstillingene.

Skriveytelse

Setting Forhåndsinnstilt Hva den kontrollerer
spark.remote.shuffle.partition.buffersize 16777216 (16 MB) Buffer per partisjon før du skriver til lagring.
spark.remote.shuffle.blocksize 8388608 (8 MB) Størrelsen på individuelle blokker lastet opp til Blob Storage.
spark.remote.shuffle.write.maxthreads cores × 16 Maksimalt antall tråder brukt til å skrive stokkingsdata.
spark.remote.shuffle.write.maxtasks 16384 Maksimale samtidige skriveoperasjoner.

Leste opptredener

Setting Forhåndsinnstilt Hva den kontrollerer
spark.remote.shuffle.read.parallel.enabled true Parallelle nedlastingsstrømmer for shuffle-lesninger.
spark.remote.shuffle.read.parallelism 4 Parallelle nedlastingsstrømmer per oppgave.
spark.remote.shuffle.read.prefetchqueuesize 250 Dybde på forhåndshentekøen under lesing.
spark.remote.shuffle.read.maxthreads cores × 4 Maksimalt antall tråder brukt til lesing.

Reliability

Setting Forhåndsinnstilt Hva den kontrollerer
spark.remote.shuffle.retries 5 Prøv på nytt på transient lagringsfeil.
spark.remote.shuffle.retrydelayms 800 Første tilbakegang mellom forsøkene.
spark.remote.shuffle.retrymaxdelayms 60000 Backoff cap.

Komprimering

Setting Forhåndsinnstilt Hva den kontrollerer
spark.remote.shuffle.compression Bruksområder spark.io.compression.codec Komprimeringsformat for fjern stokkingdata (for eksempel lz4, zstd).

Ytelsesresultater

Diagram som viser beregningskostnadsbesparelser med effektiv nedskalering aktivert versus deaktivert på en TPC-DS benchmark, som viser 54 prosent kostnadsreduksjon.

Beregn kostnadsbesparelser (TPC-DS benchmark)

Beregning Uten effektiv nedskalering Med effektiv nedskalering
Total beregning (VM-Minutes) 14,952 6,880
Kostnadsreduksjon 54%

Den totale jobbkjøretiden kan være lengre (autoscale bruker færre samtidige utøvere), men fakturert beregning kuttes med mer enn halvparten.

Beslutningslagsytelse (TPC-DS, RSM on)

Å rute små stokkinger til lokal disk og bare store stokkinger til fjernlagring gir opptil 57% kjøretidsforbedring sammenlignet med å rute hver stokking eksternt, med samme nedskaleringsfordel.

Begrensninger

  • NEE kreves. Effektiv nedskalering avhenger av den native kjøringsmotoren.
  • Azure Blob Storage only. Standard BlockBlobStorage med HNS deaktivert. Azure Data Lake Gen2 / HNS-aktiverte kontoer støttes ikke som ekstern shuffle-butikk.
  • Ikke støttet med Azure Private Link. Miljøer som bruker privat nettverk er for øyeblikket ikke kompatible.
  • Beslutningslagets granularitet er for øyeblikket per trinn. Ruting per oppgave eller per partisjon er ikke innenfor omfanget.
  • Endring i cache-atferd. Med preventShutdownExecutorWithCache=false, kan eksekutører som holder cache()/persist() data skaleres ned. Arbeidsbelastninger som er sterkt avhengige av executor-lokal cache for hot data bør valideres.