Windowing, event time and watermarks
Aggregating streams correctly: event time versus processing time, tumbling, sliding and session windows, late and out-of-order events, watermarks and allowed lateness, triggers and updates, state size, and stream joins.
Reading is half of it. See this used in a real interview: walk through Design an Ad Click Aggregator →
"Count clicks per ad per minute" sounds trivial until you notice that events arrive late, out of order, and sometimes hours after they happened (a phone that was offline). Which minute does a late click belong to, and when is a minute's count final? Stream processors answer with windows, event time and watermarks. These concepts are what interviewers probe in aggregation questions like the ad click aggregator.
Event time versus processing time
- Event time: when the event actually happened, recorded by the device or service that produced it.
- Processing time: when the stream processor sees it.
Aggregating by processing time is simple but wrong whenever there is delay: a burst of delayed events after an outage would be counted in the wrong minute, and replaying a day of data would put everything in the replay minute. For correct results, aggregate by event time, which requires handling disorder.
Window types
| Window | Shape | Example |
|---|---|---|
| Tumbling | fixed size, non-overlapping | clicks per ad per minute |
| Sliding (hopping) | fixed size, overlapping, advancing by a step | requests in the last 5 minutes, updated every 30 seconds |
| Session | per key, closed after a gap of inactivity | user sessions ending after 30 idle minutes |
| Global with triggers | one window, emitted periodically | running totals |
Sliding windows multiply work and state (each event belongs to several windows); compute them from smaller tumbling windows where possible.
Watermarks
A watermark is the processor's estimate that "all events with event time up to T have probably arrived". When the watermark passes the end of a window, the window can be emitted as complete.
- Watermarks are usually computed as the maximum event time seen minus an allowed delay (say 30 seconds), or from source-level progress.
- A watermark that is too aggressive emits results early and treats many events as late.
- A watermark that is too conservative delays every result.
The choice is a trade-off between latency and completeness.
Late events
Events arriving after the watermark passed their window:
- Allowed lateness: keep the window's state for an extra period (minutes or hours) and update the result when late events arrive, emitting a correction.
- Side output: route very late events to a separate stream for batch reconciliation.
- Drop: acceptable for some metrics, never for billing.
Downstream stores must accept updates (upsert by window key), not just appends. For billing-grade numbers, a later batch job over the full raw data produces the final answer. See counting at scale and Design an Ad Click Aggregator.
Triggers
Waiting for the watermark can be too slow for dashboards. Triggers emit early, speculative results (every 10 seconds while the window is open), then the on-time result at the watermark, then late updates. Users see fresh numbers that settle into accurate ones.
State
Every open window holds state per key: counts, sums, sketches, or buffered events. State size is roughly keys × open windows × state per window, so:
- Use compact aggregates (counters, HyperLogLog) rather than storing raw events. See probabilistic data structures.
- Bound allowed lateness; longer lateness means more open windows.
- Store state in a local embedded store (RocksDB in Flink) with periodic checkpoints to durable storage for recovery and exactly-once results. See delivery semantics.
Stream joins
- Stream to table: enrich each click with ad metadata from a changelog-backed table (kept in the processor's state).
- Stream to stream within a window: join impressions with clicks that happen within 30 minutes; both sides are buffered for the window, so state grows with the window length.
- Interval joins for event pairs like ride requested and ride accepted. See Design Uber.
Duplicates and keys
Retries produce duplicate events; deduplicate by event id within a time-bounded window before counting. Partition the stream by the aggregation key so each window's state lives on one worker, and watch for hot keys. See hot keys and skew.
In the interview
For aggregation questions: "We aggregate in one-minute tumbling windows by event time, with a watermark 30 seconds behind the latest event time, early triggers every 10 seconds for dashboards, allowed lateness of one hour with updates upserted by window key, very late events sent to a side output, and a daily batch over raw events for billing." That covers almost every follow-up. See Design a Monitoring System and Design a Leaderboard.
Checklist
- Event time for correctness; processing time only where approximations are fine.
- Window type matched to the question; sliding built from tumbling.
- Watermarks tuned for the latency versus completeness trade-off.
- Allowed lateness with updates; side outputs and batch reconciliation.
- Early triggers for fresh dashboards.
- Compact state, checkpoints, deduplication and key partitioning.