Remarque
L’accès à cette page nécessite une autorisation. Vous pouvez essayer de vous connecter ou de modifier des répertoires.
L’accès à cette page nécessite une autorisation. Vous pouvez essayer de modifier des répertoires.
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
Configuration recommandée
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
É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
BlockBlobStorageavec 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étenantcache()/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.