開始使用免費開始

在 Spark 連接中使用廣播

記住,在 Spark 中的資料表連接會分散到叢集的多個 worker 上執行。若資料不在本地,就需要進行各種 shuffle 操作,可能會對效能造成負面影響。為了改善這點,我們要使用 Spark 的 broadcast 操作,讓每個節點都擁有一份指定資料的複本。

幾個小提醒:

  • 請廣播較小的 DataFrame。DataFrame 越大,傳送到各個 worker 節點所需的時間就越長。
  • 對於很小的 DataFrame,可能更適合不要廣播,讓 Spark 自行找出最佳化方式。
  • 如果你查看查詢的執行計畫,出現 broadcastHashJoin 就代表你已成功設定廣播。

flights_dfairports_df 這兩個 DataFrame 已可供你使用。

本練習屬於課程

使用 PySpark 清理資料

檢視課程

練習說明

  • pyspark.sql.functions 匯入 broadcast() 方法。
  • 透過廣播,將 flights_dfairports_df 連接,建立新的 DataFrame broadcast_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.____()
編輯並執行程式碼