Un skew est un déséquilibre de la distribution des clés : une ou quelques clés portent une part disproportionnée des lignes. Après le shuffle, ces clés atterrissent dans la même partition, et une seule tâche traite l’essentiel des données.
Le signe qui ne trompe pas, dans la Spark UI : sur une étape de 200 tâches, la durée médiane est de 4 secondes et la durée maximale de 25 minutes. Le job n’est pas lent — il attend un seul exécuteur. Souvent, cette tâche finit en OutOfMemoryError ou en disk spill massif.
La cause la plus fréquente en production : une valeur sentinelle. null, 0, -1, "UNKNOWN", "N/A" regroupent des millions de lignes sous une clé unique.
1. Activer l’AQE avec la gestion du skew. Depuis Spark 3, c’est la première chose à vérifier, et elle règle une bonne partie des cas sans toucher au code.
spark.conf.set("spark.sql.adaptive.enabled", "true")
spark.conf.set("spark.sql.adaptive.skewJoin.enabled", "true")Spark détecte à l’exécution les partitions anormalement grosses et les découpe en sous-partitions traitées en parallèle.
2. Isoler les valeurs sentinelles. Si le skew vient de null, la meilleure correction est souvent de les exclure du join et de les traiter à part — une ligne dont la clé est null ne joint de toute façon avec rien d’utile.
val avecCle = grande.where($"cle".isNotNull)
val sansCle = grande.where($"cle".isNull)
avecCle.join(reference, Seq("cle")).unionByName(sansCle.withColumn(/* … */))3. Diffuser la petite table. Un broadcast join supprime le shuffle, donc supprime le skew avec lui. C’est la solution la plus efficace quand la table de référence est assez petite.
4. Le salting. Quand rien de ce qui précède ne s’applique : on ajoute un suffixe aléatoire à la clé chaude pour la répartir sur n partitions, et on réplique la petite table n fois pour que la jointure retombe juste.
val n = 50
val grandeSalee = grande.withColumn(
"cle_salee",
concat($"cle", lit("_"), (rand() * n).cast("int")),
)
val referenceRepliquee = reference
.withColumn("suffixe", explode(sequence(lit(0), lit(n - 1))))
.withColumn("cle_salee", concat($"cle", lit("_"), $"suffixe"))
grandeSalee.join(referenceRepliquee, Seq("cle_salee"))Le coût est réel : la table de référence est multipliée par 50. C’est pour cela que le salting vient en dernier, et qu’on le réserve aux clés chaudes identifiées plutôt qu’à toutes.
Une méthode, pas une liste de recettes. La réponse forte suit cet ordre : mesurer dans la Spark UI, identifier les clés chaudes avec un simple comptage, puis choisir la solution la plus légère qui traite ce cas précis.
grande.groupBy("cle").count().orderBy(desc("count")).show(20)Vingt lignes qui disent immédiatement si le problème est une clé sur un million ou une distribution naturellement inégale — et les deux ne se traitent pas de la même façon.