CommencezCommencez gratuitement

Utiliser la diffusion (broadcast) dans les jointures Spark

Rappelez-vous que les jointures de tables dans Spark sont réparties entre les travailleurs du grappe. Si les données ne sont pas locales, diverses opérations de redistribution (shuffle) sont nécessaires et peuvent nuire au rendement. À la place, nous allons utiliser l'opération broadcast de Spark pour donner à chaque nœud une copie des données indiquées.

Quelques conseils :

  • Diffusez le plus petit DataFrame. Plus le DataFrame est gros, plus le transfert vers les nœuds travailleurs prendra du temps.
  • Pour les petits DataFrames, il peut être préférable d'éviter la diffusion et de laisser Spark déterminer lui-même l'optimisation.
  • Si vous examinez le plan d'exécution de la requête, un broadcastHashJoin indique que vous avez bel et bien configuré la diffusion.

Les DataFrames flights_df et airports_df sont à votre disposition.

Cette activité fait partie du cours

Nettoyer des données avec PySpark

Voir le cours

Instructions de l’exercice

  • Importez la méthode broadcast() depuis pyspark.sql.functions.
  • Créez un nouveau DataFrame broadcast_df en joignant flights_df avec airports_df, en utilisant la diffusion (broadcast).
  • Affichez le plan de requête et comparez-le à l'original.

Exercice interactif pratique

Essayez cet exercice en complétant ce code d’exemple.

# Import the broadcast method from pyspark.sql.functions
from ____ import ____

# Join the flights_df and airports_df DataFrames using broadcasting
broadcast_df = flights_df.____(____(airports_df), \
    flights_df["Destination Airport"] == airports_df["IATA"] )

# Show the query plan and compare against the original
broadcast_df.____()
Modifier et exécuter le code