效能提升

使用 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 清理資料

什麼是 shuffling?

Shuffling 指為完成任務而在多個 worker 之間搬移資料。

  • 對使用者隱藏複雜度
  • 可能較慢
  • 會降低整體吞吐量
  • 常見且必要,但要盡量減少
使用 PySpark 清理資料

如何減少 shuffling?

  • 少用 .repartition(num_partitions)
    • 改用 .coalesce(num_partitions)
  • 呼叫 .join() 時要小心
  • 使用 .broadcast()
  • 有時不必特別限制
使用 PySpark 清理資料

Broadcasting

Broadcasting

  • 為每個 worker 提供物件副本
  • 避免節點間不必要的通訊
  • 可大幅加速 .join()

使用 .broadcast(<DataFrame>) 方法

from pyspark.sql.functions import broadcast
combined_df = df_1.join(broadcast(df_2))
使用 PySpark 清理資料

一起來練習吧!

使用 PySpark 清理資料

Preparing Video For Download...