Usare il broadcasting nei join di Spark
Ricorda che i join tra tabelle in Spark sono suddivisi tra i worker del cluster. Se i dati non sono locali, sono necessarie varie operazioni di shuffle che possono avere un impatto negativo sulle prestazioni. Invece, useremo le operazioni di broadcast di Spark per dare a ogni nodo una copia dei dati specificati.
Qualche suggerimento:
- Fai il broadcast del DataFrame più piccolo. Più è grande il DataFrame, più tempo serve per trasferirlo ai nodi worker.
- Su DataFrame di piccole dimensioni, potrebbe essere meglio evitare il broadcasting e lasciare che Spark individui da solo eventuali ottimizzazioni.
- Se guardi il piano di esecuzione della query, un broadcastHashJoin indica che hai configurato correttamente il broadcasting.
I DataFrame flights_df e airports_df sono a tua disposizione.
Questo esercizio fa parte del corso
Pulizia dei dati con PySpark
Istruzioni dell'esercizio
- Importa il metodo
broadcast()dapyspark.sql.functions. - Crea un nuovo DataFrame
broadcast_dfeffettuando il join traflights_dfeairports_df, usando il broadcasting. - Mostra il piano della query e valuta le differenze rispetto all'originale.
esercizio interattivo pratico
Prova questo esercizio completando questo codice di esempio.
# 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.____()