Spark の結合でブロードキャストを使う
Spark ではテーブルの結合はクラスター内のワーカーに分散されます。データがローカルにない場合はさまざまなシャッフル処理が発生し、パフォーマンスに悪影響が出ることがあります。そこで、Spark の broadcast を使って、指定したデータのコピーを各ノードに配布します。
いくつかのヒント:
- 小さい DataFrame をブロードキャストしましょう。DataFrame が大きいほど、ワーカーノードへの転送に時間がかかります。
- 小さな DataFrame では、ブロードキャストを行わずに Spark に最適化を任せたほうがよい場合があります。
- クエリ実行計画を見ると、broadcastHashJoin が表示されれば、ブロードキャストの設定に成功しています。
DataFrame の flights_df と airports_df が用意されています。
この演習はコースの一部です
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.____()