Använda broadcasting vid Spark-joins
Kom ihåg att tabelljoins i Spark fördelas mellan klustrets workers. Om data inte finns lokalt krävs olika shuffle-operationer, vilket kan påverka prestandan negativt. Istället ska vi använda Sparks broadcast-funktion för att ge varje nod en kopia av den angivna datan.
Några tips:
- Broadcasta den mindre DataFrame:n. Ju större DataFrame, desto längre tid tar det att överföra den till worker-noderna.
- För små DataFrames kan det vara bättre att hoppa över broadcasting och låta Spark sköta optimeringen på egen hand.
- Om du tittar på frågekörningsplanen indikerar en broadcastHashJoin att du har konfigurerat broadcasting korrekt.
DataFrames:en flights_df och airports_df finns tillgängliga för dig.
Den här övningen är en del av kursen
Datarensning med PySpark
Övningsinstruktioner
- Importera metoden
broadcast()frånpyspark.sql.functions. - Skapa en ny DataFrame
broadcast_dfgenom att joinaflights_dfmedairports_dfmed hjälp av broadcasting. - Visa frågekörningsplanen och fundera på skillnaderna jämfört med originalet.
Interaktiv övning med praktiskt arbete
Testa den här övningen genom att slutföra den här exempelkoden.
# 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.____()