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.
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 :
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 :
rdd.getNumPartitions // 5Ici 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éesrdd.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 silenceQue se passe-t-il si vous demandez plus de partitions que ce qui existe ?
rdd.coalesce(7).getNumPartitions
// 5Cinq, 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.
val bigger = rdd.coalesce(7)
println(s"partitions demandees: 7, obtenues: ${bigger.getNumPartitions}")
// partitions demandees: 7, obtenues: 5repartition(7) : augmenter, avec shuffle assumé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).
rdd.repartition(7).getNumPartitions
// 7repartition tient toujours sa promesse, à la différence de coalesce. Le prix, c’est le shuffle.
coalesce(7, shuffle = true) : le même travail, dit autrementrdd.coalesce(7, shuffle = true).getNumPartitions
// 7Avec 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.
writeLe scénario le plus fréquent — et le plus rentable — pour coalesce : réduire le nombre de fichiers de sortie.
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é.
| Commande | Shuffle ? | Peut augmenter ? | Peut réduire ? | Répartition équilibrée ? |
|---|---|---|---|---|
coalesce(n) avec n < actuel | Non | – | Oui | Non (hérite des tailles d’origine) |
coalesce(n) avec n > actuel | Non | Non — ignoré ! | – | – |
coalesce(n, shuffle = true) | Oui | Oui | Oui | Oui |
repartition(n) | Oui | Oui | Oui | Oui |
repartition($"col") | Oui | Oui | Oui | Oui, par colonne |
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 :
df.repartition(200) // 200 partitions égales, aléatoirement rempliesC’est le geste qui rattrape un coalesce trop agressif ou une source déséquilibrée.
coalesce(n).repartition(n).write.partitionBy → repartition($"col").coalesce(n) n’a pas augmenté le compte → vérifiez getNumPartitions et passez à repartition ou ajoutez shuffle = true.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.
Développement et déploiement de solutions de données — les premiers modules sont en accès libre.