сеньорчикОткрыть в Telegram
← вся теориятеория к собесу · Spark и большие данные

Оптимизация джобов Spark

Оптимизация Spark: меньше тасовать

Оптимизация Spark это в первую очередь сокращение шаффла и борьба с перекосом, а не подбор магических конфигов. Собес проверяет, идёшь ли ты от Spark UI и понимаешь ли, что такое skew и broadcast.

Стержень: лучший шаффл - несостоявшийся; мелкую сторону джойна вещают broadcast'ом, а стейдж, стоящий на одной таске, это skew.

// Формулировки: «джойн тормозит - что сделаешь?», «что такое перекос и как лечить?», «broadcast join - когда?».

Меньше данных и broadcast

Первым делом - меньше данных: фильтры и выбор колонок как можно раньше, партиционные фильтры при чтении. Лучший шаффл - тот, что не состоялся, потому что данных до него дошло меньше.

Broadcast join - главная оптимизация джойнов: если одна таблица мелкая, её рассылают на все ноды целиком, и шаффл гигантской стороны не нужен. Важно, чтобы «мелкая» реально влезала в память экзекьютора.

// Классическая ошибка - джойнить справочник в 10 МБ обычным способом, тасуя обе стороны. Один broadcast снимает проблему.

from pyspark.sql.functions import broadcast
# малую сторону — на все ноды, без шаффла большой
big.join(broadcast(small), "id")
broadcast hint
принудительная рассылка малой стороны джойна на все ноды
predicate pushdown
фильтр уходит в чтение - меньше данных в конвейере

Перекос (skew)

Skew - перекос по ключу: один ключ-гигант (NULL, значение по умолчанию, топ-клиент) собирает непропорционально много строк, и одна таска молотит часами, пока остальные давно закончили. В UI это стейдж, застрявший на последней таске.

Лечат тремя способами: AQE (adaptive query execution) skew join (автоматически дробит перекошенную сторону), salting (суффикс к горячему ключу размазывает его по партициям), или отделение горячих ключей в отдельную ветку обработки.

// Поэтому «Spark тупит на последнем проценте» почти всегда значит не «тупит», а перекос - надо смотреть распределение ключей, а не крутить память.

skew
перекос: один ключ собирает непропорционально много данных
salting
суффикс к горячему ключу - размазывание перекоса по партициям

Спиллы, UDF, файлы, ресурсы

Спиллы лечат больше памятью на таск (меньше параллельных тасков на экзекьютор), большим числом партиций (мельче куски) или сужением строк до шаффла. repartition(ключ) перед серией операций по ключу даёт один управляемый шаффл вместо трёх; coalesce сжимает число партиций без шаффла перед записью.

UDF (user-defined function) - тормоз и чёрный ящик для Catalyst. Порядок предпочтения: встроенные функции → pandas_udf (векторизация через Arrow) → RDD-код (resilient distributed dataset). Python-UDF ради upper() - в разы медленнее встроенной функции.

// На выходе следят за мелкими файлами: тысяча партиций даёт тысячу файлов, следующий читатель страдает. Но coalesce(1) на терабайте - другая крайность: вся запись в одну таску. Экзекьютор - 4–5 ядер, больше даёт GC-паузы (garbage collector) и контеншн.

repartition / coalesce
переразбиение с шаффлом / склейка без шаффла
pandas_udf
векторизованный UDF через Arrow, батчами

Как отвечать: «Джойн в Spark тормозит - как оптимизируешь?»

Иду от Spark UI, а не от догадок. Сначала смотрю, не перекос ли это: если стейдж висит на одной-двух тасках, пока остальные закончили, это skew по ключу джойна, и я лечу его через AQE skew join, salting горячего ключа или отдельную ветку для него. Если перекоса нет, смотрю на размеры сторон: когда одна таблица мелкая, например справочник, я делаю broadcast - рассылаю её на все ноды и убираю шаффл большой стороны совсем, это самая частая победа. Дальше сокращаю данные до джойна: фильтры и выбор колонок как можно раньше, партиционные фильтры при чтении - лучший шаффл тот, что не случился. Если видны спиллы, добавляю памяти на таск или дроблю на больше партиций. И проверяю, что джойн идёт по ключу дистрибуции, а не гоняет всё по сети. То есть последовательность: перекос → broadcast → меньше данных → память.

Почему это сильный ответ: диагностика от UI, правильный порядок (skew → broadcast → меньше данных → память) и конкретные приёмы для каждого - инженерный разбор, а не список опций.

На чём валят

  • Джойн со справочником в 10 МБ шаффлом обеих сторон - забыт broadcast.
  • Стейдж стоит на последней таске - skew, а не «Spark тупит»; смотри распределение ключей.
  • coalesce(1) на терабайте «чтобы один файл» - вся запись в одну таску.
  • Python-UDF для upper() - встроенная функция в разы быстрее.
  • Гигантские экзекьюторы на 32 ядра - GC-паузы и контеншн вместо скорости.

Проверьте себя

Пять вопросов из банка по этой подтеме. Всего их 15, остальные разбираются в тренажёре.

  1. #spark_optimization1 / 5
    Джоба «висит»: 199 задач завершились, одна работает часами. Что это и в чём причина?
    A)Skew означает, что у Spark просто закончилась оперативная память на драйвере
    B)Ключ неравномерен: одна партиция огромна, её задача тормозит всё
    C)Skew — это когда все партиции строго одинакового размера по строкам
    D)Перекос данных лечится покупкой более мощных узлов кластера
    показать ответ и разбор
    +B)Ключ неравномерен: одна партиция огромна, её задача тормозит всё

    // разбор: Классический data skew: значения ключа распределены неравномерно (99% строк с одним ключом), и после shuffle почти все данные попадают в одну партицию. Её задача-долгожитель (straggler) держит всю стадию. Лечат «солением» ключа (salting), отдельной обработкой горячих ключей, broadcast-join'ом или repartition.

  2. #spark_optimization2 / 5
    Что такое broadcast join и когда он уместен?
    A)Broadcast-join быстрее других типов join независимо от размеров таблиц
    B)Broadcast рассылает большую таблицу по узлам, чтобы обязательно вызвать shuffle
    C)Маленькую таблицу копируют на все узлы — join без shuffle большой
    D)Broadcast применим к таблицам одинакового большого размера
    показать ответ и разбор
    +C)Маленькую таблицу копируют на все узлы — join без shuffle большой

    // разбор: Если одна таблица мала (справочник), Spark рассылает её копию на все узлы, и большая таблица джойнится локально, без дорогого shuffle. Это резко ускоряет join. Но если бродкастить слишком большую таблицу, она не влезет в память исполнителей и driver'а — джоба упадёт по OOM. Порог настраивается (autoBroadcastJoinThreshold).

  3. #spark_optimization3 / 5
    Зачем в Spark вызывают cache()/persist() на DataFrame?
    A)Чтобы сохранить результат джобы в постоянное дисковое хранилище кластера после завершения
    B)Чтобы переиспользуемый датафрейм не пересчитывался с нуля при каждом action
    C)Чтобы автоматически удалить дубликаты строк из датафрейма
    D)Чтобы принудительно уменьшить число партиций ровно до одной
    показать ответ и разбор
    +B)Чтобы переиспользуемый датафрейм не пересчитывался с нуля при каждом action

    // разбор: План Spark ленивый: каждое действие (action) пересчитывает всю цепочку с нуля. Если один датафрейм используется несколькими действиями (итеративный алгоритм, ветвление), cache()/persist() держит его посчитанным в памяти/на диске, экономя пересчёты. Кэшировать имеет смысл только реально переиспользуемый результат — иначе это лишняя память.

  4. #spark_optimization4 / 5
    Executor падает по OutOfMemory на shuffle-стадии. Что происходит и что делать?
    A)Ошибка OOM на shuffle означает, что на узле кластера кончилось свободное место на диске
    B)Партиция не влезла в память executor'а; помогают больше партиций, меньше skew
    C)Executor OOM лечится уменьшением объёма входных данных вдвое
    D)OOM на shuffle никак не связан с распределением данных по ключу
    показать ответ и разбор
    +B)Партиция не влезла в память executor'а; помогают больше партиций, меньше skew

    // разбор: OOM на shuffle обычно значит, что после перераспределения в один executor попал слишком большой объём — чаще всего из-за data skew (перекос ключа) или слишком крупных партиций. Помогают: увеличить число партиций (repartition), «посолить» перекошенный ключ, включить spill на диск, поднять память executor'а. Урезание данных — крайняя мера, а не первый ход.

  5. #spark_optimization5 / 5
    Что даёт запись данных с partitionBy(колонка) в Spark/Parquet?
    A)Обеспечивает глобальный сквозной порядок строк по указанной колонке при чтении
    B)Автоматически сжимает данные вдвое сильнее прочих доступных форматов колоночного хранения
    C)Раскладывает данные по подпапкам-значениям колонки — чтение по фильтру пропускает лишнее
    D)Убирает необходимость в индексах для ускорения выполнения запросов
    показать ответ и разбор
    +C)Раскладывает данные по подпапкам-значениям колонки — чтение по фильтру пропускает лишнее

    // разбор: partitionBy при записи раскладывает данные в подпапки по значениям колонки (date=2026-01-15/...). При чтении с фильтром по этой колонке движок применяет partition pruning — читает только нужные подпапки, резко сокращая I/O. Обратная сторона: слишком дробное партиционирование (по высококардинальной колонке) плодит тысячи мелких файлов.

дальше

Теорию прочитали. Навык ставится повторением

В Сеньорчике эта подтема идёт в ежедневных сессиях: движок возвращает её, пока ответы не станут уверенными, и ведёт прогресс отдельно по каждой подтеме. Теория внутри тоже бесплатна, лимит только на количество вопросов в день.