Apache Flink deep dive: true streaming engine with stateful operators, RocksDB state backend, Chandy-Lamport distributed snapshots for exactly-once via 2PC sinks. Shows JobManager + 3 TaskManagers (source, window aggregator, sink), RocksDB local state, S3 snapshot storage, Kafka transactional sink + Postgres XA sink. Five scenarios: stateful keyed aggregation, checkpoint barrier propagation, 2PC sink commit, failure recovery with tx abort, savepoint-based version upgrade. Includes 3 ADRs comparing Flink vs Spark Structured Streaming vs Kafka Streams.
Flink упоминается в каждой второй streaming-теме, но без отдельного глубокого разбора остаётся «один из тех движков». А Flink — не «Spark Streaming на Java». Это другая модель (true event-at-a-time streaming с event-time + watermarks, не micro-batch), другая state-machine (stateful operators с RocksDB и distributed snapshots по Chandy-Lamport) и другие exactly-once гарантии (через two-phase commit sinks, а не «idempotent consumer»).
Понимать Flink надо, потому что он де-факто стандарт для mission-critical real-time пайплайнов: Uber pricing, Netflix personalization, Alibaba Singles Day, Stripe Radar. Если выбирается streaming engine для production с большим state и end-to-end exactly-once — Flink рассматривается первым.
Flink — это распределённый dataflow runtime. Операторы (map, keyBy, window, join, sink) держат local state в RocksDB или JVM heap. Периодически JobManager инжектит barriers в источники, и каждый оператор по приходу barrier снимает consistent snapshot своего state (Chandy-Lamport). На failure job рестартует с последнего checkpoint, source rewindит Kafka offset, 2PC sinks abortят pending транзакции — событие воспроизводится, но downstream видит ровно один результат.
Три ключевых отличия от micro-batch:
Реальный fraud-pipeline на 200K events/sec:
events (12 partitions).tm-1 source task, tm-2 keyBy+window aggregator (хранит 1TB feature state, шардирован по cardId), tm-3 2PC sink task.fraud-decisions (transactional producer) + Postgres feature store (XA 2PC).Edges не «направление данных», а физические каналы. JM связан со всеми TM heartbeat'ами и инжектит barriers в source. Snapshot-стрелки в S3 — async upload SST файлов и метаданных. Ответы в обратную сторону (ACK checkpoint, notifyCheckpointComplete) идут по тем же рёбрам в reverse.
Один event проходит весь pipeline. Source-task TM-1 читает offset 99812 из Kafka, сохраняет его в managed operator state (rocks-1), назначает event-time из payload (не processing time!), периодически эмитит watermark max_seen_ts - 5s. keyBy(cardId) шафлит к TM-2, который читает 5-минутное окно из RocksDB, обновляет (4 events, $570), пишет обратно, дожидается watermark > window.end и эмитит результат. TM-3 пакует в pending Kafka transaction + XA branch. Событие не visible downstream до commit на следующем checkpoint — это нормально, latency p99 ~80ms.
Каждые 60 секунд Checkpoint Coordinator на JM инжектит barrier #142 в source. Source записывает Kafka offset в state и эмитит barrier downstream. TM-2 на window aggregator выравнивает barrier по всем входам (для join'ов важно), делает RocksDB checkpoint — это hard-link SST файлов, синхронная часть ~10ms, после чего обработка продолжается, а копирование SST в S3 идёт async. TM-3 на sink flushит Kafka producer (pending tx), делает XA prepare на Postgres. Все ACK уходят на JM, JM пишет _metadata в S3, затем вызывает notifyCheckpointComplete(142) — это триггер commit-фазы 2PC.
Two-phase commit раскрывается. Phase 1: на приходе barrier sink делает producer.flush() (события на брокере, но pending — consumer с isolation.level=read_committed их не видит) + xa_prepare() на Postgres (Postgres voted YES, держит локи). Phase 2 запускается ТОЛЬКО когда checkpoint глобально complete (метаданные в S3): JM зовёт notifyCheckpointComplete(142), sink делает producer.commitTransaction() + xa_commit(branch=142). Теперь события видны downstream. Атомарность с checkpoint-ом — если что-то упадёт между phase 1 и phase 2, аборт.
TM-2 умирает (OOM, k8s eviction). JM детектит по heartbeat timeout (10s), запускает restart strategy (fixed-delay, 3 попытки). Сначала аборт pending транзакций: producer.abortTransaction() + xa_rollback() — события e600-e650, накопленные после последнего checkpoint, никогда не станут visible downstream. Затем k8s спавнит новый pod, RocksDB-2 восстанавливается из S3 SST файлов (~минуты на 340GB), source TM-1 rewindит Kafka offset на checkpoint #142 (=99812), pipeline воспроизводит те же события — state-mutations идемпотентны на этом конкретном входе. Sink открывает новую транзакцию (txn.id=flink-sink-2-144), которая commit'нется на следующем чекпоинте. End-to-end exactly-once preserved.
Savepoint — это не checkpoint, это portable full snapshot для миграций. flink stop --savepointPath инжектит MAX_WATERMARK (закрывает все окна), берёт canonical full snapshot (не incremental — savepoint должен открываться на новой версии Flink), pickleит operator UIDs. Деплоим Flink 1.19 + fraud-job-v1.3.0.jar с новым оператором. flink run -s s3://.../savepoint-abcd1234/ — новый JM матчит операторы по uid(), новый оператор стартует с пустым state, старые — с восстановленным. Sources resume с сохранённых Kafka offsets. Downtime — минуты, exactly-once intact.
Flink vs Spark Structured Streaming vs Kafka Streams. Spark отброшен из-за micro-batch latency floor (200-500ms даже на aggressive trigger — критично для блокировки до charge confirmation) и heap-based state, который не масштабируется за ~50GB per executor без OOM. Kafka Streams отброшен из-за отсутствия first-class JDBC/Iceberg sinks, отсутствия concept savepoint (нужно вручную ресетить offset) и сложности enrichment join'ов с external state на TB-scale. Flink даёт true streaming + RocksDB + 2PC + savepoints + Flink SQL в одном runtime. Платим ops-сложностью (нужен cluster, JM HA, S3 для checkpoints) и меньшей Python-экосистемой по сравнению со Spark.
RocksDB backend + incremental checkpoints вместо HashMap + full snapshots. 1TB state физически не помещается в JVM heap ни одного разумного TM (>32GB heap → severe G1 pauses). RocksDB на NVMe SSD scales TB+, доступ ~1ms (приемлемо в бюджете 500ms p99), а incremental checkpoints копируют только новые/изменённые SST файлы — typical delta 5-10% of state = 50-100GB в S3 каждые 60s вместо 1TB. RocksDB managed memory автоматически делит heap между block cache, write buffer, indexes — Flink контролирует без OOM-риска. State TTL поддержан для unbounded key-space (card-id). External state в Redis/Cassandra — antipattern, теряет exactly-once guarantee.
TwoPhaseCommitSinkFunction + Kafka transactional producer для end-to-end exactly-once. Source-side тривиально (offset в state, rewind на recovery), но replay даёт duplicate emits в downstream — double-block legitimate user после retry. At-least-once + idempotent consumer требует unique event ID + dedupe state в каждом consumer (сложно, fraud service stateless). At-least-once + UPSERT работает для Postgres feature store, но не для Kafka downstream. 2PC sink делает pre-commit на barrier (transactional send, XA prepare), commit на notifyCheckpointComplete — атомарно с Flink checkpoint. На failure до commit — abortTransaction(), дубликаты невидимы. Цена: ~5-10ms latency per checkpoint cycle, transaction.timeout.ms > checkpoint_interval × 2, consumers должны быть isolation.level=read_committed.
parallelism=1 на bottleneck-операторе — single point hot, не масштабируется.AsyncDataStream.transaction.timeout.ms < checkpoint_interval — broker abortит живую транзакцию, sink в постоянном retry loop.