Spark
Средний уровеньКак настроить dynamic allocation в Spark с нуля, а не просто подкрутить постфактум?
Dynamic allocation обычно обсуждают как ручку для тюнинга, но у неё есть жёсткие предпосылки — без них она молча ничего не делает, что выглядит точно так же, как плохо настроенная задача.
Одного spark.dynamicAllocation.enabled=true на большинстве менеджеров кластера недостаточно, чтобы это заработало. Dynamic allocation сначала требует настроенного и запущенного внешнего shuffle-сервиса — без него executor'ы нельзя безопасно убрать посреди задачи, так что Spark их просто... не убирает.
Предпосылка, о которой никто не упоминает первой
Убрать executor посреди задачи безопасно, только если shuffle-данные, которые он держит, остаются
доступны другим executor'ам после его исчезновения — это и делает внешний shuffle-сервис, работающий
независимо от жизненного цикла любого конкретного executor'а. Без него включённым Spark принимает
dynamicAllocation.enabled=true без ошибки, а затем просто никогда не масштабируется вниз,
потому что это рискнуло бы потерей shuffle-данных. На YARN и Kubernetes это отдельный сервис, который нужно
развернуть и на который нужно указать в конфиге Spark — он не появляется автоматически только потому, что
dynamic allocation включён.
Реальные границы, а не просто потолок
Оставить minExecutors на значении по умолчанию 0 означает, что задача может полностью
уйти в ноль executor'ов во время затишья, а затем платить полной задержкой выделения от менеджера кластера,
чтобы масштабироваться обратно — нормально для батч-задачи с запасом по времени, дорого для всего, что
чувствительно к задержке. maxExecutors — более знакомая ручка (это потолок по стоимости), но
это лишь половина картины без соответствующего пола:
Таймаут простоя и где всё ещё нужно суждение
spark.dynamicAllocation.executorIdleTimeout (по умолчанию 60с) контролирует, сколько
executor простаивает перед освобождением — слишком коротко, и задача с рваными, неравномерными стадиями
постоянно дёргает executor'ы вверх-вниз, добавляя задержку выделения на каждом всплеске; слишком долго — и
вы платите за простаивающую мощность между всплесками. Формулы для правильного значения не существует — всё
зависит от того, насколько на самом деле «дёрганая» нагрузка, а это именно тот паттерн, который стоит
посмотреть на таймлайне реального прогона, прежде чем выбирать число.
Где здесь место opti-pipe: движок правил читает event log завершённой задачи и отмечает признаки
недо- или перепровижининга относительно того, что реально настроено в границах — он не выставляет
minExecutors/maxExecutors за вас, поскольку это компромисс между стоимостью и
задержкой, который можете сделать только вы.
Посмотрите, как это выглядит на вашем собственном пайплайне.
Загрузите реальный event log Spark, run_results.json от dbt или экспорт метрик Flink — и получите конкретные рекомендации, которые нужно одобрить, а не ещё одно эмпирическое правило.