Začněte nyníZačněte zdarma

Použití broadcastingu při joinech ve Sparku

Joiny tabulek ve Sparku se rozdělují mezi jednotlivé workery v clusteru. Pokud data nejsou dostupná lokálně, jsou potřeba různé shuffle operace, které mohou negativně ovlivnit výkon. Místo toho využijeme Sparkovu operaci broadcast, která každému uzlu poskytne vlastní kopii zadaných dat.

Pár tipů:

  • Broadcastuj menší DataFrame. Čím větší DataFrame, tím více času zabere jeho přenos na workery.
  • U malých DataFramů může být lepší broadcasting vynechat a nechat Spark, ať si optimalizaci vyřeší sám.
  • Pokud se podíváš na plán provádění dotazu, broadcastHashJoin znamená, že se broadcasting povedlo nastavit správně.

K dispozici máš DataFramy flights_df a airports_df.

Toto cvičení je součástí kurzu

Cleaning Data with PySpark

Zobrazit kurz

Pokyny k cvičení

  • Importuj metodu broadcast() z pyspark.sql.functions.
  • Vytvoř nový DataFrame broadcast_df tak, že spojíš flights_df s airports_df pomocí broadcastingu.
  • Zobraz plán dotazu a zamysli se nad tím, čím se liší od původního.

Interaktivní cvičení na vyzkoušení si v praxi

Vyzkoušejte si toto cvičení dokončením tohoto ukázkového kódu.

# 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.____()
Upravit a spustit kód