การปรับปรุงประสิทธิภาพ

การทำความสะอาดข้อมูลด้วย PySpark

Mike Metzger

Data Engineering Consultant

อธิบาย Spark execution plan

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 ต่าง ๆ เพื่อทำงานให้เสร็จสมบูรณ์

  • ซ่อนความซับซ้อนจากผู้ใช้
  • อาจใช้เวลานาน
  • ลด throughput โดยรวม
  • มักจำเป็น แต่ควรลดให้น้อยที่สุด
การทำความสะอาดข้อมูลด้วย PySpark

จะลด shuffling ได้อย่างไร?

  • จำกัดการใช้ .repartition(num_partitions)
    • ใช้ .coalesce(num_partitions) แทน
  • ระวังการเรียกใช้ .join()
  • ใช้ .broadcast()
  • อาจไม่จำเป็นต้องจำกัดเสมอไป
การทำความสะอาดข้อมูลด้วย PySpark

Broadcasting

Broadcasting:

  • ส่งสำเนาของออบเจกต์ไปยัง worker แต่ละตัว
  • ลดการสื่อสารระหว่าง node ที่ไม่จำเป็น
  • เพิ่มความเร็วของ .join() ได้อย่างมาก

ใช้เมธอด .broadcast(<DataFrame>)

from pyspark.sql.functions import broadcast
combined_df = df_1.join(broadcast(df_2))
การทำความสะอาดข้อมูลด้วย PySpark

มาฝึกกันเถอะ!

การทำความสะอาดข้อมูลด้วย PySpark

Preparing Video For Download...