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 :
import org.apache.spark.sql.functions.broadcast
grandeTable.join(broadcast(petiteTable), Seq("cle"))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.
C’est souvent la relance, donc autant l’anticiper :
« 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.