Улучшение производительности

Очистка данных с помощью 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

Что такое перемешивание?

Перемешивание — это перемещение данных между узлами для выполнения задачи

  • Скрывает сложность от пользователя
  • Может выполняться медленно
  • Снижает общую пропускную способность
  • Часто необходимо, но лучше минимизировать
Очистка данных с помощью PySpark

Как минимизировать перемешивание?

  • Ограничьте использование .repartition(num_partitions)
    • Вместо этого применяйте .coalesce(num_partitions)
  • Используйте .join() с осторожностью
  • Применяйте .broadcast()
  • В ряде случаев ограничение не требуется
Очистка данных с помощью PySpark

Широковещательная передача

Широковещательная передача:

  • Передаёт копию объекта каждому узлу
  • Сокращает лишний обмен данными между узлами
  • Может значительно ускорить операции .join()

Используйте метод .broadcast(<DataFrame>)

from pyspark.sql.functions import broadcast
combined_df = df_1.join(broadcast(df_2))
Очистка данных с помощью PySpark

Давайте потренируемся!

Очистка данных с помощью PySpark

Preparing Video For Download...