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,
broadcastHashJoinznamená, ž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
Pokyny k cvičení
- Importuj metodu
broadcast()zpyspark.sql.functions. - Vytvoř nový DataFrame
broadcast_dftak, že spojíšflights_dfsairports_dfpomocí 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.____()