Стоимость
Средний уровеньКак снизить стоимость Flink-кластера, не сломав чекпоинты?
Счёт не подскакивает из-за одной плохой задачи — он ползёт вверх из-за нескольких настроек, заданных один раз при запуске и больше не пересмотренных. Вот с чего начать.
Счёт за Flink-кластер меняется не так, как счёт за batch-задачу — вы платите за TaskManager'ы, которые работают независимо от того, есть ли backpressure, и за чекпоинты, которые срабатывают по таймеру независимо от того, сколько реально изменилось. Поэтому дело не в одной плохой задаче, а в настройках, которые никто не пересматривал со времён запуска.
1. Интервал чекпоинтов — рычаг стоимости, который никто не пересматривает
Каждый чекпоинт — это полный проход по state backend: для RocksDB это дисковый I/O плюс, в облачном развёртывании, пачка PUT-запросов к хранилищу объектов, указанному в state.checkpoints.dir. Интервал, заданный консервативно коротким на старте, когда никто не хотел рисковать прогрессом непроверенной задачи, продолжает обходиться в ту же сумму и спустя месяцы стабильной работы. Увеличение execution.checkpointing.interval напрямую снижает эти регулярные расходы — компромисс в том, что после следующего сбоя придётся переобработать больше данных, поэтому стоит ориентироваться на реально допустимое время восстановления, а не просто на то, насколько коротким может быть число.
2. Параллелизм, рассчитанный на худший день, оплачивается каждый день
Число слотов TaskManager обычно задают один раз, исходя из оценки пиковой нагрузки, и больше не трогают — поэтому в стационарном режиме трафик оплачивает мощность, которая большую часть времени не используется. Reactive Mode или автомасштабирование на уровне платформы закрывают этот разрыв автоматически; даже без них достаточно сравнить backpressure по каждому subtask с выделенными слотами, чтобы поймать очевидный случай: параллелизм, рассчитанный на всплеск, который случается дважды в год, а оплачивается по полной каждый день между ними.
3. Неограниченный рост состояния незаметно раздувает чекпоинты
Состояние по ключу, для которого не задан TTL, растёт всё время работы задачи, и каждый чекпоинт вынужден сериализовать его целиком — поэтому длительность чекпоинтов и объём хранилища месяцами ползут вверх без единого деплоя, который можно было бы в этом обвинить. Ничего не ломается — именно поэтому это остаётся незамеченным. Применение TTL, соответствующего тому, сколько состояния реально требует бизнес-логика, обычно даёт больший выигрыш, чем кажется, именно потому что это исправляет медленный дрейф, а не одно неверное значение.
4. Не каждый TaskManager — безопасный кандидат для spot
Безопасность запуска TaskManager на spot-инстансах зависит от того, что он держит: задача, которая чисто восстанавливается из последнего чекпоинта с запасом по времени восстановления, достаточным, чтобы пережить вытеснение, — хороший кандидат, тот же бюджет времени восстановления, что и в рычаге №1. Задача с большим состоянием по ключу, на восстановление которого с нуля уходят минуты, — кандидат похуже, поскольку вытеснение обойдётся дороже, чем сэкономленные на spot вычисления.
Что opti-pipe читает из экспорта метрик Flink: длительность и размер чекпоинтов, а также backpressure по каждому subtask — никогда бизнес-логику вашего графа задач — чтобы показать, какой из этих четырёх рычагов действительно стоит использовать в вашем конкретном развёртывании, с конкретной цифрой вместо общего правила.
Посмотрите, как это выглядит на вашем собственном пайплайне.
Загрузите реальный event log Spark, run_results.json от dbt или экспорт метрик Flink — и получите конкретные рекомендации, которые нужно одобрить, а не ещё одно эмпирическое правило.