ÎncepețiÎncepe gratuit

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

Vezi cursul

Instrucțiuni pentru exercițiu

  • Importă metoda broadcast() din pyspark.sql.functions.
  • Creează un nou DataFrame numit broadcast_df prin join-ul dintre flights_df și airports_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.____()
Editează și rulează codul