coalesce vs repartition en Spark, démontré avec glom()

Deux méthodes au même argument, deux travaux différents. Démonstration avec glom(), le piège du coalesce(7) qui reste à 5, et la règle pour choisir.

Deux méthodes Spark qui prennent le même argument, coalesce(n) et repartition(n), et qui font pourtant deux choses radicalement différentes. L’une est presque gratuite, l’autre déclenche un shuffle complet ; l’une peut refuser silencieusement de faire ce que vous lui demandez, l’autre non. Chaque revue de code Spark a déjà vu passer les deux confondues — et c’est presque toujours la ligne qui explique un job dix fois trop lent, ou dix mille fichiers de 20 Ko sur le stockage objet.

On va regarder à l’intérieur des partitions avec glom() sur un RDD-jouet, pour voir exactement ce que fait chaque méthode. Les nombres sont volontairement petits — les mécanismes qu’on découvre valent aussi sur un cluster à 500 exécuteurs.

L’outil qui rend tout visible : glom()

glom() transforme un RDD en RDD de listes, où chaque liste correspond à une partition. C’est l’outil de diagnostic par excellence quand on veut voir la répartition physique :

scala
val rdd = sc.parallelize(
  List(1, 2, 3, 4, 5, 6, 7, 8, 9, 10, 11),
  5
)

rdd.glom().collect()

Résultat :

Array(
  Array(1, 2),        // partition 0
  Array(3, 4),        // partition 1
  Array(5, 6),        // partition 2
  Array(7, 8),        // partition 3
  Array(9, 10, 11)    // partition 4
)

Onze éléments répartis sur cinq partitions, quasi équitablement. Spark met deux ou trois éléments par partition, sans magie particulière — il divise et distribue.

Vérifions le compte :

scala
rdd.getNumPartitions   // 5

Ici les 5 partitions sont explicites — c’est le second argument de parallelize. Sans lui, Spark aurait utilisé sc.defaultParallelism, dérivé du nombre de cœurs disponibles : voir combien de processeurs Spark utilise-t-il vraiment.

Sur ce jouet, on va faire deux expériences : réduire à 3 partitions, puis essayer d’augmenter à 7. Regardons ce qui se passe.

coalesce(3) : réduire, sans faire voyager les données

scala
rdd.coalesce(3).glom().collect()

Résultat typique :

Array(
  Array(1, 2, 3, 4),     // ancienne P0 + ancienne P1
  Array(5, 6, 7, 8),     // ancienne P2 + ancienne P3
  Array(9, 10, 11)       // ancienne P4
)

Trois partitions au lieu de cinq. Et surtout : aucun shuffle. Spark n’a pas redistribué les données via le réseau — il a simplement dit « la partition 1 est logiquement fusionnée avec la partition 0, la partition 3 avec la 2 ». Les données ne bougent pas, seule l’étiquette change.

C’est cette absence de mouvement qui rend coalesce si économique. Sur un DataFrame de 200 Go réparti sur 800 partitions, un coalesce(200) prend quelques secondes ; un repartition(200) prendrait plusieurs minutes.

Contrepartie : coalesce fusionne les partitions telles qu’elles sont. Si une des partitions d’origine était surchargée (skew), la nouvelle partition héritière sera aussi surchargée.

coalesce(7) : le piège du silence

Que se passe-t-il si vous demandez plus de partitions que ce qui existe ?

scala
rdd.coalesce(7).getNumPartitions
// 5

Cinq, pas sept. Spark n’a pas levé d’erreur, il n’a rien affiché — il vous a simplement ignoré. C’est le comportement documenté de coalesce : il n’augmente jamais le nombre de partitions sans un shuffle, et par défaut il ne fait pas de shuffle.

Le piège est classique et douloureux à trouver : vous croyez avoir passé un job de 5 à 7 tâches parallèles, vous tournez toujours sur 5 cœurs, vous ne comprenez pas pourquoi le doublement du cluster ne change rien. La bonne pratique : vérifier getNumPartitions à chaque étape critique.

scala
val bigger = rdd.coalesce(7)
println(s"partitions demandees: 7, obtenues: ${bigger.getNumPartitions}")
// partitions demandees: 7, obtenues: 5

repartition(7) : augmenter, avec shuffle assumé

scala
rdd.repartition(7).glom().collect()

Résultat (l’ordre exact varie d’une exécution à l’autre — c’est un shuffle) :

Array(
  Array(4, 9),
  Array(1, 6),
  Array(3, 11),
  Array(5),
  Array(2, 10),
  Array(7),
  Array(8)
)

Sept partitions, chacune contenant une ou deux valeurs distribuées aléatoirement — le partitionnement se fait via un RoundRobinPartitioner interne qui vise l’équilibre. Le shuffle a bien eu lieu : les données ont traversé le réseau (dans notre cas local, elles ont juste changé d’étiquette, mais sur un vrai cluster elles auraient fait un aller-retour).

scala
rdd.repartition(7).getNumPartitions
// 7

repartition tient toujours sa promesse, à la différence de coalesce. Le prix, c’est le shuffle.

coalesce(7, shuffle = true) : le même travail, dit autrement

scala
rdd.coalesce(7, shuffle = true).getNumPartitions
// 7

Avec shuffle = true, coalesce fait exactement ce que fait repartition : shuffle complet, augmentation autorisée, répartition équilibrée. En pratique, repartition(n) est coalesce(n, shuffle = true) — c’est l’implémentation interne, un simple alias explicite.

L’intérêt de connaître les deux formes : dans un code où le nombre n est calculé dynamiquement, coalesce(n, shuffle = n > current) permet d’exprimer « redécoupe seulement si nécessaire », sans avoir à if-else autour d’une méthode ou d’une autre.

Le vrai cas d’usage : juste avant write

Le scénario le plus fréquent — et le plus rentable — pour coalesce : réduire le nombre de fichiers de sortie.

scala
val df = spark.read.parquet("s3://in/events/")
  .filter($"level" === "ERROR")
  .select(colonnesUtiles: _*)
  // ...

df.coalesce(10).write.parquet("s3://out/errors/")

Sans le coalesce(10), un DataFrame issu d’un shuffle a typiquement 200 partitions par défaut ; on obtient 200 fichiers Parquet, souvent minuscules après filtrage — c’est le small files problem que payent tous les lecteurs suivants. coalesce(10) regroupe sans shuffle et produit dix fichiers propres.

Attention : ne jamais utiliser coalesce(1) sur un gros DataFrame. Cela force toute la sortie à passer par un unique exécuteur — le job se retrouve monothread, la mémoire de cet exécuteur explose, et vous avez perdu tout l’intérêt d’avoir un cluster. Pour un fichier unique, utilisez repartition(1) (avec shuffle) ou, mieux, concaténez en aval avec un outil dédié.

Table de vérité, à imprimer et à afficher au mur

CommandeShuffle ?Peut augmenter ?Peut réduire ?Répartition équilibrée ?
coalesce(n) avec n < actuelNonOuiNon (hérite des tailles d’origine)
coalesce(n) avec n > actuelNonNon — ignoré !
coalesce(n, shuffle = true)OuiOuiOuiOui
repartition(n)OuiOuiOuiOui
repartition($"col")OuiOuiOuiOui, par colonne
Et pour équilibrer sans changer le nombre de partitions ?

repartition($"col") redistribue selon les valeurs d’une colonne, en gardant le nombre de partitions à la valeur de spark.sql.shuffle.partitions. C’est ce qu’on veut typiquement avant d’écrire avec partitionBy — chaque valeur de la colonne se retrouve dans une partition unique en mémoire, ce qui produit un seul fichier par valeur sur disque.

Pour équilibrer sans regrouper par colonne :

scala
df.repartition(200)   // 200 partitions égales, aléatoirement remplies

C’est le geste qui rattrape un coalesce trop agressif ou une source déséquilibrée.

La règle simple à retenir

  • Vous voulez moins de partitions, vos données sont déjà à peu près équilibrées → coalesce(n).
  • Vous voulez plus de partitions, ou vous voulez rééquilibrer des tailles inégales → repartition(n).
  • Vous voulez une partition par valeur d’une colonne avant write.partitionByrepartition($"col").
  • Vous vous demandez pourquoi votre coalesce(n) n’a pas augmenté le compte → vérifiez getNumPartitions et passez à repartition ou ajoutez shuffle = true.

Ce qu’il faut retenir

coalesce et repartition ne sont pas deux façons de faire la même chose : ce sont deux outils qui répondent à deux problèmes différents. coalesce est l’outil de la réduction bon marché — sans shuffle, sans mouvement de données, avec ses limites (peut refuser d’augmenter, peut laisser des partitions inégales). repartition est l’outil du redécoupage propre — shuffle assumé, tailles équilibrées, promesse tenue à tous les coups.

Cet article approfondit un point brièvement abordé dans le tutoriel plus large sur les partitions et le shuffle en Spark. Pour comprendre l’impact concret sur les fichiers écrits, direction Parquet et Delta Lake, écrire proprement — la ligne coalesce(10).write.parquet(...) y prend tout son sens.

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