Zacznij terazZacznij za darmo

Używanie broadcastingu w złączeniach Sparka

Pamiętaj, że złączenia tabel w Sparku są rozdzielane między węzły robocze klastra. Jeśli dane nie są dostępne lokalnie, konieczne są różne operacje shuffle, które mogą negatywnie wpłynąć na wydajność. Zamiast tego użyjemy operacji broadcast Sparka, aby przekazać każdemu węzłowi kopię wskazanych danych.

Kilka wskazówek:

  • Wykonuj broadcast na mniejszym DataFrame. Im większy DataFrame, tym więcej czasu zajmuje jego przesłanie do węzłów roboczych.
  • W przypadku małych DataFrame'ów lepiej może być pominąć broadcasting i pozwolić Sparkowi samodzielnie dobrać optymalizację.
  • Jeśli przejrzysz plan wykonania zapytania, broadcastHashJoin oznacza, że broadcasting został poprawnie skonfigurowany.

DataFrame'y flights_df i airports_df są dostępne w środowisku.

To ćwiczenie jest częścią kursu

Czyszczenie danych w PySpark

Zobacz kurs

Instrukcje do ćwiczenia

  • Zaimportuj metodę broadcast() z pyspark.sql.functions.
  • Utwórz nowy DataFrame broadcast_df, łącząc flights_df z airports_df przy użyciu broadcastingu.
  • Wyświetl plan zapytania i zastanów się, czym różni się od poprzedniego.

Interaktywne ćwiczenie praktyczne

Spróbuj tego ćwiczenia, uzupełniając ten przykładowy kod.

# 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.____()
Edytuj i uruchom kod