Un shuffle est la redistribution des données entre partitions, déclenchée quand une opération a besoin de rassembler toutes les valeurs d’une même clé sur la même machine — un groupBy, un join, un reduceByKey, un repartition, un distinct, un orderBy.
Il est coûteux parce qu’il cumule les trois opérations les plus lentes d’un cluster :
Et il crée une barrière : l’étape suivante ne peut pas commencer avant que la précédente soit entièrement terminée. Un seul exécuteur lent retient donc tout le job.
Que vous savez distinguer les opérations qui shufflent de celles qui ne shufflent pas — c’est la compétence pratique derrière la question.
Une transformation étroite (narrow) produit chaque partition de sortie à partir d’une seule partition d’entrée : map, filter, select, withColumn, coalesce. Aucun mouvement de données, tout s’exécute dans la même étape.
Une transformation large (wide) a besoin de plusieurs partitions d’entrée pour chaque partition de sortie : groupByKey, reduceByKey, join, repartition, distinct. Elle impose un shuffle, donc une nouvelle étape.
La conséquence à énoncer : le nombre d’étapes d’un job est le nombre de shuffles plus un. C’est ainsi qu’on lit une interface Spark UI.
Quatre leviers, du plus efficace au plus fin :
where placé avant le join réduit le volume déplacé.broadcast) pour supprimer entièrement le shuffle d’un join.reduceByKey à groupByKey : l’agrégation partielle se fait côté source, donc il y a moins à transférer.spark.sql.shuffle.partitions : 200 par défaut, ce qui est absurde sur 10 Go comme sur 10 To.« Comment voyez-vous qu’un shuffle est le problème ? »
Dans la Spark UI : une étape dont les métriques Shuffle Read et Shuffle Write sont élevées, et dont la durée de la tâche médiane est très inférieure à celle de la tâche la plus longue — signe d’un déséquilibre des clés.
Le détail du dimensionnement est développé dans partitions et shuffle en Spark : bien les régler.