С большой
Люблю мемы про дата саентистов, которые ставят свой sql-запрос и уходят пить кофе на три года. Смешно, но не жизненно. Мы, так называемые дата инженеры, слушаем и не осуждаем, а помогаем!
Задача: сделать инференс модели по батчам на Spark.
Загоняем препроцессинг в mapInPandas, итерируем по списку файлов в parquet и дело с концом. Среднее время обработки батча - 10 минут… А батчей десятки тысяч… Никуда не годится! Расскажу вам сегодня о нескольких способах ускорения вычислений на Spark.
Во-первых, выставляем параметр cores на 2 или больше (сколько совесть позволяет).
SparkConf()
.setAppName(“name”)
.setMaster(“yarn”)
.set(“spark.executor.cores”, 2)
Если другие юзеры кластера не пришли к вам с вилами, то идем дальше.
Экзекьюторов тоже побольше.
.set(“spark.dynamicAllocation.maxExecutors”, 50)
Помним, что общее количество ядер = cores*executors, с ресурсами тоже не наглейте…
Во-вторых, те таблицы, что мы используем несколько раз: кэшируем.
Cache обязательно используем с action-функцией, иначе не выполнится.
DF.cache().count()
В-третьих, перед сохранением делаем repartiton датафрейма. Функция repartition в Apache Spark перераспределяет данные между заданным числом разбиений (partitions), обеспечивая равномерность распределения и улучшая производительность вычислений. Хорошо работает, но только для больших таблиц.
Я задала константное количество партиций - 15. Оцените размеры таблицы самостоятельно и поперебирайте этот параметр.
NEW_DF = DF.repartiton(15)
Этими тремя нехитрыми способами удалось снизить время обработки батча с 10 минут до 30 секунд, но этого недостаточно, поэтому идем дальше.
В-четвертых, заброадкастим наш датасет на все узлы кластера. Функция broadcast в Spark используется для распределения неизменяемых объектов на все узлы кластера, чтобы избежать их повторной отправки при каждой задаче. Это особенно полезно для больших наборов данных, которые нужно использовать в вычислениях, так как позволяет экономить сетевой трафик и повышать производительность. Объекты, передаваемые через broadcast, хранятся в памяти на узлах, что позволяет быстро к ним обращаться. Ключевое различие между cache и broadcast (это любят спрашивать на собеседованиях): cache хранит промежуточные данные в памяти для ускорения повторных вычислений, а broadcast эффективно распространяет неизменяемые объекты на все узлы для уменьшения сетевых затрат при выполнении задач.
from pyspark.sql import functions as F
NEW_DF = F.broadcast(DF)
В-пятых, чтение таблицы с явным указанием схемы.
DF = spark.read.schema(input_schema).parquet(path)
И последний пункт: в-шестых, использовать reduce для агрегации в финальный датасет. Функция reduce в Spark выполняет агрегацию данных, объединяя все элементы в одной RDD (Resilient Distributed Dataset) или DataFrame с помощью заданной функции. Она применяет указанную бинарную операцию к элементам, последовательно сводя их к одному значению. Например, можно использовать reduce для суммирования чисел или нахождения наибольшего значения в наборе данных.
Функция reduce ускоряет обработку данных благодаря:
1. Снижению объема данных, сводя их к одному значению.
2. Параллелизму, позволяя выполнять операции на разных узлах кластера.
3. Уменьшению потребления памяти, сводя набор данных к меньшему числу значений.
4. Оптимизации вычислений, улучшая выполнение DAG в Spark.
Пример:
from functools import reduce
from pyspark.sql import DataFrame
results = []
for part_path in partition_paths: results.append(process_part(part_path))
if results:
final_df = reduce(DataFrame.union, results)
Уже использовав все описанные способы, мне удалось ускорить время обработки батча до 0.3 секунды.
#data_science #spark #hadoop #hdfs
