Utilizarea broadcasting în join-urile Spark
Reține că join-urile de tabele în Spark sunt distribuite între nodurile worker ale clusterului. Dacă datele nu sunt locale, sunt necesare diverse operațiuni de tip shuffle, care pot afecta negativ performanța. În schimb, vom folosi operațiunea broadcast din Spark pentru a oferi fiecărui nod o copie a datelor specificate.
Câteva sfaturi utile:
- Aplică broadcast pe DataFrame-ul mai mic. Cu cât un DataFrame este mai mare, cu atât durează mai mult transferul către nodurile worker.
- Pentru DataFrame-uri mici, poate fi mai bine să sari peste broadcasting și să lași Spark să gestioneze optimizările singur.
- Dacă analizezi planul de execuție al interogării, prezența unui broadcastHashJoin confirmă că broadcasting-ul a fost configurat cu succes.
DataFrame-urile flights_df și airports_df sunt disponibile.
Acest exercițiu face parte din cursul
Curățarea datelor cu PySpark
Instrucțiuni pentru exercițiu
- Importă metoda
broadcast()dinpyspark.sql.functions. - Creează un nou DataFrame numit
broadcast_dfprin join-ul dintreflights_dfșiairports_df, folosind broadcasting. - Afișează planul de interogare și observă diferențele față de varianta originală.
Exercițiu interactiv practic
Încearcă acest exercițiu completând acest cod de exemplu.
# 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.____()