性能优化

使用 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)指为完成任务在各个工作节点间移动数据

  • 对用户屏蔽复杂性
  • 可能较慢
  • 降低整体吞吐量
  • 常常必要,但应尽量减少
使用 PySpark 进行数据清洗

如何减少洗牌?

  • 少用 .repartition(num_partitions)
    • 优先用 .coalesce(num_partitions)
  • 谨慎调用 .join()
  • 使用 .broadcast()
  • 不一定需要限制
使用 PySpark 进行数据清洗

广播

广播(Broadcasting):

  • 为每个工作节点提供对象副本
  • 避免节点间不必要通信
  • 可大幅加速 .join()

使用 .broadcast(<DataFrame>) 方法

from pyspark.sql.functions import broadcast
combined_df = df_1.join(broadcast(df_2))
使用 PySpark 进行数据清洗

让我们来练习!

使用 PySpark 进行数据清洗

Preparing Video For Download...