🛹5 шагов от Python к PySpark и 10 лучших практик настройки Spark-заданий
Узнайте, как быстро конвертировать Python-скрипты в задания PySpark, эффективно используя всю мощь распределенных вычислений Apache Spark.
1. Преобразуйте локальный датафрейм Pandas в Spark Dataframe через Apache Arrow (независимый от языка столбчатый формат в памяти) или Koalas (API Pandas в Apache Spark)
2. Напишите пользовательскую функцию PySpark (UDF) для функции Python. UDF PySpark принимают столбцы и применяет логику построчно для создания нового столбца
3. Загрузите датасет в Spark RDD или DataFrame
4. Избегайте циклов, используя преобразование map() для каждого элемента RDD с использованием функции, возвращающей новый RDD.
5. Учитывайте взаимозависимость датафреймов – если новое значение столбца DataFrame зависит от других таких же структур данных, объедините их через JOIN и вызовите UDF, чтобы получить новое значение столбца.
Чтобы по максимуму использовать все возможности кластера, перед запуском Spark-заданий помните о следующих рекомендациях:
1. Избегайте слишком больших структур данных (RDD, DataFrames) и помните про форматы (Avro и Parquet лучше, чем TXT, CSV или JSON)
2. Для уменьшения накладных расходов на параллельную обработку данных используйте coalesce(), чтобы сократить количество разделов
3. Сокращайте неиспользуемые ресурсы (ядра в кластере), распределяя данные с помощью repartition()
4. Используйте reduceByKey вместо groupByKey, настраивая уровень параллелизма и задавая количество разделов при вызове операций перетасовки данных (shuffle)
5. Избегайте перетасовки больших объемов данных, настроив spark.sql.shuffle.partitions для указания количества разделов при перетасовке для объединений или агрегатов.
6. Отфильтруйте данные перед обработкой, убрав лишнее
7. Используйте Broadcast-переменные, подобные распределенному кэшу в Hadoop, чтобы повысить производительность, сделав данные доступными для всех исполнителей и уменьшив их перетасовку
8. Если RDD или DataFrame используется более одного раза, кэшируйте их, чтобы избежать повторного вычисления и повысить производительность
9. Следите за пользовательским интерфейсом Spark для настройки своего приложения
10. Используйте динамическое размещение (spark.dynamicAllocation.enabled), чтобы масштабировать количество исполнителей в приложении в зависимости от рабочей нагрузки
Post #190
1.08K