SysDesignPrep.com
Study guide 100 of 183

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

WindowShapeExample
Tumblingfixed size, non-overlappingclicks per ad per minute
Sliding (hopping)fixed size, overlapping, advancing by a steprequests in the last 5 minutes, updated every 30 seconds
Sessionper key, closed after a gap of inactivityuser sessions ending after 30 idle minutes
Global with triggersone window, emitted periodicallyrunning 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.

Open in your browser to sign in

Google does not allow sign-in inside this app's built-in browser. Open this page in Safari and sign in there. The link opens this same page.

Tap the ⋯ or share button at the top or bottom of the screen, then Open in browser. Or copy the link and paste it into Safari.