始める無料で始める

Spark の結合でブロードキャストを使う

Spark ではテーブルの結合はクラスター内のワーカーに分散されます。データがローカルにない場合はさまざまなシャッフル処理が発生し、パフォーマンスに悪影響が出ることがあります。そこで、Spark の broadcast を使って、指定したデータのコピーを各ノードに配布します。

いくつかのヒント:

  • 小さい DataFrame をブロードキャストしましょう。DataFrame が大きいほど、ワーカーノードへの転送に時間がかかります。
  • 小さな DataFrame では、ブロードキャストを行わずに Spark に最適化を任せたほうがよい場合があります。
  • クエリ実行計画を見ると、broadcastHashJoin が表示されれば、ブロードキャストの設定に成功しています。

DataFrame の flights_dfairports_df が用意されています。

この演習はコースの一部です

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.____()
コードを編集して実行