Skip to content
Learning Platform

Фраза «exactly-once» в маркетинге стриминговых движков означает примерно ничего, пока вы не разберётесь, что именно гарантируется и какой ценой. Apache Flink даёт одну из самых честных реализаций этой семантики в индустрии, и за ней стоит красивый механизм: распределённый снапшот состояния, который снимается без остановки потока. Разберём его до уровня, на котором становится понятно, почему растёт latency при backpressure, зачем нужен RocksDB и почему один лишь чекпоинт не делает ваш пайплайн exactly-once от источника до приёмника.

Что такое state и почему exactly-once трудно

Стриминговое приложение редко бывает stateless. Подсчёт уникальных пользователей за окно, дедупликация, джойн двух потоков, машина состояний для CEP — всё это требует, чтобы оператор помнил что-то между событиями. Это и есть state. В Flink он бывает keyed (привязан к ключу после 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 по всем входным каналам, он:

  1. снимает снапшот своего локального состояния;
  2. отправляет в JobManager подтверждение (acknowledgement), что состояние записано в durable-хранилище;
  3. пробрасывает барьер 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.

СвойствоAlignedUnaligned
Что в снапшотетолько состояние операторовсостояние + in-flight данные в каналах
Поведение при backpressureвремя чекпоинта растёт, риск таймаутапочти постоянное время
Размер чекпоинтаменьшебольше (плюс буферы каналов)
Latency обработкистрадает при alignmentстабильнее
Когда применятьстабильный поток, нет backpressureустойчивый backpressure, нужен надёжный чекпоинт
Семантикаexactly-onceexactly-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 хранит состояние в RocksDB на локальном диске TaskManager’а. Размер ограничен диском (десятки терабайт на NVMe), но каждое чтение и запись проходит через сериализацию и обращение к диску — медленнее. Зато открывается ключевая возможность.

КритерийHashMapStateBackendEmbeddedRocksDBStateBackend
Где состояниеJVM heap (объекты)локальный диск (LSM, off-heap)
Предел размераразмер heapразмер диска
Latency доступаминимальнаявыше (сериализация + диск)
Инкрементальные чекпоинтынет (всегда полный)да
РискOOM при большом stateмедленнее на маленьком state

Дело в том, что RocksDB — это LSM-дерево, и его данные на диске лежат иммутабельными SST-файлами. Между двумя чекпоинтами меняется лишь часть файлов. Это позволяет делать инкрементальные чекпоинты: в хранилище уезжают только новые и изменившиеся SST-файлы, а не весь массив состояния целиком. Для job’ов с состоянием в сотни гигабайт разница колоссальная — полный чекпоинт занимал бы минуты, инкрементальный укладывается в секунды и почти не нагружает сеть.

Расплата отложена на 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). Протокол привязывает коммит данных к жизненному циклу чекпоинта:

  1. Pre-commit. Между чекпоинтами синк пишет данные в транзакцию, которая ещё не видна потребителям (например, Kafka-транзакция или временный файл). При получении барьера N синк завершает текущую транзакцию — flush’ит данные, но не коммитит — и фиксирует её идентификатор в своём состоянии (которое попадёт в чекпоинт N).
  2. 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 с инкрементами.

  1. Прошёл чекпоинт N: offset’ы источников, состояние операторов и ID транзакций приёмника согласованно записаны в S3, чекпоинт завершён.
  2. После него job обработал ещё 50 000 заказов, обновил состояние, синк открыл транзакцию N+1 и накопил в ней результаты.
  3. В этот момент падает 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.

Ещё в направлении · Data Engineering

Все материалы направления →