Un join qui met une heure au lieu de trois minutes, c’est la stratégie choisie. Broadcast, sort-merge, skew : quand chacun s’applique et comment forcer.
Un join en Spark ressemble à un join en SQL — jusqu’à ce que vous regardiez l’onglet Stages. Là, un job de trois minutes annoncées passe une heure entière, ou coince à 99 % avec une seule tâche qui refuse de finir. Le problème n’est presque jamais votre SQL : c’est la stratégie de jointure choisie par Spark, et parfois la distribution des clés.
Cet article donne les quatre stratégies que Spark connaît, quand chacune s’applique, comment lire le plan pour savoir laquelle a été retenue, et les deux gestes qui rattrapent 90 % des joins pathologiques.
| Stratégie | Coût | Quand Spark la choisit |
|---|---|---|
| Broadcast Hash Join | Le moins cher, sans shuffle sur la petite table | Une table est « suffisamment petite » (par défaut ≤ 10 Mo) |
| Sort-Merge Join | Shuffle des deux côtés + tri | Les deux tables sont grandes, clé triable |
| Shuffle Hash Join | Shuffle des deux côtés, pas de tri | Rare, souvent forcé manuellement |
| Broadcast Nested Loop Join | Catastrophique | Aucune clé d’égalité, ou fallback |
Le nom qui apparaît dans l’UI Spark dit exactement ce qui s’est passé. C’est la première chose à lire quand un job traîne :
big.join(small, Seq("id")).explain(true)
// == Physical Plan ==
// *(2) BroadcastHashJoin [id#0], [id#5], Inner, BuildRight ← type de join
// :- *(2) Filter isnotnull(id#0)
// : +- *(2) FileScan parquet ...
// +- BroadcastExchange HashedRelationBroadcastMode ...Une des deux tables tient en mémoire de chaque exécuteur ? Spark la copie sur tous les nœuds, puis fait la jointure localement — sans aucun shuffle de la grosse table. Sur un job qui joint un fait de 500 Go avec une dimension de 5 Mo, la différence est du simple au dix.
import org.apache.spark.sql.functions.broadcast
val enriched = fact.join(broadcast(dim), "id")broadcast() hint Spark : « fais un broadcast join, la petite table tient. » C’est explicite, ça survit aux optimisations de Catalyst, et c’est presque toujours ce qu’il faut écrire quand vous savez qu’une table est petite.
Sans hint, Spark décide en fonction d’un seuil configurable :
spark.conf.set("spark.sql.autoBroadcastJoinThreshold", 100 * 1024 * 1024) // 100 MoLe défaut (10 Mo) est conservateur. Sur un cluster moderne où chaque exécuteur a 8 à 16 Go de mémoire, monter à 100 ou 200 Mo est raisonnable et débloque beaucoup de joins qui auraient été des sort-merge inutilement.
Depuis Spark 3, l’AQE peut convertir un sort-merge en broadcast à la volée si le filtre en amont réduit une des tables sous le seuil. Vous verrez alors AQEShuffleRead puis BroadcastHashJoin dans le plan final. C’est un cas où ne rien forcer laisse Spark faire le bon choix — à condition que l’AQE soit actif.
Deux tables trop grosses pour un broadcast ? Spark :
C’est robuste et scalable, mais coûteux : deux shuffles complets. Pour l’accélérer sans changer le type de join, deux leviers :
repartition("id") ou écrire avec bucketBy("id")),select avant le join) — le shuffle transporte moins d’octets.Comme le sort-merge, mais sans le tri : Spark construit une hash-map côté droite après le shuffle. Économique en CPU (pas de tri), plus gourmand en mémoire. Spark ne le choisit presque jamais tout seul — vous pouvez le forcer avec le hint SHUFFLE_HASH, mais dans 95 % des cas, laissez-le tranquille : le sort-merge est plus prévisible.
Un join sur user_id où un seul utilisateur représente 40 % des événements ? Après le shuffle, une seule partition hérite de la moitié du travail. 199 tâches finissent en 2 minutes, la 200ᵉ tourne pendant 3 heures. C’est le skew, et c’est la cause de la moitié des tickets « le job est lent » en production.
Trois signaux dans l’UI Spark :
SortMergeJoin dont un côté a une répartition très déséquilibrée.Trois remèdes, du plus simple au plus lourd :
Activer l’AQE skew join — sur Spark 3, un seul flag suffit dans la majorité des cas :
spark.conf.set("spark.sql.adaptive.enabled", "true")
spark.conf.set("spark.sql.adaptive.skewJoin.enabled", "true")Spark détecte les partitions plus grandes que la médiane et les découpe en sous-partitions au vol.
Salter la clé — technique manuelle pour les cas que l’AQE ne rattrape pas :
import org.apache.spark.sql.functions._
val salted = big
.withColumn("salt", (rand() * 10).cast("int"))
.withColumn("id_salt", concat($"id", lit("_"), $"salt"))
val exploded = small
.withColumn("salt", explode(sequence(lit(0), lit(9))))
.withColumn("id_salt", concat($"id", lit("_"), $"salt"))
salted.join(exploded, "id_salt")On ajoute un « sel » aléatoire à la grosse table, on duplique la petite table sur toutes les valeurs de sel, on rejoint sur la clé composée. Le travail est réparti sur 10 partitions au lieu d’une seule. Coûteux à écrire, mais imbattable sur un skew extrême.
Filtrer les valeurs pathologiques et les traiter à part — quand un null ou une clé sentinelle représente 30 % des lignes, le plus simple reste de la retirer, faire le join sur le reste, et gérer le cas particulier séparément.
Si vous devez retenir deux choses de tout cet article :
Ces deux réglages sont gratuits, presque toujours sûrs, et coupent des heures d’exécution sur les pipelines qui n’ont jamais été audités. C’est ce que fait un DBA Spark expérimenté dans les cinq premières minutes d’une revue de performance.
Le SQL du join est identique dans tous les cas ; c’est la stratégie physique qui décide de tout. Lisez l’UI Spark avant de suspecter votre logique : le nom du join (BroadcastHashJoin, SortMergeJoin) est écrit noir sur blanc, et il vous dit à lui seul ce qu’il faut ajuster. Sur un pipeline sain, la plupart des jointures sont soit broadcast, soit sort-merge propres — pas de tâche à 45 minutes qui bloque le reste.
Sujet lié directement : les partitions et le shuffle, qui déterminent le coût réel de chaque join que vous écrivez.
Développement et déploiement de solutions de données — les premiers modules sont en accès libre.