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 une réduction plus rapide de la taille des clusters, un coût de calcul plus faible et des tâches plus résilientes – sans aucun changement dans vos requêtes, notebooks 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

  • Activez le moteur d’exécution natif (NEE).
  • Activer la mise à l’échelle automatique (recommandé). Une réduction efficace fonctionne également sans auto-scale via les configurations Spark décrites plus loin dans cet article.
  • Runtime 1.3 (Apache Spark 3.5) ou version ultérieure.

Fonctionnement

Lorsque Spark traite une requête, il redistribue souvent les données entre les étapes - un shuffle. Normalement, chaque executor stocke les données de shuffle sur son disque local, ce qui lie les executors à ces données. Les exécuteurs ne peuvent pas être libérés tant que chaque consommateur n’a pas fini de lire. Ce couplage est la principale raison pour laquelle les clusters ne peuvent pas être réduits rapidement, et explique pourquoi la perte d’un exécuteur entraîne de coûteuses nouvelles tentatives 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 ensuite être libéré, la migration aléatoire déplace les blocs vers des pairs ou vers un stockage de secours en arrière-plan.
  • La couche de décision choisit le bon chemin par étape à 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écision)

La couche de décision évalue chaque échange de redistribution 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 est ensuite mis hors service, la migration shuffle prend le relais.

La couche de décision redistribue automatiquement les données et ne nécessite aucune action de votre part. 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.

Cette conception é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

AQE Shuffle Write permet à l’exécution adaptative des requêtes de façonner la phase d’écriture shuffle elle-même, en choisissant un partitionnement que l’AQE en aval consomme sans recoalescence, et en produisant des blocs moins nombreux et de taille plus appropriée pour le stockage distant. Combiné à la couche décision, vous obtenez un temps d’horloge mural plus rapide sur les requêtes grandes/déséquilibrées et une latence inchangée pour les requêtes petites.

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é à laquelle la couche de décision achemine le brassage. 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 en regroupement lors du désarmement

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é.

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

Acheminer de petits mélanges vers le disque local et seulement de grands mélanges vers le stockage distant offre jusqu’à 57% d’amélioration à l’exécution par rapport à chaque mélange à distance, avec le même avantage de réduction de volume.

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.