Stream processing concept page (/concepts/stream-processing). Source -> operators -> sink pipeline. Stateless operators (filter, map) vs stateful (keyBy, tumbling window aggregate, RocksDB local state). Time semantics: event-time vs processing-time. Watermarks and late events with side output. Backpressure via credit-based flow control. Checkpoint-restore via Chandy-Lamport. Includes 5 scenarios (simple ETL, stateful windowed count, late event handling, backpressure cascade, checkpoint recovery) and 2 ADRs (stream-vs-batch-vs-lambda, event-time-vs-processing-time).
Batch processing (Hadoop, Spark batch, nightly Airflow DAGs) хорош для отчётов «вчера»: aggregate логи за день к 9 утра, посчитать MAU за месяц к началу следующего. Но бизнес 2026-го всё чаще хочет «сейчас»: alert на fraud transaction через секунду после авторизации, ETA для Uber-driver через 200 ms, recommendation tile, обновлённый в момент клика, dashboard со свежими CTR с лагом 5 секунд. Это требует другой парадигмы — данные приходят непрерывно, обработка идёт непрерывно, результаты доступны через миллисекунды-секунды.
Второй мотив — state. Современный stream processing — это не «трансформация без памяти», а stateful computation: running totals, joins двух потоков по окну, sessionization, deduplication, windowed aggregations. Flink держит терабайты state на одном job (RocksDB-backed). Это превращает streaming engine из «трубы» в distributed database с push-обновлениями: те же данные, что лежат в БД, но материализация идёт continuously, и query == subscription.
Третий — унификация batch и streaming. До 2015 у тебя был Hadoop для batch и Storm для streaming — две кодовые базы, две модели данных, два сета операторов. Beam (Dataflow), Flink, Spark Structured Streaming унифицировали API: код один, данные могут быть bounded (batch) или unbounded (stream). Это снижает operational и cognitive overhead в разы — а batch получается как частный случай streaming с финитным водяным знаком.
Без понимания stream processing невозможно проектировать ::concept{slug="exactly-once-semantics"}, ::concept{slug="windowing-watermarks"}, real-time ML inference pipelines, материализованные views поверх ::concept{slug="kafka-deep"}, и весь класс продуктов «live analytics» (ClickHouse + Kafka + Flink стандарт-стек 2026-го).
Stream processing — это continuous computation на потенциально-бесконечных потоках событий, с управляемым state, корректной обработкой времени (event-time + watermarks) и фолт-толерантностью через distributed snapshots (checkpoints).
Ключевые отличия от batch одной таблицей:
| Batch | Streaming | |
|---|---|---|
| Input | bounded | unbounded |
| Latency | минуты-часы | ms-секунды |
| Failure model | rerun job | checkpoint + replay |
| State | derived при каждом run | persistent across events |
| Time semantics | snapshot | event-time + watermarks |
| Output | one-shot | continuous (append/update/retract) |
В DataFlow Model (Akidau, Google, VLDB 2015) batch — частный случай streaming: bounded stream = stream с конечным водяным знаком, который рано или поздно достигнет +∞.
На канвасе — типичный Flink-style DAG job, разрезанный на три логических участка.
Слева — Kafka source topic clicks на 8 партиций (kafka-in). Это unbounded input: события льются непрерывно, partition-параллелизм даёт horizontal scaling source-операторов.
В центре — stream processing job, разделён на две подгруппы:
filter-op → map-op) — операторы без памяти. Filter отсеивает ботов по User-Agent, map парсит JSON и обогащает по IP. Параллелизм тривиален: каждый task processит свою партицию независимо, нет shuffle, операторы chain в один thread.keyby-op → window-op + rocksdb) — операторы со state. keyBy(userId) — это network shuffle: события репартиционируются по hash(key) так, что весь трафик одного userId оседает на одном worker. Window agregator копит running counts в локальном RocksDB, эмитит при закрытии окна (когда watermark > end_of_window).rocksdb — embedded LSM-tree, держит state на NVMe-диске worker'а. Это позволяет state size расти до терабайтов на job без upload-overhead в БД на каждый event.
Справа — S3 для checkpoints (асинхронные snapshots state для recovery), Kafka sink topic agg-clicks (агрегированный output, пишется в Kafka транзакции для exactly-once), и отдельный late-events side output (события, опоздавшие после watermark, не дропаются молча, а попадают в специальный stream для special handling + alerting).
ADR-001 живёт на filter-op («когда stream, когда batch, когда lambda»), ADR-002 — на window-op («event-time vs processing-time для бизнес-агрегатов»). Это два решения, которые архитектор стрим-системы принимает раньше любого кода.
1. Simple ETL (simple-etl) — stateless pipeline без shuffle. Source → filter (отсеваем bot UA, ~30% drop) → map (parse JSON, geo-обогащение по in-memory lookup table) → sink. Операторы chain'атся в один thread, нет сетевых hops, latency ~10 ms p99 на 500K events/s с одного 4-core TM. Учит: «не каждый pipeline нуждается в Flink-уровне сложности; stateless ETL — это просто и быстро, и можно гонять на Kafka Streams в JVM-микросервисе без отдельного кластера».
2. Stateful windowed count (stateful-windowed-count) — каноничный «count clicks per user per minute». keyBy(userId) шафлит поток по сети (~1 Gbps между TM при 100K events/s), tumbling 1-min window накапливает count в RocksDB. RocksDB lookup на каждый event — sub-millisecond за счёт block cache. Когда watermark пересекает end_of_window, оператор эмитит результат как (userId, minute, count) в выходной Kafka transaction. 1M unique users × ~100 байт state = 100 MB на TM — комфортно для RocksDB. Учит: «stateful streaming = distributed DB с push-обновлениями; state локален, shuffle дорогой, размер state определяет parallelism».
3. Late event handling (late-event-handling) — production-grade pitfall. Mobile-юзер был в метро, его click с event_time=15:42 приехал в 15:48 (mobile reconnect). Watermark системы уже прошёл 15:43, window 15:42-15:43 закрыт и эмитнут. Что делать с поздним событием? Drop (default Flink) — silent data loss. Allowed lateness +30s — окно ещё доступно. Side output — late event попадает в отдельный stream с alert на counter. Учит: «event-time без late-handling = data loss; для mobile/IoT нужно bounded out-of-orderness 5 минут + allowed lateness; never drop silently — всегда side output + alert».
4. Backpressure (backpressure) — каскадное замедление от медленного sink. Kafka cluster притормаживает (network blip, p99 = 800 ms на ack). window-op output buffer заполняется, credits для upstream обнуляются. keyby-op ждёт credits от window-op, тоже стопает. filter-op ждёт keyby-op. Весь pipeline стоит, source перестаёт читать Kafka, consumer lag растёт +10K msg/s. Учит: «backpressure — это feature, а не bug: credit-based flow control предохраняет от OOM. Лечится sink parallelism + batch size, НЕ bump'ом checkpoint interval (это маскировка симптома)».
5. Checkpoint + recovery (checkpoint-recovery) — distributed snapshot (Chandy-Lamport) + exactly-once recovery. JM инжектит barrier id=42, операторы синхронно snapshot'ят state в S3 (8 GB на 10 Gbps link ≈ 30 секунд async upload). TM с window operator падает OOM на skewed key (1 user = 80% events). Local RocksDB на упавшем TM недоступен. New TM качает snapshot с S3, kafka-source seek'ает к offset из checkpoint 42, replay начинается. Kafka transactions для replayed событий aborted, clean re-emit — no duplicates. Учит: «EOS требует координации трёх вещей: source replayable (Kafka offsets), state snapshotable (RocksDB → S3), sink transactional (Kafka txn 2PC); если хоть одно не выполнено — at-least-once максимум».
Context. Команды путают «нужно real-time» с «нужно streaming». Реально 80% данных можно гонять nightly batch (Spark, Trino over Iceberg, Airflow по cron), это в разы дешевле и проще, чем держать Flink cluster 24/7 с 3 JobManager'ами в HA и регулярными RocksDB rebalancing'ами. Streaming оправдан там, где TTV (time-to-value) события измеряется секундами: fraud detection (block transaction до auth), real-time pricing (Uber surge), live dashboards (CTR за минуту), feature engineering для real-time ML inference.
Decision. Дерево решений: первое — какая задержка нужна от события до результата? Часы/дни → batch (Spark + Iceberg + Airflow), кончили. Секунды/минуты → streaming. Второе — нужен ли historical reprocessing (бэкфилл новой ML-фичи на год данных)? Если да → Kappa: храни raw events в Kafka 30+ дней (или в Iceberg) и replay через тот же Flink job. Третье — есть ли legacy batch pipeline, который нельзя выкинуть? Тогда Lambda: stream layer для fresh, batch layer для accurate, merge в serving layer (но это дорого в operations). Четвёртое — сложность state? Stateless map/filter → Kafka Streams в JVM-микросервисе. Большой keyed state (>100 GB), сложные windows, CEP → Flink. SQL-only people с materialized views → Materialize/RisingWave/ksqlDB.
Trade-off.
Когда batch достаточно: finance reports, BI dashboards с лагом часа, ML training, compliance reports. Не плати за streaming, если можешь обойтись без.
Когда streaming необходим: fraud, real-time pricing, live analytics, real-time ML inference, ::concept{slug="event-sourcing"} с CQRS read models, ::concept{slug="cdc"} в downstream replicas.
Context. Соблазн писать window по processing-time — проще, watermarks не нужны, late events не существуют «по определению». Но любая mobile аппа отправляет события батчами при reconnect (метро, самолёт, плохой 4G), GPS-точки приходят с задержкой 30-300 секунд из-за power-saving, IoT-датчики копят и шлют раз в N минут. Если посчитать «sales за час 10:00-11
» по processing-time, при rerun на бэкфилле получишь другие числа: события, приехавшие в 11 в первый раз попали в окно 11:00-12, а при бэкфилле — в правильное 10:00-11.Decision. Использовать event-time + watermarks для всего, что бизнес читает или сравнивает между runs. Watermark стратегия: bounded out-of-orderness 5-30 секунд для backend events, 1-5 минут для mobile/IoT. Allowed lateness ещё +30s для grace, поздние события — в side output с alert на counter (не drop silently). Processing-time только для операционных метрик (CPU per worker, GC pauses), debug, low-importance dashboards с явной пометкой «processing-time».
Trade-off.
Context. Flink state можно держать в JVM heap (FsStateBackend) или в embedded RocksDB на disk. Heap быстрее (нет serialization), но ограничен JVM memory (~30 GB разумного предела, дальше GC pauses становятся неприемлемыми). RocksDB медленнее (~10× на access), но позволяет state size в терабайты на TM за счёт NVMe-диска.
Decision. Default — RocksDB для любого production job с state > 1 GB или unbounded ростом (windowed aggregations по pop user_id, session state, dedupe stores). Heap — только для тестов, прототипов, или сильно ограниченных по latency hot paths с small state (< 100 MB).
Trade-off.
Context. Чем чаще checkpoints, тем меньше re-processing после сбоя, но тем больше overhead (RocksDB snapshot, S3 upload, координация barrier'ов). 10-секундный interval даёт <10s data re-process, но требует 8 GB upload каждые 10s = ~6 Gbps на TM. 60-секундный interval даёт более скромный overhead, но при сбое replay 60 секунд = ~30M events заново.
Decision. 10-60 секунд в зависимости от state size и SLA. State < 10 GB → 10 секунд OK. State 100+ GB → 30-60 секунд, иначе incremental checkpoint upload насытит диск/сеть. Никогда < 5 секунд — overhead на координацию barrier'ов начинает доминировать.
Trade-off. Совсем не делать checkpoints — RPO = от создания job (потеряешь всё state при сбое). Делать каждую секунду — половина CPU/IO уйдёт на snapshots. Sweet spot — 30 секунд для большинства production jobs.
Apache Flink на Uber (Marmaray + AthenaX) — десятки тысяч Flink jobs для real-time matching ETA, dynamic pricing, fraud detection. Marmaray — internal data pipeline framework на Flink, AthenaX — SQL-интерфейс поверх. Skewed-key проблема (горячий driver-ID) решается салтингом + pre-aggregation.
Netflix Keystone — Flink-based real-time stream processing platform, обрабатывает ~1 трлн events/день. Real-time alerting на video playback quality, ML feature engineering для personalization. До 2018 был Mantis (Reactive Streams), мигрировали на Flink за state management и EOS.
Alibaba Blink → Apache Flink — Alibaba форкнули Flink (Blink), доработали для своих масштабов, потом смерджили обратно в upstream. На Singles' Day 2020 пиково обрабатывали 4 млрд events/sec на Flink-jobs для real-time inventory, recommendations, fraud.
Stripe Sigma + Radar — Spark Structured Streaming для fraud scoring. Каждая транзакция scored в real-time через ML model, decision в <100 ms. Micro-batch trigger ~100 ms — компромисс latency vs throughput.
Pinterest — ksqlDB + Kafka Streams для feature engineering и ad targeting. Embedded в JVM микросервисах, нет отдельного Flink-кластера. Идеально для команд, которые не хотят operational overhead Flink.
DoorDash — Flink для дельта-вычислений ETA. Обработка GPS-событий курьеров (миллионы per second), windowed aggregations, ::concept{slug="cep"}-style pattern matching на «late delivery alert».
Cloudflare Workers + Flink — analytics pipeline для DDoS detection. Cloudflare Logpush → Kafka → Flink → ClickHouse, latency от события до dashboard ~5 секунд.
Materialize / RisingWave — streaming SQL с incremental view maintenance. CREATE MATERIALIZED VIEW continuously поддерживается актуальным, query == subscription. Альтернатива Flink для команд, которые хотят SQL вместо Java DSL.
processing-time для бизнес-семантики. Sales aggregate по processing-time даст разные числа при rerun, аналитики потеряют доверие к данным. Всегда event-time + watermarks для всего, что бизнес читает.keyBy(orderId) где orderId уникальный. Каждый orderId — свой state slot, миллион orderId = миллион state entries без шанса GC. Решение: агрегировать по более широкому ключу (user_id, region, category) или явно cleanup'ить state через TTL.key + random(1..N)), потом merge на втором этапе.processing.guarantee=exactly_once_v2 при EOS требованиях. Default — at-least-once, дубликаты возможны и они проявятся на проде через неделю под нагрузкой.back_pressure per operator. Если HIGH — bottleneck. Лечится parallelism, batch size, sink optimization. Bump checkpoint interval — это маскировка, не fix.KeyedState API явно.CREATE MATERIALIZED VIEW.