Partitions et shuffle en Spark : bien les régler

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.

6 min de lecturesparkscalapartitionsshufflebig-data

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.

Qu’est-ce qu’une partition, concrètement

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à :

  • Nombre de partitions < nombre de cœurs ⇒ ressources gaspillées.
  • Partitions énormes ⇒ une tâche traîne (le fameux straggler) pendant que 199 sont finies.

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.

Le chiffre magique : spark.sql.shuffle.partitions = 200

Aprè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é :

scala
spark.conf.set("spark.sql.shuffle.partitions", "200")  // valeur par défaut

200 a été choisi en 2015 pour un cluster de laboratoire. Aujourd’hui :

  • Sur un petit cluster (16 cœurs, 100 Go de données), 200 est trop : les partitions font 500 Mo, un seul cœur en digère une à la fois, ça déborde en spill.
  • Sur un gros cluster (500 cœurs, 10 To de données), 200 est ridicule : chaque partition pèse 50 Go, 300 cœurs restent inactifs.

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 partitions

Et sur un cluster à 400 cœurs, on arrondit au multiple supérieur (5 200 par exemple) pour absorber les tâches lentes.

Le shuffle, ce qu’il fait vraiment

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érationShuffle ?
select, filter, withColumn (avec une expression pure)Non
map, flatMapNon
groupBy(...).agg(...)Oui
join (sauf broadcast)Oui
distinctOui
orderBy globalOui — le plus coûteux
repartition(n)Oui, volontairement
coalesce(n) avec n < partitions actuellesNon — voir plus bas

Cherchez Exchange dans le plan explain : chaque Exchange est un shuffle.

scala
df.explain(true)
// == Physical Plan ==
// *(3) HashAggregate(...)
// +- Exchange hashpartitioning(prenom#7, 200)   ← shuffle ici
//    +- *(2) HashAggregate(...)
//       +- *(1) LocalTableScan

repartition vs coalesce, la nuance qui compte

Deux 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().

AQE : Spark 3, la fin (partielle) du réglage manuel

Depuis Spark 3, l’Adaptive Query Execution ajuste le nombre de partitions en cours de job, à partir des statistiques observées :

scala
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 :

  1. Il combine les partitions trop petites après un shuffle, sans qu’on ait à toucher spark.sql.shuffle.partitions.
  2. Il détecte le skew — quand une clé de jointure représente 80 % des données — et découpe la partition fautive automatiquement.
  3. Il change le plan de join à la volée si une des tables devient assez petite pour être broadcast (voir l’article dédié sur les joins).

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.

Comment vérifier concrètement

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.

Ce que ça donne, mesuré

Sur un job de 800 Go, jointure + agrégation, exécuté sur un cluster de 128 cœurs :

ConfigurationDuréeCœurs actifs
spark.sql.shuffle.partitions = 200 (défaut)47 min200/128 → tâches en attente
= 800 + AQE activé18 min128/128 en continu, coalescing à la fin
= 3 200 + AQE + broadcast forcé11 min128/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.

Ce qu’il faut retenir

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.

Ce sujet fait partie d’un cours complet

Développement et déploiement de solutions de données — les premiers modules sont en accès libre.

Voir le plan du cours

Continuer sur le même sujet