Cải thiện hiệu năng

Làm sạch dữ liệu với PySpark

Mike Metzger

Data Engineering Consultant

Giải thích kế hoạch thực thi của 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>
Làm sạch dữ liệu với PySpark

Shuffling là gì?

Shuffling là việc di chuyển dữ liệu giữa các worker để hoàn thành tác vụ

  • Ẩn bớt độ phức tạp với người dùng
  • Có thể chậm
  • Giảm thông lượng tổng thể
  • Thường cần thiết, nhưng nên giảm thiểu
Làm sạch dữ liệu với PySpark

Cách giảm shuffling

  • Hạn chế dùng .repartition(num_partitions)
    • Thay bằng .coalesce(num_partitions)
  • Cân nhắc khi gọi .join()
  • Dùng .broadcast()
  • Có thể không cần giới hạn
Làm sạch dữ liệu với PySpark

Broadcasting

Broadcasting:

  • Cung cấp một bản sao đối tượng cho mỗi worker
  • Tránh giao tiếp thừa giữa các nút
  • Có thể tăng tốc đáng kể các phép .join()

Dùng .broadcast(<DataFrame>)

from pyspark.sql.functions import broadcast
combined_df = df_1.join(broadcast(df_2))
Làm sạch dữ liệu với PySpark

Cùng luyện tập!

Làm sạch dữ liệu với PySpark

Preparing Video For Download...