0%
Data-Intensive Applications
Foundations of Data Systems
Distributed Data
Encoding and Evolution
Batch Processing
Data Quality and Governance
Operational Patterns
Stream Processing Patterns
Stream processing systems face a fundamental problem: streams are unbounded, but most useful computations require boundaries. You cannot compute "average order value" over an infinite stream. You need to group events into finite chunks and compute over those chunks. Windows are how you draw those boundaries.
Before the strategies, the number that decides whether any of them is affordable. A sliding window keeps every open pane in state simultaneously, so the live entry count is
for keys, window length , and slide . A million keys on a 60 second window sliding every 10 seconds is 6 million live entries, which is comfortable. The same million keys on a one hour window sliding every minute is 60 open panes and 60 million entries, which is not. Shortening the slide from 60 s to 10 s to make a dashboard smoother multiplies state by six, and that is usually how a stream job discovers its memory limit: not from a traffic increase but from a reporting request.
Every windowing strategy answers the same two questions: which events belong to this window, and when should the window emit its result? Different strategies answer these questions differently, and the choice has direct consequences for correctness, latency, and resource consumption. The three fundamental window types are tumbling (fixed, non-overlapping), sliding (fixed, overlapping), and session (dynamic, data-driven).
Tumbling Windows
Tumbling windows are fixed-size, non-overlapping time intervals. A 5-minute tumbling window groups all events between 10:00-10:05 into one window, 10:05-10:10 into the next, and so on. Every event belongs to exactly one window. No gaps, no overlaps. The window boundaries are aligned to the epoch (or a configurable offset), so all processor instances agree on where windows start and end without coordination.
This is the simplest windowing strategy and the right default for periodic aggregations. "Count the clicks per minute." "Sum revenue per hour." "Find the max temperature per day." Each event lands in one window, each window fires one result.
The tradeoff is boundary sensitivity. An event at 10:04:59 and an event at 10:05:01 are two seconds apart but land in different windows. If you are counting login failures to detect brute force attacks, a burst of 10 failures split across a window boundary looks like two benign windows of 5 instead of one alarming burst of 10.
Tumbling window size selection involves a latency-accuracy tradeoff. Smaller windows (1 minute) produce results faster but are more susceptible to variance: a single outlier event heavily skews a 1-minute average. Larger windows (1 hour) smooth out variance but delay results. The right size depends on your use case: a real-time dashboard showing trends might use 1-minute windows for responsiveness, while a billing system computing daily totals uses 24-hour windows for accuracy.
Sliding (Hopping) Windows
Sliding windows solve the boundary problem by overlapping. A 10-minute window that slides every 2 minutes produces overlapping windows: 10:00-10:10, 10:02-10:12, 10:04-10:14, and so on. Each event belongs to multiple windows. An event at 10:05 appears in five windows simultaneously.
This is the right choice when you need smooth aggregations that do not miss patterns at boundaries. Moving averages, rolling sums, and trend detection all benefit from sliding windows. The cost is computation: each event triggers updates to multiple windows. A 10-minute window with a 1-minute slide means each event updates 10 windows.
The amplification factor (window size divided by slide interval) is the key metric for capacity planning. A 1-hour window with a 1-minute slide means each event updates 60 windows. At 10,000 events per second, that is 600,000 window updates per second. The state backend must handle this write amplification, and every checkpoint must persist all active windows. Teams often start with aggressive slide intervals for precision, discover the cost in production, and widen the slide to reduce load.
When an interviewer asks about real-time anomaly detection or rate limiting, reach for sliding windows. Tumbling windows create blind spots at boundaries that attackers or anomalies can exploit. Sliding windows give you a continuous view of recent activity, which is what you actually need for detecting patterns that do not respect clock boundaries.
Session Windows
Session windows are the only window type driven by data, not by clocks. A session window groups events that arrive close together and closes after a gap of inactivity. Set the gap to 30 minutes, and a user who clicks at 10:00, 10:05, 10:12, and 10:45 produces two sessions: one from 10:00-10:12 and another starting at 10:45.
Session windows are essential for user behavior analysis. "How long was the user active?" "How many pages did they view per session?" "What was their checkout funnel path?" These questions do not have fixed time boundaries because user sessions vary in length.
Choosing the right gap duration requires understanding your domain. For a web application, 30 minutes is the industry standard (this is what Google Analytics uses). For a mobile game, 5 minutes might be appropriate because users switch between apps frequently. For IoT sensors, the gap might represent device sleep cycles (minutes to hours). Setting the gap too short splits single sessions into many, inflating session counts and deflating duration metrics. Setting it too long merges distinct sessions into one, hiding return visits.
The implementation challenge is that sessions can merge. If a late event arrives at 10:08 with an allowed lateness of 30 minutes, it might bridge two sessions that were previously separate, forcing the system to merge them and recompute the aggregate. This makes session windows the most expensive to maintain in stateful stream processing.
Choosing the Right Window
The choice between window types is driven by the question you are answering:
- "How many X per fixed period?" Use tumbling windows. Dashboard counts, hourly summaries, daily aggregates. Simple, predictable, low cost.
- "Is X trending up or down right now?" Use sliding windows. The overlap smooths out noise and eliminates boundary artifacts. Accept the amplification cost.
- "How does one user's session look?" Use session windows. The dynamic boundaries match actual behavior patterns. Accept the merge complexity.
A common mistake is reaching for the most sophisticated window type by default. Session windows are powerful but expensive. If your question has fixed time boundaries ("revenue per hour"), use a tumbling window. The simplest correct window type is always the best choice because it minimizes state, reduces checkpoint size, and simplifies debugging.
Count-Based vs Time-Based Windows
All the windows described above are time-based: they group events by timestamps. Count-based windows group events by quantity: "every 100 events" or "every 1,000 messages." Count-based tumbling windows emit a result after exactly N events, regardless of how long it takes to accumulate them.
Count-based windows are useful for batch-oriented operations within a stream. A database write buffer that flushes every 500 records. A machine learning feature pipeline that retrains after 10,000 new training examples. A log aggregator that batches 200 log lines before sending to storage.
The key difference from time-based windows is that count-based windows have variable latency: if events arrive slowly, you might wait minutes for a window to fill. If events arrive quickly, windows fill in milliseconds. This makes them unsuitable for latency-sensitive aggregations but ideal for throughput-sensitive batch operations.
In practice, many systems combine both approaches: use a count-based trigger with a time-based fallback. "Flush after 500 records OR 10 seconds, whichever comes first." This ensures that high-throughput streams batch efficiently while low-throughput streams do not stall indefinitely waiting for the count threshold. Flink's trigger API supports these compound conditions natively, allowing you to build custom triggering logic for any window type.