Фраза «exactly-once» в маркетинге стриминговых движков означает примерно ничего, пока вы не разберётесь, что именно гарантируется и какой ценой. Apache Flink даёт одну из самых честных реализаций этой семантики в индустрии, и за ней стоит красивый механизм: распределённый снапшот состояния, который снимается без остановки потока. Разберём его до уровня, на котором становится понятно, почему растёт latency при backpressure, зачем нужен RocksDB и почему один лишь чекпоинт не делает ваш пайплайн exactly-once от источника до приёмника.
Что такое state и почему exactly-once трудно
Стриминговое приложение редко бывает stateless. Подсчёт уникальных пользователей за окно, дедупликация, джойн двух потоков, машина состояний для CEP — всё это требует, чтобы оператор помнил что-то между событиями. Это и есть keyBy) и operator (привязан к экземпляру оператора).
Теперь представьте job из нескольких операторов с параллелизмом, скажем, четыре. Это десятки независимых subtask’ов, разбросанных по TaskManager’ам, у каждого свой кусок состояния в своей JVM. Падает один TaskManager — Flink должен перезапустить job и восстановить состояние так, будто сбоя не было: ни одно событие не потеряно и ни одно не учтено дважды.
Сложность в слове «согласованно». Нельзя просто скопировать состояние каждого subtask’а в произвольный момент: пока вы снимаете снапшот оператора A, событие уже могло уйти в оператор B, и тогда в снапшоте оно окажется учтённым дважды или потерянным. Нужен глобально согласованный срез распределённого состояния — фотография системы в один логический момент, при том что физически часы на узлах не синхронизированы, а события летят по сети асинхронно.
Наивное решение — stop-the-world: заморозить все subtask’и, скопировать состояние, разморозить. Для потока с миллионами событий в секунду это убивает latency и пропускную способность. Flink идёт другим путём.
Алгоритм чекпоинтов: барьеры Chandy-Lamport
В основе лежит алгоритм распределённых снапшотов Чанди-Лампорта (1985), адаптированный Flink под потоки и названный ABS (Asynchronous Barrier Snapshotting). Идея в том, чтобы не останавливать обработку, а пустить по графу специальные маркеры.
Координирует процесс JobManager (точнее, его CheckpointCoordinator). По таймеру он инициирует чекпоинт номер N и инъектит барьер в каждый источник. N, относится к чекпоинту N, всё, что после, — к следующему.
Дальше происходит магия. Барьер течёт по графу операторов ровно так же, как обычные данные, в порядке очереди. Когда оператор получает барьер N по всем входным каналам, он:
- снимает снапшот своего локального состояния;
- отправляет в JobManager подтверждение (acknowledgement), что состояние записано в durable-хранилище;
- пробрасывает барьер
Nдальше по всем выходным каналам.
Поскольку барьеры встроены в сам поток, граница «зачекпоинчено / не зачекпоинчено» получается естественно согласованной по всему графу, без единой точки синхронизации. Источник, дойдя до барьера, фиксирует свой offset (например, позицию в Kafka-партиции); приёмник, получив барьер, знает, что выше по течению всё уже снято. Когда JobManager собрал подтверждения от всех subtask’ов, чекпоинт N объявляется завершённым (completed) и его метаданные записываются в хранилище.
Ключевое слово — asynchronous. Запись состояния в распределённое хранилище (S3, HDFS) идёт в фоне: оператор делает быструю копию состояния, отдаёт барьер дальше и продолжает обрабатывать события, пока байты улетают на диск. Поток не стоит.
При восстановлении после сбоя Flink берёт последний завершённый чекпоинт, заливает состояние во все операторы, перематывает источники на сохранённые offset’ы и продолжает с этой точки. Никакой потери — потому что offset и состояние согласованы; никакого двойного счёта — по той же причине.
Barrier alignment: выровненные и невыровненные чекпоинты
Тонкость в том, что у оператора обычно несколько входных каналов (например, после keyBy или джойна). Барьер N приходит по ним не одновременно: один upstream-subtask быстрее, другой медленнее. Что делать?
В выровненном (aligned) чекпоинте оператор, получив барьер N по одному каналу, перестаёт читать этот канал и буферизует приходящие из него события, продолжая обрабатывать остальные каналы. Только когда барьер N пришёл по всем входам, оператор снимает снапшот и отдаёт барьер дальше. Это и есть alignment — выравнивание барьеров. Снимок получается чистым: состояние не «загрязнено» событиями из-за барьера.
Проблема — backpressure. Если канал заблокирован (медленный downstream, перекос данных), барьер по нему ползёт долго, а оператор всё это время держит alignment и копит буфер. Время чекпоинта раздувается, latency растёт, в худшем случае чекпоинты начинают таймаутить.
Невыровненные (unaligned) чекпоинты, появившиеся во Flink 1.11, решают это. Как только барьер N приходит по первому каналу, оператор немедленно «перепрыгивает» его в начало выходной очереди и снимает снапшот, включая в него in-flight данные — события, что уже в буферах входных и выходных каналов, но ещё не обработаны. Барьер не ждёт остальных. Состояние снапшота теперь шире (оператор + содержимое каналов), но чекпоинт завершается за время, почти не зависящее от backpressure.
| Свойство | Aligned | Unaligned |
|---|---|---|
| Что в снапшоте | только состояние операторов | состояние + in-flight данные в каналах |
| Поведение при backpressure | время чекпоинта растёт, риск таймаута | почти постоянное время |
| Размер чекпоинта | меньше | больше (плюс буферы каналов) |
| Latency обработки | страдает при alignment | стабильнее |
| Когда применять | стабильный поток, нет backpressure | устойчивый backpressure, нужен надёжный чекпоинт |
| Семантика | exactly-once | exactly-once |
Обратите внимание: и aligned, и unaligned дают exactly-once. Выбор — это компромисс между размером чекпоинта и его устойчивостью к заторам, а не между гарантиями. Включается опционально:
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
env.enableCheckpointing(10_000); // чекпоинт каждые 10 секунд
env.getCheckpointConfig()
.setCheckpointingMode(CheckpointingMode.EXACTLY_ONCE);
env.getCheckpointConfig().enableUnalignedCheckpoints();
State backends: где живёт состояние и инкрементальные чекпоинты
Снапшот снапшотом, но состояние сначала где-то лежит в runtime. За это отвечает state backend, и выбор между двумя вариантами — одно из главных архитектурных решений.
HashMapStateBackend держит состояние как обычные Java-объекты в JVM heap. Доступ мгновенный, без сериализации, но размер упёрт в heap: переоцените — получите OOM. Подходит, когда состояние помещается в память (гигабайты на subtask).
EmbeddedRocksDBStateBackend хранит состояние в
| Критерий | HashMapStateBackend | EmbeddedRocksDBStateBackend |
|---|---|---|
| Где состояние | JVM heap (объекты) | локальный диск (LSM, off-heap) |
| Предел размера | размер heap | размер диска |
| Latency доступа | минимальная | выше (сериализация + диск) |
| Инкрементальные чекпоинты | нет (всегда полный) | да |
| Риск | OOM при большом state | медленнее на маленьком state |
Дело в том, что RocksDB — это
Расплата отложена на recovery: чтобы восстановиться, Flink собирает состояние из цепочки инкрементов, и восстановление может быть дольше, чем при полном чекпоинте. Классический компромисс: дешёвый штатный путь против более дорогого редкого восстановления. Для HashMap-бэкенда инкрементов нет — он всегда снимает полный снапшот.
End-to-end exactly-once: чекпоинты плюс транзакционные синки
Здесь живёт главное недопонимание. Чекпоинты дают exactly-once внутри Flink: состояние и offset’ы источников согласованы, при сбое всё откатывается к последнему чекпоинту. Но между двумя чекпоинтами job уже что-то отправил в приёмник — в Kafka, базу, файл. При откате к чекпоинту источник перемотается назад и повторно обработает события, а значит, заново отправит в приёмник то, что уже туда ушло. Это дубликаты на стороне приёмника — то есть at-least-once end-to-end, даже при exactly-once внутри.
Чтобы дотянуть гарантию до приёмника, нужен синк, умеющий двухфазный коммит (2PC), согласованный с чекпоинтами. В Flink это абстракция TwoPhaseCommitSinkFunction (и её современная реализация в KafkaSink с DeliveryGuarantee.EXACTLY_ONCE). Протокол привязывает коммит данных к жизненному циклу чекпоинта:
- Pre-commit. Между чекпоинтами синк пишет данные в транзакцию, которая ещё не видна потребителям (например, Kafka-транзакция или временный файл). При получении барьера
Nсинк завершает текущую транзакцию — flush’ит данные, но не коммитит — и фиксирует её идентификатор в своём состоянии (которое попадёт в чекпоинтN). - Commit. Когда JobManager сообщает, что чекпоинт
Nзавершён глобально (notifyCheckpointComplete), синк коммитит транзакцию, и данные становятся видимыми.
Логика восстановления вытекает отсюда сама. Если сбой случился до завершения чекпоинта — транзакция не закоммичена, при рестарте она абортится, дубликатов нет. Если сбой случился после завершения чекпоинта, но до получения сигнала о коммите — Flink при восстановлении знает ID транзакции из состояния и повторяет коммит (commit идемпотентен). В обоих случаях каждое событие материализуется в приёмнике ровно один раз.
KafkaSink<String> sink = KafkaSink.<String>builder()
.setBootstrapServers(brokers)
.setRecordSerializer(serializer)
.setDeliveryGuarantee(DeliveryGuarantee.EXACTLY_ONCE)
.setTransactionalIdPrefix("orders-sink")
.build();
Цена end-to-end exactly-once — латентность видимости данных. Потребитель Kafka видит записи только после коммита транзакции, то есть после завершения чекпоинта. Если чекпоинт раз в 10 секунд, downstream видит данные с задержкой до этих 10 секунд. Поэтому exactly-once — не «бесплатно лучше», а осознанный выбор: для биллинга и финансов он обязателен, для дашборда метрик at-least-once с идемпотентным UPSERT нередко дешевле и достаточно.
Подытожим, что складывается в end-to-end гарантию:
- exactly-once внутри Flink: чекпоинты Chandy-Lamport плюс перематываемые источники с offset’ами;
- exactly-once на выходе: транзакционный 2PC-синк, коммитящий по завершении чекпоинта;
- если выпадает хотя бы одно звено — фактическая гарантия деградирует до at-least-once, как бы ни был настроен сам Flink.
Эта механика — барьеры, alignment, бэкенды, 2PC — детально разбирается с запускаемыми примерами в бесплатном курсе по Apache Flink: можно прогнать чекпоинт, уронить TaskManager и посмотреть, как состояние восстанавливается из снапшота.
Сквозной пример: что происходит при сбое
Соберём всё вместе на конкретном сценарии. Job читает заказы из Kafka, считает в keyed state сумму по пользователю и пишет агрегаты в Kafka exactly-once. Чекпоинты — раз в 10 секунд, бэкенд RocksDB с инкрементами.
- Прошёл чекпоинт
N: offset’ы источников, состояние операторов и ID транзакций приёмника согласованно записаны в S3, чекпоинт завершён. - После него job обработал ещё 50 000 заказов, обновил состояние, синк открыл транзакцию
N+1и накопил в ней результаты. - В этот момент падает TaskManager — например, OOM или вытеснение пода в Kubernetes.
Что делает Flink. JobManager обнаруживает потерю TaskManager’а и перезапускает job из последнего завершённого чекпоинта N. Источники перематываются на offset’ы из N, состояние операторов заливается из снапшота (для RocksDB — собирается из цепочки инкрементальных SST-файлов), а транзакция N+1, не успевшая закоммититься, абортится на стороне Kafka. Те 50 000 заказов будут прочитаны и обработаны заново — но в приёмник они уйдут впервые, потому что прежняя транзакция отменена и её данные никогда не были видны потребителям.
Результат: внутреннее состояние согласовано (счётчики не задвоились), на выходе ни одного дубликата. Именно совместная работа трёх механизмов — перематываемого источника, согласованного снапшота состояния и транзакционного приёмника — превращает «exactly-once внутри» в «exactly-once end-to-end». Уберите любой: при at-least-once-источнике потеряете гарантию на входе, при обычном (не транзакционном) синке — задвоите на выходе.
Отдельно стоит понимать, что exactly-once здесь — не про физическую доставку каждого байта ровно один раз, а про effectively-once на уровне эффектов: события могут физически обрабатываться повторно после отката, но их видимый эффект (изменение состояния, запись в приёмник) проявляется ровно один раз. Это различие снимает большую часть «магии» из маркетингового термина.
Что дальше
Exactly-once в Flink — это не одна кнопка, а согласованная цепочка: состояние операторов, барьеры, выбор бэкенда и транзакционный приёмник. Понимание того, как барьер течёт по графу и почему unaligned-чекпоинт спасает под backpressure, отличает инженера, который чинит таймауты чекпоинтов в проде, от того, кто просто включил галочку и надеется.
Хотите пройти путь от первого DataStream до production-grade exactly-once с runnable-песочницей прямо в браузере — это в бесплатном курсе по Apache Flink на нашей платформе. Все материалы открыты, на русском, с разбором внутренностей до железа. Полный каталог инженерии данных — в направлении Data Engineering.