Windowing & Watermarks concept page. Streaming pipeline: Kafka source feeds a Flink job (Watermark Assigner -> keyBy -> Window Operator -> RocksDB state). Window emissions go to an aggregates sink, late events go to a side-output sink. Four scenarios: tumbling 1m window aggregation, late event with allowedLateness=2m re-fires window, session window with gap=15min for user activity, plus an ADR contrasting event-time + watermark vs processing-time simplicity.
Потоки бесконечны, а агрегации требуют ограниченных кусков: «продажи за последний час», «сессия пользователя», «top-N кликов за минуту». Windowing — это примитив, который разрезает infinite stream на bounded windows. Без него streaming-движок умеет только map/filter — для GROUP BY нужно решить, где заканчивается «час».
Главный вопрос — какое время считаем. Wall-clock на воркере (processing time) даёт простой код, но неправильную бизнес-семантику: rerun на тех же данных даст другие числа, потому что aggregates зависят от того, когда события дошли до кластера. Event time (timestamp в payload) — правильное, но события приходят поздно, не по порядку, и иногда теряются.
Чтобы корректно обрабатывать out-of-order поток, нужны watermarks — heuristic «всё, события до момента T уже арривнули, можно закрывать window и эмитить результат». Watermark — это компромисс: слишком оптимистичный → много late events drop'нуто и метрики дырявые; слишком консервативный → окна закрываются с задержкой и latency растёт.
На этом срезе различаются production-ready streaming engines (Flink, Beam) от toy-проектов: first-class watermarks с per-partition tracking, allowed lateness, side outputs, idle source detection. Spark Structured Streaming умеет watermarks, но ограниченно. Storm и старый Spark — практически нет.
Windowing = разрезание infinite stream на bounded chunks. Event time = когда событие реально произошло. Watermark = граница «всё до T уже здесь». Late event = пришёл после watermark — drop, accept с grace или side output. Trigger = когда эмитить результат (на close, рано, или повторно для late events).
Три времени, которые нельзя путать:
payload.timestamp). Канонический для бизнеса.Watermark W(t) = max_seen_event_time - max_lateness (bounded out-of-orderness, самый ходовой вариант). Window [start, end) закрывается, когда W >= end.
event_time в payload.boundedOutOfOrderness=30s) — следит за max_seen_event_time, выдаёт W = max − 30s.Tumbling 1m, eventTime, allowedLateness=2m. Назначает событие в окно [start, end), держит per-window state.Edges рисуют физическую топологию data plane; ответы в анимации идут reverse по тем же edges — отдельных edge для emit/cleanup в обратку нет.
tumbling — happy path. Tumbling 1m, события в порядке. Watermark двигается с каждым новым max_seen_event_time, и когда W >= window_end — окно фaйрится, эмитится агрегат, state очищается. Это базовый цикл: assign → append → wait for W → fire → emit → cleanup.
late-event — allowedLateness. Окно уже закрыто и initial-эмит сделан, но allowedLateness=2m держит state живым ещё 2 минуты. Опоздавшее событие (event_time < W, но в пределах grace) принимается, state обновляется, downstream получает повторный emit для того же window key. Когда W уходит за end + lateness, state окончательно дропается; всё, что пришло после — летит в side output. Это ключевой паттерн: downstream должен быть idempotent, иначе counts задвоятся.
session — gap-based window. Session window закрывается после gap=15min без событий для key. Каждое новое событие продлевает окно до event_time + gap. Watermark, который ушёл за конец session, триггерит fire и эмит «session_duration, event_count». Новое событие после gap стартует новую сессию. В отличие от tumbling, границы окон разные для каждого key.
adr-event-vs-processing — выбор time semantics. Option A (processing time): простой код, ноль out-of-order логики, минимальная latency — но non-deterministic, replay невоспроизводим, mobile/network gap кладёт события в чужой бакет. Option B (event time + watermark + allowedLateness): детерминизм, корректная бизнес-семантика, чистый replay — ценой watermark heuristic, idempotent sink и обязательного idle detection. Правило: если rerun должен давать идентичные аггрегаты — только event time, исключений нет. Для ops-дашбордов processing time допустим.
| Решение | Pros | Cons | Когда |
|---|---|---|---|
| Processing time | Простой код, нулевая ceremony, минимальная latency | Non-deterministic, replay даёт другие числа, mobile gap → события в чужом окне | Ops-дашборды, non-critical alerting, internal metrics |
| Event time + watermark | Детерминизм, корректные бизнес-аггрегаты, чистый replay, late handling | Heuristic watermark, нужны idle detection и idempotent sink | Биллинг, аналитика, любые числа в отчётах |
| Tumbling window | Простой state (one bucket per window), каждое событие в одном окне | Грубая дискретизация, нет «скользящих» метрик | «Sales every 5 min», billing periods |
| Sliding window | Гладкая trailing-метрика, обновление каждую минуту | State explosion: size/slide буферов per key | «Trailing 10-min avg every 1 min» — только если slide ≥ size/10 |
| Session window | Естественные пользовательские сессии, gap определяет границу | Unbounded growth при активном user, нужен max-length cap | User engagement, trip detection, fraud sessions |
| allowedLateness=0 | Простой downstream (single emit per window), низкий state | Late events дропаются, метрики дырявые | Когда поток заведомо ordered (CDC из одной БД) |
| allowedLateness=N min | Late events учитываются, повторные emit с уточнением | Downstream обязан быть idempotent / upsert, state живёт дольше | Mobile-клиенты, события через flaky network |
| Discarding mode | Каждый emit = delta, легко агрегировать downstream | Сложно реконструировать «полное» окно | Stream pipelines с downstream-aggregation |
| Accumulating mode | Каждый emit = running total | Дубли, если downstream считает за дельту | Materialized views, dashboards |
| Accumulating + retracting | Корректные updates: retraction + new | Сложнее sink (нужен MERGE / upsert) | Materialize, Flink Table API upsert, kSQL |
Decision rule: «event time + watermark + allowedLateness + idempotent sink» — default для всего, что попадает в отчёты. Processing time оставляем для observability, где «приблизительно» допустимо.
size=1h, slide=1s × 1M keys = state explosion. Используйте hierarchical aggregation.withIdleness(1min) — обязательная строка.