Поліпшення продуктивності

Очищення даних у 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 — це переміщення даних між виконавцями для завершення завдання

  • Приховує складність від користувача
  • Може бути повільним
  • Знижує загальну пропускну здатність
  • Часто необхідний, але його слід мінімізувати
Очищення даних у PySpark

Як зменшити shuffling?

  • Обмежуйте використання .repartition(num_partitions)
    • Натомість використовуйте .coalesce(num_partitions)
  • Будьте обережні з викликами .join()
  • Використовуйте .broadcast()
  • Можливо, обмежувати не потрібно
Очищення даних у PySpark

Broadcasting

Broadcasting:

  • Надає копію об'єкта кожному виконавцю
  • Запобігає зайвій комунікації між вузлами
  • Може суттєво прискорити операції .join()

Використовуйте метод .broadcast(<DataFrame>)

from pyspark.sql.functions import broadcast
combined_df = df_1.join(broadcast(df_2))
Очищення даних у PySpark

Давайте потренуємось!

Очищення даних у PySpark

Preparing Video For Download...