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 выполняет stateful stream graph и восстанавливается из consistent snapshots. Главная граница: exactly-once state effects внутри job не превращает Kafka, PostgreSQL и внешний API в одну transaction.
| Компонент | Ответственность |
|---|---|
| Replayable Source | Предоставляет partition offsets и replay after restore. |
| Checkpoint Coordinator | Запускает numbered barriers и принимает successful acknowledgements. |
| Stateful Operator A | Обрабатывает keyed events и snapshots local/operator state. |
| Stateful Operator B | Aligns multiple inputs or captures in-flight buffers for unaligned mode. |
| Durable Checkpoint Store | Хранит recoverable snapshot metadata/files. |
| Transactional Kafka Sink | Может commit output with checkpoint through its connector contract. |
| Idempotent Database Sink | Применяет stable event IDs/versions; не атомарен с Kafka sink автоматически. |
| Job Recovery and Savepoints | Restarts from checkpoint or controlled savepoint with stable operator IDs. |
Barriers divide pre/post-checkpoint events. A multi-input operator pauses already-barriered channels until the others catch up, then snapshots state.
Проверяемый исход: Checkpoint completion is a recovery point, not a claim that business output waited one checkpoint interval to become visible.
After a task failure, Flink restores the latest completed checkpoint and rewinds replayable sources. Some records execute again, but managed state reflects each logical input once under exactly-once mode.
Проверяемый исход: Duplicate execution is expected; sink contract determines whether duplicate external effects appear.
Kafka output can participate via connector transactions; a database sink may be idempotent or transactional under a separate contract. Writing both does not create one cross-system atomic commit.
Проверяемый исход: Every sink reports its own checkpoint/idempotency proof; reconciliation handles partial cross-sink success.
Slow downstream work delays aligned barriers. Unaligned checkpoints can capture in-flight data and shorten checkpoint completion, but do not remove processing latency. Savepoints with stable operator IDs support controlled upgrades.
Проверяемый исход: Operations distinguish checkpoint health, record latency and savepoint compatibility.
Числа выше — учебные inputs или размерностные формулы. Их нельзя выдавать за benchmark или SLA конкретного продукта.
[CONCEPT]exactly-once-semantics
Диаграмма показывает причинные границы и recovery contracts, а не скрытую реализацию конкретного managed-сервиса. Любая stronger guarantee действует только в явно названной transaction/checkpoint/acknowledgement boundary.
Введите числа или выберите пресет