Qu’est-ce qu’un broadcast join et quand l’utiliser ?

Questions d’entrevue Apache Spark

Intermédiairejoinbroadcastshuffleperformance

La réponse courte

Un broadcast join copie intégralement la petite table sur chaque exécuteur, ce qui permet de joindre localement, partition par partition, sans aucun shuffle de la grande table.

C’est la stratégie la plus rapide quand elle s’applique, parce qu’elle supprime l’opération la plus coûteuse du job. Spark la choisit automatiquement quand il estime la petite table sous le seuil de spark.sql.autoBroadcastJoinThreshold, qui vaut 10 Mo par défaut.

On peut la forcer :

scala
import org.apache.spark.sql.functions.broadcast

grandeTable.join(broadcast(petiteTable), Seq("cle"))

Ce que l’intervieweur vérifie

Que vous connaissez la limite, parce que forcer un broadcast à l’aveugle est une façon classique de faire tomber un cluster.

La table diffusée doit tenir en mémoire sur chaque exécuteur, et elle passe d’abord par le driver, qui la collecte avant de la redistribuer. Diffuser une table de 2 Go sur 50 exécuteurs, c’est 2 Go dans le driver puis 100 Go de mémoire consommée au total. Le symptôme typique est un OutOfMemoryError sur le driver, pas sur les exécuteurs.

L’ordre de grandeur raisonnable : jusqu’à quelques centaines de mégaoctets en montant le seuil, au-delà il faut une autre stratégie.

Les autres stratégies, pour situer

C’est souvent la relance, donc autant l’anticiper :

  • Sort-Merge Join — le choix par défaut sur deux grandes tables. Les deux côtés sont shufflés par clé puis triés. Robuste, mais coûteux.
  • Shuffle Hash Join — shuffle des deux côtés, puis table de hachage en mémoire sur le plus petit. Utilisé quand un côté est nettement plus petit sans être diffusable.
  • Broadcast Nested Loop Join — le repli quand il n’y a pas de condition d’égalité (une jointure sur un intervalle, par exemple). Complexité quadratique : si votre join met une heure sans raison apparente, c’est souvent celui-là qu’il faut chercher dans le plan.

La relance probable

« Pourquoi Spark n’a-t-il pas choisi le broadcast alors que ma table est petite ? »

Parce qu’il se fie aux statistiques, pas à la taille réelle. Sur une table sans statistiques à jour — un CSV, un DataFrame issu de transformations complexes — son estimation peut être largement fausse. D’où l’utilité de ANALYZE TABLE ... COMPUTE STATISTICS, ou du broadcast() explicite quand on sait ce qu’on fait.

Et depuis Spark 3, l’Adaptive Query Execution corrige une partie de ces erreurs en cours d’exécution, en convertissant un sort-merge en broadcast quand la taille réelle mesurée à l’exécution le permet.

Le détail des stratégies est dans joins en Spark : broadcast, shuffle et skew.

Toutes les questions Apache Spark