Flink

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

Как настроить state backend и хранилище чекпоинтов Flink перед выходом в прод?

Настройки Flink по умолчанию заточены под то, чтобы задача быстро запустилась в разработке, а не под то, чтобы пережить реальный рестарт с реальным объёмом состояния. Почти всё решают две настройки.

Читать на:
EmbeddedRocksDBStateBackend HashMapStateBackend Кастомный state backend
Как продакшен-задачи Flink с заметным объёмом состояния обычно делятся между backend'ами, примерно — не измеренная статистика.

Каждой стейтфул-задаче Flink нужен state backend (где живёт состояние оператора, пока задача работает) и хранилище чекпоинтов (куда пишутся консистентные снапшоты этого состояния). Flink нормально работает на настройках по умолчанию — ровно до тех пор, пока состояние не вырастет за пределы того, что комфортно помещается в heap JVM.

HashMapStateBackend против EmbeddedRocksDBStateBackend

HashMapStateBackend хранит всё состояние как Java-объекты в heap — быстро, но ограничено тем, сколько heap вы выделили task manager'у, и каждый чекпоинт вынужден сериализовать всё записываемое состояние целиком. EmbeddedRocksDBStateBackend хранит состояние на локальном диске (с настраиваемым кэшем в памяти) и сериализует только то, что реально нужно на конкретный чекпоинт, ценой некоторой задержки при каждом обращении к состоянию. После нескольких сотен МБ состояния на слот задачи RocksDB обычно должен быть выбором по умолчанию, а не оптимизацией «на потом».

# flink-conf.yaml state.backend.type: rocksdb state.backend.incremental: true

state.backend.incremental: true важен почти так же сильно, как сам выбор backend'а — без него каждый чекпоинт заново переписывает полный снапшот состояния вместо того, чтобы записать только изменившееся с прошлого раза, а это быстро становится дорого по мере роста состояния.

Хранилище чекпоинтов должно быть надёжным и общим, а не локальным

Хранилище чекпоинтов — это отдельная от state backend настройка: это то, куда реально попадают файлы чекпоинта, и оно должно быть доступно любому task manager'у, который восстанавливает упавшую задачу, а на практике это означает объектное хранилище, а не локальный диск:

# flink-conf.yaml state.checkpoint-storage: filesystem state.checkpoints.dir: s3://my-flink-checkpoints/prod/ state.checkpoints.num-retained: 3

num-retained намеренно хранит больше одного чекпоинта — если самый свежий окажется повреждён или недописан в момент гибели task manager'а, Flink откатится к предыдущему вместо того, чтобы восстанавливаться не с чего.

Настройка, которую больно менять потом

Переключение state backend'а после того, как у задачи уже есть реальное продакшен-состояние, означает, что новый backend не может прочитать формат чекпоинтов старого — миграции на месте нет, есть только savepoint, снятый под старым backend'ом и восстановленный под новым, с реальным простоем на это время. Выбрать RocksDB перед запуском, если есть хоть какой-то шанс роста состояния, не стоит ничего; переносить с HashMapStateBackend под нагрузкой — это плановое окно обслуживания.

Вне области opti-pipe: state backend и хранилище чекпоинтов — это конфигурация кластера/задачи, задаваемая до старта задачи — движок правил только читает метрики, которые экспортирует уже работающая или завершённая задача, и не трогает, как и где сохраняется состояние.

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

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