在 Spark 連接中使用廣播
記住,在 Spark 中的資料表連接會分散到叢集的多個 worker 上執行。若資料不在本地,就需要進行各種 shuffle 操作,可能會對效能造成負面影響。為了改善這點,我們要使用 Spark 的 broadcast 操作,讓每個節點都擁有一份指定資料的複本。
幾個小提醒:
- 請廣播較小的 DataFrame。DataFrame 越大,傳送到各個 worker 節點所需的時間就越長。
- 對於很小的 DataFrame,可能更適合不要廣播,讓 Spark 自行找出最佳化方式。
- 如果你查看查詢的執行計畫,出現 broadcastHashJoin 就代表你已成功設定廣播。
flights_df 和 airports_df 這兩個 DataFrame 已可供你使用。
本練習屬於課程
使用 PySpark 清理資料
練習說明
- 從
pyspark.sql.functions匯入broadcast()方法。 - 透過廣播,將
flights_df與airports_df連接,建立新的 DataFramebroadcast_df。 - 顯示查詢計畫,並思考與原本版本的差異。
動手互動練習
試著完成這個範例程式碼,體驗一下這個練習。
# 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.____()