Stream Windowing & Watermarks Lab (Interactive)
Tune watermark delay on an out-of-order stream and watch accuracy trade against staleness. Run 48 out-of-order events through tumbling, sliding, or session windows; a watermark heuristic decides which stragglers still count and when each result fires.
Watermarks & Windowing under Out-of-Order Data
A 10-minute stream of 48 events (some mobile stragglers) flows through your engine — tune the watermark delay and watch accuracy trade against result latency.
Event completeness
85.4%
Late events dropped
7 of 48
Avg window fire lag
53.5s after end
Observed max disorder
123s
watermark = max(eventTime seen) − 30s. Balanced: most stragglers still make their window; a few slide-memberships silently undercount. Stragglers that missed every window: 7.
▸ record-by-record, <10ms — window fires the instant the watermark passes (Stripe-style velocity checks).
A watermark asserts "no further events with event-time < T will arrive." Set it too aggressively and out-of-order stragglers are diverted to side output, quietly corrupting the aggregate; set it too generously and every result is late. Production engines pair watermarks with allowed-lateness side outputs and RocksDB-backed state checkpoints so a corrected straggler can re-fire an already-emitted window.
How It Works Under the Hood
Event-time tells you when something happened; processing-time tells when the server saw it, and mobile networks make them disagree. Stream engines therefore close windows on a watermark assertion—"no further events earlier than T will arrive"—balancing completeness against result latency. A short delay fires fast but diverts stragglers to side output, silently undercounting; a generous delay captures everything but serves stale numbers. Engine choice rides on top: Flink and Kafka Streams evaluate record-by-record, Spark micro-batches add half a second per result.
Core Architectural Principles
- watermark = max(eventTime observed) − allowed-lateness; windows close the moment it passes their end.
- Events arriving after their window closed are dropped or corrected—the completeness versus freshness dial, computed per window.
- Session windows extend on activity and fire only gap + watermark after the last event; sliding windows double-count into adjacent spans.
Use Stripe’s velocity-check shape: sliding 60-second windows for fraud features under 20ms via Flink. Define event-time versus processing-time, explain watermarks as a probabilistic assertion, and quantify your chosen lateness allowance against observed disorder—that sentence separates seniors from students.
Aggressive watermarks deliver fresh but lossy aggregates; conservative ones deliver complete but stale results—side outputs bridge the gap.