Vos jobs tournent avec 200 partitions par défaut, et c’est presque toujours faux. La règle des 128 Mo, repartition, coalesce et l’AQE de Spark 3.
Un job Spark qui traîne, c’est presque toujours un problème de partitions. Trop peu : le cluster attend qu’un seul exécuteur finisse un fichier de 10 Go. Trop : Spark passe plus de temps à orchestrer des micro-tâches qu’à calculer. Et au milieu de tout ça, un chiffre magique inscrit dans la config par défaut — 200 — que 90 % des utilisateurs n’ont jamais changé, alors qu’il n’est presque jamais le bon.
Cet article donne les deux ou trois règles qui, correctement appliquées, coupent souvent les temps d’exécution par deux ou trois — sans changer une ligne de logique.
Une partition est l’unité de parallélisme de Spark : un morceau de données assigné à une tâche, exécuté par un cœur d’un exécuteur. Si votre DataFrame a 200 partitions et votre cluster 400 cœurs disponibles, la moitié de vos cœurs regarderont les autres travailler.
Combien de cœurs Spark croit-il avoir, au fait ? C’est sc.defaultParallelism qui répond, et il fixe le nombre de partitions initiales de tous vos RDD — le sujet de l’article combien de processeurs Spark utilise-t-il vraiment.
Deux règles simples découlent directement de là :
Attention : plus de partitions que de cœurs n’exécute pas tout en parallèle pour autant — Spark travaille par vagues successives, expliqué dans pourquoi 100 partitions ne font pas 100 processeurs.
spark.sql.shuffle.partitions = 200Après chaque groupBy, join, distinct, orderBy, Spark redécoupe les données selon la clé pour les rassembler. Le nombre de partitions produit est fixé par une seule propriété :
spark.conf.set("spark.sql.shuffle.partitions", "200") // valeur par défaut200 a été choisi en 2015 pour un cluster de laboratoire. Aujourd’hui :
La règle usuelle, tirée de dix ans de retours terrain :
128 à 256 Mo de données par partition, en visant un multiple du nombre de cœurs disponibles.
Pour 1 To de données brutes après filtre, cela donne :
1 000 000 Mo / 200 Mo = 5 000 partitionsEt sur un cluster à 400 cœurs, on arrondit au multiple supérieur (5 200 par exemple) pour absorber les tâches lentes.
Un shuffle, c’est la sérialisation d’une partie du DataFrame, l’écriture sur le disque local des exécuteurs, la lecture réseau par les exécuteurs de destination, puis la désérialisation. Concrètement : le disque local et le réseau du cluster tournent à plein régime, pendant que le CPU attend.
Ce qui déclenche un shuffle :
| Opération | Shuffle ? |
|---|---|
select, filter, withColumn (avec une expression pure) | Non |
map, flatMap | Non |
groupBy(...).agg(...) | Oui |
join (sauf broadcast) | Oui |
distinct | Oui |
orderBy global | Oui — le plus coûteux |
repartition(n) | Oui, volontairement |
coalesce(n) avec n < partitions actuelles | Non — voir plus bas |
Cherchez Exchange dans le plan explain : chaque Exchange est un shuffle.
df.explain(true)
// == Physical Plan ==
// *(3) HashAggregate(...)
// +- Exchange hashpartitioning(prenom#7, 200) ← shuffle ici
// +- *(2) HashAggregate(...)
// +- *(1) LocalTableScanrepartition vs coalesce, la nuance qui compteDeux méthodes qui modifient le nombre de partitions ; une seule fait un shuffle. repartition(n) shuffle et équilibre. coalesce(n) fusionne sans shuffle — quasi gratuit — mais ne peut réduire qu’à condition de partir de plus de partitions que demandé ; il refuse silencieusement d’augmenter, et il hérite des déséquilibres de départ.
Piège classique : coalesce(200) en amont d’un gros groupBy laisse les tailles inégales, et une seule tâche supporte tout le poids. Le bon usage est en aval, juste avant l’écriture (.coalesce(10).write.parquet(...)) pour éviter de cracher 200 petits fichiers.
La démonstration complète, partition par partition avec glom() et le piège du coalesce(7) qui reste à 5, vit dans un article dédié : coalesce vs repartition, vu avec glom().
Depuis Spark 3, l’Adaptive Query Execution ajuste le nombre de partitions en cours de job, à partir des statistiques observées :
spark.conf.set("spark.sql.adaptive.enabled", "true")
spark.conf.set("spark.sql.adaptive.coalescePartitions.enabled", "true")
spark.conf.set("spark.sql.adaptive.skewJoin.enabled", "true")L’AQE fait trois choses très concrètes :
spark.sql.shuffle.partitions.Depuis Spark 3.2, l’AQE est activé par défaut. Vérifiez-le sur vos environnements plus anciens : sur Spark 2.4 et Spark 3.0/3.1, c’est off par défaut, et ce seul flag change souvent les temps d’exécution du simple au triple.
Dans l’UI Spark, l’onglet SQL montre le plan finalement exécuté, avec les optimisations d’AQE annotées (AQEShuffleRead, CoalescedShuffleRead). Si vous voyez Exchange en volume mais aucun AQEShuffleRead, l’AQE ne joue pas — vérifiez spark.sql.adaptive.enabled et la version de Spark.
Sur un job de 800 Go, jointure + agrégation, exécuté sur un cluster de 128 cœurs :
| Configuration | Durée | Cœurs actifs |
|---|---|---|
spark.sql.shuffle.partitions = 200 (défaut) | 47 min | 200/128 → tâches en attente |
= 800 + AQE activé | 18 min | 128/128 en continu, coalescing à la fin |
= 3 200 + AQE + broadcast forcé | 11 min | 128/128, skew détecté et découpé |
Aucune de ces trois lignes ne change la logique du calcul. Elles changent uniquement comment Spark découpe le travail. C’est là que vit la performance.
Une partition, c’est une tâche. Un shuffle, c’est un redécoupage coûteux. Le défaut 200 est une valeur historique qui ne convient à presque rien. Réglez spark.sql.shuffle.partitions pour viser 128–256 Mo par partition, activez l’AQE, utilisez coalesce uniquement pour réduire à l’écriture — et vous venez de gagner, sans effort, la majorité des optimisations que les articles Spark promettent en dix pages chacune.
Ce réglage est le préalable au suivant : les joints Spark et le pattern broadcast, sujet du prochain article de la série.
Développement et déploiement de solutions de données — les premiers modules sont en accès libre.