Spark

Средний уровень

Как настроить dynamic allocation в Spark с нуля, а не просто подкрутить постфактум?

Dynamic allocation обычно обсуждают как ручку для тюнинга, но у неё есть жёсткие предпосылки — без них она молча ничего не делает, что выглядит точно так же, как плохо настроенная задача.

Читать на:
Внешний shuffle-сервис Границы min/max executors Тюнинг таймаута простоя
Примерно порядок, в котором эти настройки требуют внимания — shuffle-сервис не опционален, остальное — тюнинг поверх него.

Одного spark.dynamicAllocation.enabled=true на большинстве менеджеров кластера недостаточно, чтобы это заработало. Dynamic allocation сначала требует настроенного и запущенного внешнего shuffle-сервиса — без него executor'ы нельзя безопасно убрать посреди задачи, так что Spark их просто... не убирает.

Предпосылка, о которой никто не упоминает первой

Убрать executor посреди задачи безопасно, только если shuffle-данные, которые он держит, остаются доступны другим executor'ам после его исчезновения — это и делает внешний shuffle-сервис, работающий независимо от жизненного цикла любого конкретного executor'а. Без него включённым Spark принимает dynamicAllocation.enabled=true без ошибки, а затем просто никогда не масштабируется вниз, потому что это рискнуло бы потерей shuffle-данных. На YARN и Kubernetes это отдельный сервис, который нужно развернуть и на который нужно указать в конфиге Spark — он не появляется автоматически только потому, что dynamic allocation включён.

# spark-defaults.conf spark.shuffle.service.enabled true spark.dynamicAllocation.enabled true

Реальные границы, а не просто потолок

Оставить minExecutors на значении по умолчанию 0 означает, что задача может полностью уйти в ноль executor'ов во время затишья, а затем платить полной задержкой выделения от менеджера кластера, чтобы масштабироваться обратно — нормально для батч-задачи с запасом по времени, дорого для всего, что чувствительно к задержке. maxExecutors — более знакомая ручка (это потолок по стоимости), но это лишь половина картины без соответствующего пола:

# spark-defaults.conf spark.dynamicAllocation.minExecutors 2 spark.dynamicAllocation.maxExecutors 40 spark.dynamicAllocation.initialExecutors 2

Таймаут простоя и где всё ещё нужно суждение

spark.dynamicAllocation.executorIdleTimeout (по умолчанию 60с) контролирует, сколько executor простаивает перед освобождением — слишком коротко, и задача с рваными, неравномерными стадиями постоянно дёргает executor'ы вверх-вниз, добавляя задержку выделения на каждом всплеске; слишком долго — и вы платите за простаивающую мощность между всплесками. Формулы для правильного значения не существует — всё зависит от того, насколько на самом деле «дёрганая» нагрузка, а это именно тот паттерн, который стоит посмотреть на таймлайне реального прогона, прежде чем выбирать число.

Где здесь место opti-pipe: движок правил читает event log завершённой задачи и отмечает признаки недо- или перепровижининга относительно того, что реально настроено в границах — он не выставляет minExecutors/maxExecutors за вас, поскольку это компромисс между стоимостью и задержкой, который можете сделать только вы.

Посмотрите, как это выглядит на вашем собственном пайплайне.

Загрузите реальный event log Spark, run_results.json от dbt или экспорт метрик Flink — и получите конкретные рекомендации, которые нужно одобрить, а не ещё одно эмпирическое правило.