Как только поток нужно не просто фильтровать, а агрегировать, джойнить или дедуплицировать, появляется состояние — и вместе с ним вопрос, что будет при падении задачи. Наивный ответ «пересчитаем с начала» не работает на бесконечном потоке: у него нет начала, к которому можно вернуться. Flink делает состояние гражданином первого класса — оно живёт внутри операторов, локально к ключу; переживает рестарты через распределённые снимки; и — при аккуратной настройке стоков — даёт exactly-once на всём пути наружу, а не только внутри движка. Эта статья про то, как устроено состояние, как из периодических снимков вырастает восстановление после сбоя, и как то же самое двухфазное подтверждение доводит exactly-once до внешней Kafka.
Вторая статья серии «Apache Flink: глубокое погружение». Опирается на модель времени и окон из первой статьи — здесь окно TUMBLE из неё превращается в состояние, которое нельзя потерять.
В статье
- Зачем состояние и где оно живёт
- State backends: heap против RocksDB
- Checkpoints: снимок на ходу
- Восстановление после падения TaskManager
- Exactly-once до самой Kafka: 2PC
- Контраст: at-least-once и повторы
- Savepoints против checkpoints
- Эксплуатационные ошибки
- Checklist: состояние и гарантии в проде
- Демо и версии
- Документация и первоисточники
Зачем состояние и где оно живёт
Любая операция сложнее «отфильтровать и переслать» помнит что-то между сообщениями. Оконная сумма копит частичный итог, пока окно не закроется. Джойн двух потоков держит буфер одной стороны в ожидании совпадения с другой. Дедупликация хранит множество уже виденных ключей. Счётчик уникальных пользователей — структуру, аппроксимирующую 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 (rocksdb → hashmap); на нашем маленьком окне разница по скорости незаметна, а вот инкрементальность у hashmap пропадёт сразу.
Checkpoints: снимок на ходу
Состояние локально и переживает отдельные события — но чтобы оно пережило падение узла, его надо периодически выгружать в надёжное хранилище. Именно это делает checkpoint: согласованный снимок состояния всех операторов на один и тот же логический момент потока. «Согласованный» — ключевое слово: снимок должен соответствовать состоянию, как если бы все операторы остановились ровно после одних и тех же входных событий. Наивно это потребовало бы остановить весь конвейер (stop-the-world). Flink так не делает.
Механизм — асинхронные барьерные снимки (asynchronous barrier snapshotting), инженерная адаптация алгоритма распределённых снимков Chandy–Lamport. Идея:
- JobManager периодически (в демо — раз в 10 секунд,
execution.checkpointing.interval: 10s) впрыскивает в источники специальную метку — барьер с номером чекпоинта. Барьер течёт по dataflow-графу вместе с данными, не обгоняя их. - Когда оператор получает барьер по всем входам, он знает: «всё, что было до барьера, я уже обработал; всё, что после, — ещё нет». В этот момент он снимает свою часть состояния и асинхронно выгружает её в хранилище чекпоинтов, продолжая при этом обрабатывать поток. Барьер он пропускает дальше по графу.
- Когда барьер дошёл до всех стоков и все операторы подтвердили выгрузку, 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. Как это работает по фазам:
- Prepare. Между двумя чекпоинтами сток пишет результаты в Kafka внутри открытой транзакции. Данные физически ложатся в лог топика, но помечены как незакоммиченные — потребитель с
read_committedих не видит. - 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-стока. Результат оконной агрегации TUMBLE — append-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 taskmanager → up -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 |
Документация и первоисточники
- Apache Flink 1.20 — документация (state backends, checkpointing, fault tolerance).
- Flink 1.20 — State & Fault Tolerance — чекпоинты, инкрементальность, exactly-once-семантика.
- Flink 1.20 — Kafka connector, delivery guarantees —
sink.delivery-guarantee, transactional.id, 2PC. - K. M. Chandy, L. Lamport. Distributed Snapshots: Determining Global States of Distributed Systems (1985) — теоретическая основа асинхронных барьерных снимков Flink.
- Соседние статьи серии: модель, event-time и watermarks (#1), SQL, коннекторы и эксплуатация (#3). Смежное: exactly-once в Kafka, гарантии доставки и идемпотентность, LSM- и B-tree-движки храненияСкоро.
Комментарии