Réduction efficace et gestionnaire de shuffle à distance

S'applique à :✅ Fabric Data Engineering and Data Science

La réduction efficace de capacité est une fonctionnalité de Microsoft Fabric Spark qui dissocie les données de shuffle Spark de la durée de vie des exécuteurs. Au lieu de stocker la sortie du shuffle sur les disques locaux des exécuteurs, Fabric Spark achemine les données de shuffle vers Stockage Blob Azure (ou les y migre à la demande) et permet à l’exécution adaptative des requêtes (AQE) de définir la manière dont l’écriture est effectuée. Le résultat est un scale-down de cluster plus rapide, un coût de calcul inférieur et des travaux plus résilients, sans aucune modification de vos requêtes, blocs-notes ou pipelines.

Overview

Une mise à l’échelle efficace repose sur quatre fonctionnalités de coopération :

Capacité Qu’est-ce que cela fait ?
Gestionnaire de shuffle à distance (RSM) Écrit et lit les données de shuffle dans Stockage Blob Azure au lieu des disques locaux de l’exécuteur.
Migration aléatoire Retire les blocs de shuffle d’un exécuteur avant sa mise hors service, au lieu de les supprimer.
Couche de décision Routage à l’exécution à chaque étape, qui garde les petites redistributions en local et déporte les grandes redistributions vers un stockage distant.
Écriture aléatoire AQE Permettre à l’exécution adaptative des requêtes de participer à la phase d’écriture du shuffle afin que le partitionnement soit correct dès le départ.

Prerequisites

  • Le moteur d’exécution natif (NEE) doit être activé.
  • Mise à l’échelle automatique activée (recommandé). Une mise à l’échelle efficace fonctionne également sans mise à l’échelle automatique via les configurations Spark ci-dessous.
  • Runtime 1.3 (Apache Spark 3.5) ou version ultérieure.

Fonctionnement

Lorsque Spark traite une requête, il redistribue souvent des données entre les étapes — un shuffle. Normalement, les données de shuffle sont stockées sur le disque local de chaque exécuteur, ce qui lie les exécuteurs à ces données. Ils ne peuvent pas être libérés tant que chaque consommateur n’a pas terminé la lecture. Ce couplage est la principale raison pour laquelle les clusters ne peuvent pas diminuer rapidement en taille et pour laquelle la perte d’un exécuteur entraîne de coûteuses relances d’étape.

Une réduction d’échelle efficace rompt ce couplage :

  • Large shuffles accédez directement à Stockage Blob Azure via le Gestionnaire de shuffle distant.
  • Petites opérations de mélange restent sur le disque local pour des raisons de performance. Si leur exécuteur doit être libéré ultérieurement, shuffle Migration déplace les blocs vers des homologues ou vers le stockage de secours en arrière-plan.
  • La couche de décision choisit le chemin d’accès approprié par étape au moment de l’exécution.
  • AQE Shuffle Write garantit que le processus d’écriture produit un partitionnement que l’AQE en aval consomme sans nouveau regroupement, évitant ainsi des E/S inutiles.
                ┌───────────────────────────┐
   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

Routage intelligent (couche décisionnaire)

La couche de décision évalue chaque échange de shuffle et décide :

  • Brassages volumineux → Stockage Blob Azure. Avantage maximal de scale-down et de tolérance de panne.
  • Petits remaniements → disque local. Aucune surcharge d’E/S dans le cloud pour les transferts de petite taille. Si l’exécuteur désactive ultérieurement, Shuffle Migration prend le relais.

Le routage est automatique et ne nécessite aucune entrée utilisateur. La granularité recommandée est par étape.

Principaux avantages

Coûts inférieurs : payer uniquement pour le calcul que vous utilisez

Avec une mise à l’échelle efficace, les exécuteurs sont libérés dès que leur travail est effectué. Ils ne restent plus inactifs à conserver des données de shuffle que les tâches en aval pourraient éventuellement lire.

  • Réduction plus rapide. La mise à l’échelle automatique supprime les nœuds immédiatement après l’achèvement de la tâche.
  • Moins de ressources de calcul inutilisées. Aucun exécuteur « zombie » n’est maintenu en vie uniquement pour assurer son shuffle local.
  • Aucun surapprovisionnement de disque. Les opérations de redistribution volumineuses sont dirigées vers le stockage Blob au lieu de nécessiter de grands disques locaux.
  • Coût de stockage limité. Le stockage de secours est nettoyé automatiquement lorsque les blocs ne sont plus nécessaires.

Travaux plus résilients

Lorsque les données de shuffle résident uniquement sur le disque local, une panne d’un exécuteur signifie que ces données sont perdues et que Spark doit les recalculer. Avec une réduction d’échelle efficace, les données se trouvent déjà dans le stockage Blob ou y sont migrées avant la disparition de l’exécuteur.

Scénario Sans réduction d’échelle efficace Avec une réduction efficace
Blocage de l’exécuteur Données de shuffle perdues ; étapes réexécutées Les données sont sécurisées dans le stockage ; aucune recomputation
Préemption de nœud Données supprimées, nouvelles tentatives coûteuses Les données survivent ; le travail continue normalement
Mise hors service en douceur Lecture aléatoire désactivée à l’arrêt Blocs migrés vers le stockage homologue ou de secours
Blips réseau lors de l’extraction En cascade FetchFailedException Les lectures proviennent du stockage et ne sont pas affectées.

Cela élimine la cause la plus fréquente de FetchFailedException en production.

Montée en charge plus rapide et véritablement élastique

Sans réduction efficace, l’autoscaler ne peut pas libérer un nœud tant qu’un exécuteur sur ce nœud détient encore des données de shuffle ou des données en cache. Une mise à l’échelle efficace dissocie les deux :

  • Les données de shuffle se trouvent dans le stockage Blob (ou y sont migrées lors de l’arrêt).
  • Le cache n’immobilise plus les exécuteurs. Les caches reproductibles, tels que le cache d’instantanés Delta, sont exclus de la protection contre la réduction de capacité.

Le générateur de mise à l’échelle automatique peut supprimer librement les nœuds inactifs et redimensionner le cluster en réponse aux modifications de charge de travail.

Meilleures performances pour les opérations de redistribution déséquilibrées et volumineuses

L’écriture de shuffle d’AQE permet à l’exécution adaptative des requêtes de façonner l’écriture de shuffle elle-même, en choisissant un partitionnement que l’AQE en aval peut exploiter sans recoalescence supplémentaire, et en produisant moins de blocs, mais de meilleure taille, pour le stockage distant. Associé à la couche de décision, vous bénéficiez d’un temps d’exécution réel réduit pour les requêtes volumineuses ou déséquilibrées, tout en conservant une latence inchangée pour les petites requêtes.

Commencez

Appliquez cette configuration pour activer l’ensemble complet et efficace de réduction d’échelle :

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

Aucune modification du code n’est requise. Vous pouvez également les définir dans vos propriétés Spark d’environnement.

Référence de configuration

Gestionnaire de shuffle à distance (RSM)

Setting Recommandé Qu’est-ce qu’il contrôle ?
spark.remote.shuffle.enabled true Active un scale-down efficace. Les données de shuffle sont stockées dans Stockage Blob Azure au lieu des disques locaux des exécuteurs.

Couche de décision

Setting Recommandé Qu’est-ce qu’il contrôle ?
spark.sql.rsm.decisionlayer.enabled.level stage Granularité au niveau de laquelle la couche de décision achemine le shuffle. stage évalue chaque étape Spark indépendamment.

Écriture aléatoire AQE

Setting Recommandé Qu’est-ce qu’il contrôle ?
spark.sql.adaptive.shuffleWrite.enabled true Permet à AQE de participer à la phase d’écriture aléatoire. Produit un partitionnement utilisé par l’AQE en aval sans recoalescence.

Note

L’AQE lui-même (spark.sql.adaptive.enabled) doit être activé. Elle est activée par défaut dans Fabric Spark.

Migration aléatoire lors de la mise hors service

Setting Recommandé Qu’est-ce qu’il contrôle ?
spark.storage.decommission.shuffleBlocks.enabled true Migre les blocs de shuffle d’un exécuteur en cours de mise hors service, au lieu de les abandonner.
spark.storage.decommission.shuffleBlocks.cleanup true Nettoie les blocs de shuffle sur l’exécuteur source après une migration réussie.
spark.storage.decommission.shuffleBlocks.migrateToFallbackStorage true Si aucun exécuteur homologue ne peut accepter les blocs, les migre vers le stockage de secours (Stockage Blob Azure).
spark.storage.decommission.fallbackStorage.cleanUp true Supprime les blocs de shuffle du stockage de secours une fois qu’ils ne sont plus nécessaires, ce qui permet de limiter les coûts de stockage.

Allocation dynamique prenant en compte le cache

Setting Recommandé Qu’est-ce qu’il contrôle ?
spark.dynamicAllocation.preventShutdownExecutorWithCache false Permet l’allocation dynamique pour libérer les exécuteurs même lorsqu’ils contiennent des blocs mis en cache.
spark.dynamicAllocation.excludeDeltaSnapshotCache true Ignore le cache d’instantané Delta lorsque vous décidez si un exécuteur contient toujours un cache utile. Le cache d’instantanés delta est reproductible et ne doit pas bloquer la réduction du dimensionnement.

Réglage avancé (RSM)

La plupart des utilisateurs n’ont pas besoin de modifier ces valeurs par défaut.

Performances en écriture

Setting Par défaut Qu’est-ce qu’il contrôle ?
spark.remote.shuffle.partition.buffersize 16777216 (16 Mo) Mémoire tampon par partition avant d’écrire dans le stockage.
spark.remote.shuffle.blocksize 8388608 (8 Mo) Taille des blocs individuels chargés sur Stockage Blob.
spark.remote.shuffle.write.maxthreads cores × 16 Nombre maximal de threads utilisés pour écrire les données de shuffle.
spark.remote.shuffle.write.maxtasks 16384 Nombre maximal d’opérations d’écriture simultanées.

Lire les performances

Setting Par défaut Qu’est-ce qu’il contrôle ?
spark.remote.shuffle.read.parallel.enabled true Flux de téléchargement parallèles pour les lectures aléatoires.
spark.remote.shuffle.read.parallelism 4 Flux de téléchargement parallèles par tâche.
spark.remote.shuffle.read.prefetchqueuesize 250 Profondeur de la file d’attente de prérécupération lors des lectures.
spark.remote.shuffle.read.maxthreads cores × 4 Nombre maximal de threads utilisés pour la lecture.

Fiabilité

Setting Par défaut Qu’est-ce qu’il contrôle ?
spark.remote.shuffle.retries 5 Nouvelles tentatives sur les erreurs de stockage temporaires.
spark.remote.shuffle.retrydelayms 800 Délai initial entre deux tentatives.
spark.remote.shuffle.retrymaxdelayms 60000 Limite de repli.

Compression

Setting Par défaut Qu’est-ce qu’il contrôle ?
spark.remote.shuffle.compression Utilise spark.io.compression.codec Format de compression pour les données de shuffle distantes (par exemple, lz4, zstd).

Résultats des performances

Graphique montrant les économies de calcul avec une mise à l’échelle efficace activée ou désactivée sur un benchmark de TPC-DS, montrant une réduction des coûts de 54 %.

Économies de calcul (TPC-DS benchmark)

Metric Sans réduction d’échelle efficace Avec une réduction efficace
Capacité de calcul totale (VM-Minutes) 14,952 6,880
Réduction des coûts 54%

Le runtime total du travail peut être plus long (la mise à l’échelle automatique utilise moins d’exécuteurs simultanés), mais le calcul facturé est réduit de plus de la moitié.

Performances de la couche de décision (TPC-DS, RSM activé)

L’acheminement des petits shuffles vers le disque local et des shuffles volumineux uniquement vers le stockage distant permet d’améliorer le temps d’exécution jusqu’à 57 % par rapport à l’acheminement de tous les shuffles vers le stockage distant, avec le même avantage en matière de réduction à l’échelle.

Limitations

  • NEE obligatoire. Une mise à l’échelle efficace dépend du moteur d’exécution natif.
  • Stockage Blob Azure uniquement. Standard BlockBlobStorage avec HNS désactivé. Les comptes Azure Data Lake Gen2 / avec HNS activé ne sont pas pris en charge comme magasin de shuffle distant.
  • Non pris en charge avec Azure Private Link. Les environnements utilisant la mise en réseau de liaison privée ne sont pas actuellement compatibles.
  • La granularité de la couche de décision est actuellement par étape. Le routage par tâche ou par partition ne fait pas partie du périmètre.
  • Modification du comportement du cache. Avec preventShutdownExecutorWithCache=false, les exécuteurs détenant cache()/persist() données peuvent être réduits. Les charges de travail qui dépendent fortement du cache local de l’exécuteur pour les données chaudes doivent être validées.