パフォーマンスの改善

PySpark でデータをクレンジングする

Mike Metzger

Data Engineering Consultant

Spark の実行計画を読む

voter_df = df.select(df['VOTER NAME']).distinct()
voter_df.explain()
== Physical Plan ==
*(2) HashAggregate(keys=[VOTER NAME#15], functions=[])
+- Exchange hashpartitioning(VOTER NAME#15, 200)
   +- *(1) HashAggregate(keys=[VOTER NAME#15], functions=[])
      +- *(1) FileScan csv [VOTER NAME#15] Batched: false, Format: CSV, Location: 
      InMemoryFileIndex[file:/DallasCouncilVotes.csv.gz], 
      PartitionFilters: [], PushedFilters: [], 
      ReadSchema: struct<VOTER NAME:string>
PySpark でデータをクレンジングする

シャッフルとは?

「シャッフル」とは、タスク完了のためにデータを複数ワーカー間で移動することです。

  • ユーザーから複雑さを隠す
  • 完了に時間がかかることがある
  • スループットを低下させる
  • しばしば必要だが、最小化を心がける
PySpark でデータをクレンジングする

シャッフルを抑える方法

  • .repartition(num_partitions)の多用は避ける
    • 代わりに.coalesce(num_partitions)を使用
  • .join()の呼び出しに注意
  • .broadcast()を使用
  • 制限不要な場合もある
PySpark でデータをクレンジングする

ブロードキャスト

ブロードキャスト:

  • 各ワーカーにオブジェクトのコピーを配布
  • ノード間の不要な通信を防ぐ
  • .join()を大幅に高速化できる

.broadcast(<DataFrame>) を使用

from pyspark.sql.functions import broadcast
combined_df = df_1.join(broadcast(df_2))
PySpark でデータをクレンジングする

練習しましょう!

PySpark でデータをクレンジングする

Preparing Video For Download...