Joins en Spark : broadcast, shuffle et skew

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.

6 min de lecturesparkscalajoinshufflebig-data

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.

Les quatre stratégies

StratégieCoûtQuand Spark la choisit
Broadcast Hash JoinLe moins cher, sans shuffle sur la petite tableUne table est « suffisamment petite » (par défaut ≤ 10 Mo)
Sort-Merge JoinShuffle des deux côtés + triLes deux tables sont grandes, clé triable
Shuffle Hash JoinShuffle des deux côtés, pas de triRare, souvent forcé manuellement
Broadcast Nested Loop JoinCatastrophiqueAucune 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 :

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

Broadcast hash join : le graal, quand c’est possible

Une des deux tables tient en mémoire de chaque exécuteur ? Spark la copie sur tous les nœuds, puis fait la jointure localementsans 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.

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

scala
spark.conf.set("spark.sql.autoBroadcastJoinThreshold", 100 * 1024 * 1024)  // 100 Mo

Le 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.

Attention : broadcast + AQE = décision dynamique

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.

Sort-merge join : le cas général des grandes tables

Deux tables trop grosses pour un broadcast ? Spark :

  1. shuffle les deux tables par la clé de jointure,
  2. trie chaque partition,
  3. balaye les deux côtés en parallèle pour émettre les paires qui matchent.

C’est robuste et scalable, mais coûteux : deux shuffles complets. Pour l’accélérer sans changer le type de join, deux leviers :

  • Pré-partitionner les tables par la clé de join (repartition("id") ou écrire avec bucketBy("id")),
  • Réduire les colonnes portées (select avant le join) — le shuffle transporte moins d’octets.

Le shuffle hash join, presque jamais automatique

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.

Le skew, le vrai tueur de performance

Un join sur user_idun 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 :

  1. Onglet Stages → la médiane des tâches est à 10 s, le max à 45 min.
  2. La durée totale du stage est presque égale à la durée de la tâche max.
  3. L’onglet SQL montre un 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 :

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

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

Les deux gestes qui règlent 90 % des cas

Si vous devez retenir deux choses de tout cet article :

  1. Broadcast quand une table est petite — hint explicite, seuil relevé si nécessaire.
  2. AQE activé — le skew join et le coalesce automatique règlent la majorité du reste sans intervention.

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.

Ce qu’il faut retenir

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.

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