Состояние и exactly-once в Flink: checkpoints, savepoints, 2PC

Стрим-процессор без надёжного состояния бесполезен: агрегаты, джойны, дедуп — всё это состояние, которое нельзя потерять при падении. Как Flink это решает: keyed vs operator state, state backends (heap vs RocksDB для состояния больше памяти), распределённые снапшоты через барьеры (алгоритм Chandy-Lamport), инкрементальные checkpoints и savepoints для апгрейда, и exactly-once end-to-end через two-phase commit в стоки

Как только поток нужно не просто фильтровать, а агрегировать, джойнить или дедуплицировать, появляется состояние — и вместе с ним вопрос, что будет при падении задачи. Наивный ответ «пересчитаем с начала» не работает на бесконечном потоке: у него нет начала, к которому можно вернуться. Flink делает состояние гражданином первого класса — оно живёт внутри операторов, локально к ключу; переживает рестарты через распределённые снимки; и — при аккуратной настройке стоков — даёт exactly-once на всём пути наружу, а не только внутри движка. Эта статья про то, как устроено состояние, как из периодических снимков вырастает восстановление после сбоя, и как то же самое двухфазное подтверждение доводит exactly-once до внешней Kafka.

Вторая статья серии «Apache Flink: глубокое погружение». Опирается на модель времени и окон из первой статьи — здесь окно TUMBLE из неё превращается в состояние, которое нельзя потерять.

Белки Apache Flink прячут жёлуди-состояние в нору-хранилище RocksDB, камера-checkpoint снимает снимок стейта; ветка TaskManager (FAILED) обломилась, состояние восстанавливается из последнего снимка (NOTHING LOST, NOTHING DUPLICATED), транзакционный ящик EXACTLY-ONCE DELIVERY COMMITTED уходит в Kafka

В статье

Зачем состояние и где оно живёт

Любая операция сложнее «отфильтровать и переслать» помнит что-то между сообщениями. Оконная сумма копит частичный итог, пока окно не закроется. Джойн двух потоков держит буфер одной стороны в ожидании совпадения с другой. Дедупликация хранит множество уже виденных ключей. Счётчик уникальных пользователей — структуру, аппроксимирующую cardinality. Всё это — состояние, и на бесконечном потоке его нельзя ни выбросить, ни пересчитать с нуля: истории «с начала» не существует.

Flink различает два вида состояния.

  • Keyed state — привязано к ключу потока. После keyBy (в SQL — GROUP BY или PARTITION BY) каждый ключ попадает всегда на один и тот же экземпляр оператора, и его состояние живёт локально там же. Это и есть причина, по которой ключ определяет распределение: Flink шардирует состояние по ключу, как БД шардирует строки. Типы keyed state — ValueState (одно значение на ключ), ListState, MapState, ReducingState/AggregatingState (сворачивают на лету). Оконная агрегация из демо — ровно этот случай: на каждый user_id и каждое открытое окно движок держит накапливаемую сумму.
  • Operator state — привязано к экземпляру оператора, а не к ключу. Классический пример — оффсеты, которые Kafka-источник помнит, докуда дочитал каждую партицию. Его немного, и обычно оно скрыто внутри коннекторов.

Ключевой факт: состояние локально. Оно не в общей БД, не в Redis рядом — оно в памяти или на диске того самого TaskManager, который обрабатывает данный ключ. Это то, что делает Flink быстрым (нет сетевого хопа за состоянием на каждое событие) и одновременно то, что нужно надёжно снимать в снимок, иначе падение узла унесёт состояние с собой. Где именно физически лежит это локальное состояние — определяет state backend.

State backends: heap против RocksDB

State backend — это подсистема, отвечающая на два вопроса: где хранить рабочее состояние между событиями и как отдавать его в чекпоинт. В Flink 1.20 их два, и выбор между ними — это выбор «скорость против объёма».

hashmap (heap) rocksdb
Где живёт рабочее состояние объекты на JVM-heap TaskManager встроенный LSM-движокСкоро на локальном диске
Доступ к состоянию без сериализации, максимально быстро (де)сериализация на каждый доступ, медленнее
Потолок объёма ограничено heap → риск OOM на больших окнах/множестве ключей ограничено диском, состояние больше RAM
Инкрементальные чекпоинты недоступны (пишется весь снимок) доступны — пишется только дельта
Когда брать небольшое состояние, минимальная латентность большое состояние, длинные окна, много ключей

RocksDB под капотом — LSM-дерево: записи ложатся в memtable, флашатся в отсортированные SST-файлы на диске, периодически компактятся. Отсюда два свойства, которые Flink переиспользует напрямую. Первое — состояние держится на диске, поэтому оно может быть кратно больше оперативной памяти узла (тот самый сценарий, где hashmap уходит в OOM). Второе, важнейшее для чекпоинтов, — SST-файлы иммутабельны: раз записанный файл больше не меняется. Значит, для очередного снимка достаточно сохранить только те SST-файлы, которые появились с прошлого раза, — это и есть инкрементальные чекпоинты.

В демо стоит state.backend.type: rocksdb и execution.checkpointing.incremental: true. На вкладке Checkpoints в дашборде это видно по размеру снимков: первый большой (полный набор SST), последующие — маленькие дельты. С hashmap инкрементальность недоступна в принципе: heap-снимок целостен только целиком, дельту из него не вырезать, каждый чекпоинт пишет всё состояние заново.

Переключить backend в стенде — одна строка в docker-compose.yml (rocksdbhashmap); на нашем маленьком окне разница по скорости незаметна, а вот инкрементальность у hashmap пропадёт сразу.

Checkpoints: снимок на ходу

Состояние локально и переживает отдельные события — но чтобы оно пережило падение узла, его надо периодически выгружать в надёжное хранилище. Именно это делает checkpoint: согласованный снимок состояния всех операторов на один и тот же логический момент потока. «Согласованный» — ключевое слово: снимок должен соответствовать состоянию, как если бы все операторы остановились ровно после одних и тех же входных событий. Наивно это потребовало бы остановить весь конвейер (stop-the-world). Flink так не делает.

Механизм — асинхронные барьерные снимки (asynchronous barrier snapshotting), инженерная адаптация алгоритма распределённых снимков Chandy–Lamport. Идея:

  1. JobManager периодически (в демо — раз в 10 секунд, execution.checkpointing.interval: 10s) впрыскивает в источники специальную метку — барьер с номером чекпоинта. Барьер течёт по dataflow-графу вместе с данными, не обгоняя их.
  2. Когда оператор получает барьер по всем входам, он знает: «всё, что было до барьера, я уже обработал; всё, что после, — ещё нет». В этот момент он снимает свою часть состояния и асинхронно выгружает её в хранилище чекпоинтов, продолжая при этом обрабатывать поток. Барьер он пропускает дальше по графу.
  3. Когда барьер дошёл до всех стоков и все операторы подтвердили выгрузку, JobManager помечает чекпоинт завершённым.

Барьер как бы разрезает бесконечный поток на «до» и «после» одной консистентной линией через весь граф — без остановки обработки. У оператора с несколькими входами (после keyBy, джойна) барьеры по разным входам приходят не одновременно, и здесь возникает выравнивание барьеров (barrier alignment): оператор, получивший барьер по одному входу, придерживает этот вход, пока барьер не придёт по остальным, — только так снимок остаётся согласованным. Ценой служит рост задержки под нагрузкой; для линейных конвейеров вроде нашего демо выравнивать нечего.

Чекпоинты хранятся вне TaskManager — в проде это S3/HDFS, в стенде для простоты локальный том (state.checkpoints.dir: file:///tmp/flink-checkpoints). Смысл один: снимок должен пережить смерть узла, который его породил.

Восстановление после падения TaskManager

Проверка, ради которой всё и строилось: убить TaskManager посреди работы и увидеть, что состояние восстановилось, а не потерялось и не задвоилось.

docker compose kill taskmanager
docker compose up -d taskmanager

В дашборде джоба уходит в RESTARTING, затем RUNNING. Что происходит под капотом: JobManager замечает пропажу TaskManager, останавливает джобу и перезапускает её с последнего успешного чекпоинта. Каждый оператор поднимает свою часть состояния из снимка, источники перематывают оффсеты Kafka на позицию, зафиксированную в том же чекпоинте, — и обработка продолжается ровно оттуда.

Критично, что происходит с частично накопленным окном. Открытое, но ещё не закрытое окно TUMBLE — это живое keyed state в RocksDB. При рестарте оно восстанавливается из чекпоинта, а не пересчитывается с нуля: сумма по user_id, накопленная до момента снимка, возвращается как была. Ни один результат окна не исчезает и ни один не появляется дважды — окна не задваиваются и не теряются. Именно за это отвечает связка «оффсеты источника + состояние оператора сняты одним согласованным барьером»: перемотка входа и восстановленное состояние всегда соответствуют одной и той же логической точке потока.

Важно не путать это восстановление внутри Flink с гарантией наружу. То, что описано выше, — exactly-once эффекта на состояние: после сбоя внутреннее состояние такое, будто сбоя не было. Но результат уже утёк в Kafka-сток, и без отдельного механизма при рестарте сток мог бы переписать те же окна повторно. За exactly-once на выходе отвечает следующий раздел.

Exactly-once до самой Kafka: 2PC

Exactly-once бывает двух разных обещаний, и их постоянно смешивают:

  • Внутри Flink — состояние после восстановления согласовано (предыдущий раздел). Это Flink даёт «из коробки», как только включены чекпоинты.
  • End-to-end — каждый входной эффект отражён во внешнем стоке ровно один раз, даже после сбоев. Это требует, чтобы сток умел транзакции и участвовал в протоколе вместе с чекпоинтами.

Kafka-сток такое умеет. В демо это две строки в определении таблицы:

CREATE TEMPORARY TABLE user_totals (
    user_id      INT,
    window_end   TIMESTAMP(3),
    total_amount BIGINT
) WITH (
    'connector' = 'kafka',
    'topic' = 'user-totals',
    'properties.bootstrap.servers' = 'kafka:9092',
    'format' = 'json',
    'sink.delivery-guarantee' = 'exactly-once',
    'sink.transactional-id-prefix' = 'user-totals-tx'
);

sink.delivery-guarantee = 'exactly-once' включает двухфазное подтверждение (two-phase commit), привязанное к чекпоинтам, а sink.transactional-id-prefix даёт префикс transactional.id для транзакций Kafka. Как это работает по фазам:

  1. Prepare. Между двумя чекпоинтами сток пишет результаты в Kafka внутри открытой транзакции. Данные физически ложатся в лог топика, но помечены как незакоммиченные — потребитель с read_committed их не видит.
  2. Commit. Когда чекпоинт успешно завершается (все операторы сняли состояние), Flink коммитит транзакцию Kafka. Ровно в этот момент записи становятся видимыми для read_committed-потребителей.

Барьер чекпоинта играет роль команды «commit» распределённой транзакции: транзакция Kafka фиксируется тогда и только тогда, когда зафиксирован снимок состояния. При падении между чекпоинтами незакоммиченная транзакция откатывается (Kafka сама отбросит её по transactional.id), джоба стартует с прошлого чекпоинта и переписывает окна заново — но старой, откаченной, записи снаружи никто не видел. В итоге каждое закрытое окно попадает в топик ровно один раз.

Проверяется это isolation-level-ом потребителя:

docker compose exec kafka /opt/kafka/bin/kafka-console-consumer.sh \
  --bootstrap-server kafka:9092 --topic user-totals --from-beginning \
  --isolation-level read_committed

С read_committed в топике — по одной записи на каждое (user_id, window_end), без дублей даже после kill TaskManager.

Отдельно стоит понять, почему здесь хватает обычного Kafka-стока. Результат оконной агрегации TUMBLEappend-only: каждое окно закрывается один раз и порождает ровно одну финальную строку, которую больше не пересматривают. Апдейтов существующих строк нет, поэтому достаточно транзакционного append + 2PC. Если бы запрос был обновляющим (например, бегущий GROUP BY без окна, который на каждое событие уточняет текущий итог по ключу), поток стал бы changelog с ретракциями, и понадобился бы upsert-сток (upsert-kafka с ключом): там exactly-once обеспечивается не транзакцией на весь батч, а тем, что ключ+последнее значение идемпотентны при переигрывании. Наш случай проще именно потому, что окно даёт финальные, неизменяемые результаты.

Цена exactly-once — латентность. Результат окна виден потребителю не в момент закрытия окна, а в момент следующего успешного чекпоинта: до коммита он лежит в топике невидимым. С интервалом 10 секунд это добавляет к задержке до ~10 секунд сверх времени окна. Хотите ниже задержку — уменьшайте интервал чекпоинта (ценой накладных расходов на снимки) либо смиритесь с at-least-once. Эта развилка — частный случай общего компромисса из гарантий доставки и идемпотентности; та же логика транзакционной публикации живёт и в паттерне transactional outbox, и в событийных системах вроде Event Sourcing на практике, где проекции обязаны быть идемпотентны по той же причине.

Контраст: at-least-once и повторы

Чтобы увидеть, что exactly-once не бесплатная магия, а конкретный механизм, достаточно его выключить. Меняем в стоке одну настройку:

'sink.delivery-guarantee' = 'at-least-once'

Теперь сток пишет в Kafka без транзакций, коммитя оффсеты по чекпоинтам, но не пряча записи до коммита. Повторяем сценарий: docker compose kill taskmanagerup -d taskmanager. После рестарта с последнего чекпоинта сток заново записывает окна, которые он успел отправить между последним чекпоинтом и падением, — но которые уже физически в топике. Дублей.

Увидеть их можно, читая без фильтра по коммиту:

docker compose exec kafka /opt/kafka/bin/kafka-console-consumer.sh \
  --bootstrap-server kafka:9092 --topic user-totals --from-beginning \
  --isolation-level read_uncommitted

С read_uncommitted для at-least-once-стока видны повторы одних и тех же окон. Разница между двумя режимами — ровно разница между «сток пишет и надеется» и «сток пишет в транзакции, а коммитит по барьеру чекпоинта». at-least-once дешевле по латентности (результат виден сразу), но перекладывает дедупликацию на потребителя: тому придётся делать применение идемпотентным по ключу окна.

Savepoints против checkpoints

И checkpoint, и savepoint — снимки состояния джобы, но у них разное назначение, и путать их дорого.

Checkpoint Savepoint
Кто инициирует Flink автоматически, по расписанию человек/оператор, вручную
Зачем восстановление после сбоя осознанная остановка: апгрейд кода, смена параллелизма, миграция
Формат оптимизирован под движок (в т.ч. инкрементальный) самодостаточный, переносимый снимок
Жизненный цикл обычно чистятся автоматически живут, пока их явно не удалят
Стоимость снятия дёшево (инкрементально) дороже (полный, канонический)

Смысловая разница: checkpoint — про отказоустойчивость (машина сама поднимается с последнего снимка, как в разделе про kill TaskManager), savepoint — про управляемую эволюцию джобы. Хотите выкатить новую версию кода, увеличить параллелизм оператора или перенести джобу в другой кластер — снимаете savepoint, останавливаете джобу, запускаете новую версию из этого savepoint. Состояние переносится, окна не теряются, обработка продолжается с той же точки.

Заметьте: в стенде стоит execution.checkpointing.externalized-checkpoint-retention: RETAIN_ON_CANCELLATION — чекпоинты не удаляются при отмене джобы, что позволяет вручную стартовать с них как с savepoint. Это удобный мостик, но полноценная работа с savepoints — совместимость состояния между версиями оператора, uid-ы, изменение параллелизма, эволюция схемы состояния — это тема третьей статьи серии про эксплуатацию Flink. Здесь достаточно запомнить границу: авто-снимок для восстановления — checkpoint, ручной снимок для апгрейда — savepoint.

Эксплуатационные ошибки

  • Ждать exactly-once наружу без транзакционного стока. Чекпоинты дают согласованность состояния внутри Flink, но не мешают стоку записать результат дважды при рестарте. End-to-end exactly-once требует sink.delivery-guarantee = exactly-once (или идемпотентного upsert-стока) — сами по себе чекпоинты его не обеспечивают.
  • Читать exactly-once-топик с read_uncommitted. Транзакционный сток кладёт данные в топик до коммита. Потребитель без isolation-level=read_committed увидит и незакоммиченные, и откаченные записи — то есть те самые «дубли», от которых exactly-once и защищает. Гарантия существует только в паре со read_committed на стороне читателя.
  • Не задавать/переиспользовать transactional-id-prefix. Префикс transactional.id должен быть уникальным на джобу. Две джобы с одинаковым префиксом будут воевать за одни transactional.id в Kafka и ломать транзакции друг друга.
  • Таймаут транзакции стока больше, чем разрешает брокер. Flink-сток запрашивает transaction.timeout.ms в час, а Kafka по умолчанию ограничивает transaction.max.timeout.ms пятнадцатью минутами. Сток не может инициализировать продюсер, джоба уходит в бесконечный RESTARTING, а в логе TaskManager — The transaction timeout is larger than the maximum value allowed by the broker. Лечится на стороне брокера (в стенде KAFKA_TRANSACTION_MAX_TIMEOUT_MS), а не занижением таймаута стока: транзакция, протухшая раньше, чем джоба успела восстановиться, откатывается — и результат окна теряется совсем. Таймаут должен покрывать интервал чекпоинта плюс реалистичное время простоя.
  • Слишком редкий интервал чекпоинта при exactly-once. Результат виден потребителю только после коммита, привязанного к чекпоинту. Интервал в минуту означает задержку публикации до минуты. Латентность стока — это интервал чекпоинта, а не время окна; это осознанная настройка, а не «поставил и забыл».
  • hashmap-backend на большом состоянии. Heap быстрый, но конечен. Длинные окна, session-окна, много ключей, крупный джойн-буфер — и TaskManager падает с OOM. Большое состояние — это RocksDB на диске; heap оставляют маленьким горячим состояниям.
  • Ждать инкрементальные чекпоинты от heap. execution.checkpointing.incremental: true работает только с RocksDB (иммутабельные SST-файлы). С hashmap каждый чекпоинт полный — на большом состоянии это тяжёлые периодические снимки.
  • Слишком дорогое выравнивание барьеров под нагрузкой. На широких графах с джойнами barrier alignment под backpressure раздувает задержку чекпоинта: оператор придерживает опередившие входы. Если чекпоинты стали долгими — смотреть в сторону unaligned checkpoints (за рамками демо) и в саму причину backpressure.
  • Путать checkpoint и savepoint при апгрейде. Пытаться выкатить новую версию кода «просто перезапуском с чекпоинта» — хрупко: формат чекпоинта оптимизирован под движок и не гарантирует переносимость между версиями. Для апгрейдов — savepoint.

Checklist: состояние и гарантии в проде

  • State backend выбран под объём: rocksdb для большого состояния/длинных окон, hashmap — только для маленького горячего.
  • Инкрементальные чекпоинты включены при RocksDB (execution.checkpointing.incremental: true).
  • Интервал чекпоинта осознанно сбалансирован: латентность стока при exactly-once ≈ интервал чекпоинта.
  • Чекпоинты пишутся в надёжное внешнее хранилище (S3/HDFS), переживающее смерть узла, а не в локальный том узла.
  • Для exactly-once наружу: транзакционный сток (sink.delivery-guarantee = exactly-once) с уникальным transactional-id-prefix.
  • Потребители exactly-once-топиков читают с isolation-level = read_committed.
  • Выбор append-only-стока против upsert-стока соответствует характеру запроса (оконный финал против бегущего changelog).
  • Проверено восстановление: kill TaskManager → джоба поднимается с последнего чекпоинта, результаты не задваиваются и не теряются.
  • Для апгрейдов кода/параллелизма используется savepoint, а не рестарт с чекпоинта.
  • uid-ы операторов заданы явно, чтобы состояние маппилось между версиями джобы (детали — в статье про эксплуатацию).

Демо и версии

Стенд: digital-cookbook/flink/01-state/ — оконная агрегация с состоянием в RocksDB, чекпоинты каждые 10 секунд, восстановление после kill taskmanager и exactly-once-сток в Kafka. Полный SQL — в sql/exactly-once.sql; сценарии kill/restore и переключение read_committed/read_uncommitted — в README стенда. Важно честно: «нет дублей» в стенде сейчас проверяется вручную (глазами в read_committed-консюмере после kill/restore), без автоматического smoke-теста или ассерта на счётчики — для регрессии такую проверку стоит добавить.

Версии (прогнано живьём):

Компонент Версия
Apache Flink 1.20.1 (Java 17)
Apache Kafka 3.9.0 (KRaft)
Flink SQL Kafka connector 3.3.0-1.20
State backend RocksDB, инкрементальные чекпоинты
Интервал чекпоинта 10 s

Документация и первоисточники

Обсуждение в Telegram

Присоединиться →

Комментарии