Windowing — Tumbling, Hopping, Session & Global Windows with Triggers
A window turns an infinite stream into finite chunks you can aggregate. The window type is a product decision, not a tuning knob: tumbling windows answer per minute, hopping windows answer over the last five minutes updated every minute, session windows answer per visit, and a global window answers ever, until you say so. Triggers decide when a window emits, and the accumulation mode decides what each emission means downstream. Getting these wrong produces numbers that look plausible and are silently wrong.
- 1Gist
- 2Maps
- 3Q&A
- 4Sandbox
Voice readout needs Web Speech Synthesis in this browser.
Question ladder
L1
Tumbling versus hopping. When do you need hopping?
Answer
When the question is over the last N minutes and you need it refreshed more often than N. Tumbling gives one result per period and hides bursts that straddle a boundary.
L2
How are session windows implemented if events arrive out of order?
Answer
Each event creates a proto-window from t to t plus the gap. The engine merges overlapping windows for the same key and merges their state, so a late event can bridge two sessions. That needs mergeable accumulators and, in SQL, retractions.
L3
What is a trigger versus a watermark?
Answer
The watermark is the engine's estimate of event-time completeness. A trigger is the policy that decides when to emit a pane, often when the watermark passes the window end, plus early and late firings.
L4
Accumulating versus discarding. Which for a key-value sink?
Answer
Accumulating with an upsert keyed by the key and the window start. Each pane overwrites the previous value. Discarding fits sinks that add deltas and breaks if anyone treats a pane as the full value.
L5
How would you compute a 24 hour rolling unique-user count updated every minute?
Answer
Do not hop a one-minute slide over raw events. Pre-aggregate per-minute HyperLogLog sketches in tumbling windows, then merge the last 1440 sketches. Sketches are mergeable.
L6
Why can a global window be useful?
Answer
For count-based or punctuation-based logic, such as flush every 100 events or when an order-closed record arrives, where time boundaries are meaningless. You must supply the trigger and a cleanup policy.
L7
Why are timezone-naive daily windows wrong?
Answer
Tumbling windows align to epoch UTC unless you set an offset. A day for America/Chicago needs that offset and a plan for daylight saving time, or Monday's bucket swallows Sunday evening.
Failure modes
Tiny hop on a wide window
A one-hour window sliding every second creates 3600 windows per event.
Sessions on bots
A key that never goes idle never closes, and state grows without a max duration.
Global window and no trigger
The window is assigned and never emits.
Sum of accumulating panes
Early 15 plus on-time 22 plus late 25 is 62, not 25, unless the sink overwrites or retracts.
Misconceptions
Sliding, hopping, and Spark's window with a slide are three ideas.
Flink sliding windows, Kafka Streams hopping windows, and Spark window(col, size, slide) are the same overlapping assigner. Kafka Streams also has a separate SlidingWindows type.
The watermark is the trigger.
The watermark is a completeness estimate. The trigger is the policy that fires a pane, often because of the watermark.
A pane is the final answer.
Early and late panes are refinements. The accumulation mode says whether they replace, add, or retract.
Interviewer traps
Recommending a one-minute slide over raw events for a daily unique count.
Tumble into mergeable sketches, then combine the last 1440 buckets. Point at the sketch lessons for the merge.
Calling Kafka Streams' default emit a final result.
It emits on every update. suppress until the window closes when downstream wants one row.
Design scenario
Same prompt for every reader.
Requirements
Three assigners, one trigger policy, and a sink that does not sum accumulating panes. Bots must not hold a session open all day.
Traffic / scale
Tens of thousands of events per second, with a few hot keys.
Latency
Speculative dashboard numbers within 30 seconds. Final panes when the watermark passes the window.
Consistency
A replay produces the same final totals. Early panes may be overwritten.
Availability
A late event inside allowed lateness updates the pane. Anything later is a side output, not a silent drop with no metric.
Failure assumptions
- Events for one user arrive out of order.
- A bot key never pauses.
Constraints
- Do not hop raw events for a 24 hour unique count.
- Daily windows for Chicago are not UTC midnight unless someone asked for UTC.
Prompt
Product wants per-minute order totals, a five-minute rate refreshed every minute, and session counts that close after ten seconds of silence. Dashboards may show an early number. Finance must not double-count.
API
What is the upsert key for an accumulating pane?
Data
Which assigner owns per minute, last five minutes, and a visit?
Architecture
Where does a retraction go if a downstream job sums by region?
What a second firing means
Prefer
The sink knows the accumulation mode
Discarding panes are deltas to add. Accumulating panes overwrite one row. Retractions let a downstream sum stay equal to the latest total.
- Upsert key is the business key plus the window start.
- A retraction subtracts the previous pane before adding the new one.
- Session merges need mergeable state, and SQL needs the retraction.
Alternative
Every pane is appended and summed
Early 15, on-time 22, and late 25 become 62. The dashboard looks precise and the number is wrong.
- Accumulating mode double-counts if anyone treats panes as additive.
- Discarding mode breaks if anyone treats a pane as the full value.
- A global window with no trigger emits nothing, forever.
From one event to a pane
The assigner is the product decision. The trigger is when. The mode is what the number means.
- 1
Event at time t for key k
The timestamp is event time. Processing time is a different window. - 2
Assigner stamps windows
Tumbling lands in one bucket. Hopping lands in size divided by slide buckets. A session opens or merges. Global is the only window. - 3
Trigger fires a pane
Watermark past the end, a processing-time timer, a count, or a punctuation record. - 4
Accumulation mode shapes output
Discard the delta, emit the total so far, or retract the previous total and emit the new one.
Window types
| Window | Definition | Event belongs to | State cost | Best for | Watch out |
|---|---|---|---|---|---|
| Tumbling | Fixed size, no overlap | Exactly one window | Low: one accumulator per key per window | Per-minute metrics, billing buckets | A burst split across two windows hides a spike |
| Hopping | Fixed size, advances by a slide smaller than the size | size/slide windows | size/slide times tumbling | Moving averages, last-five-minutes alerts | Write amplification. A tiny slide is a huge cost |
| Session | Per key, closes after a gap of inactivity. Windows merge | One window that grows and merges | Unpredictable. A never-idle key never closes | User visits, device bursts, conversations | Bots never go idle. Merging needs mergeable state |
| Global | One window for all time per key | The only window | Grows forever without eviction | Custom triggers: count, punctuation | Never fires without a custom trigger |
| Count | Every N elements per key | Implemented as global plus a count trigger | Low | Micro-aggregation, batching to sinks | Not time-meaningful. Slow keys never flush |
Terminology trap: Flink's sliding window is Kafka Streams' hopping window is Spark's window(col, "10 minutes", "5 minutes"). Kafka Streams also has a separate sliding window, SlidingWindows, for windowed aggregations. Joins use JoinWindows, defined by a max time difference between records. Partitions and keys underneath this stay on Apache Kafka.
Decisions
- 1
1. Event at time t, key k
- next2. Which window?
- ?
2. Which window?
- one bucket3a. Tumbling
- overlap3b. Hopping
- gap3c. Session or global
- 3
3a. Tumbling
- next4. Trigger emits a pane
- 4
3b. Hopping
- next4. Trigger emits a pane
- 5
3c. Session or global
- next4. Trigger emits a pane
- 6
4. Trigger emits a pane
- next5. Discard, accumulate, retract
- 7
5. Discard, accumulate, retract
Lesson map
Windowing — Tumbling, Hopping, Session & Global Windows with Triggers
A window turns an infinite stream into finite chunks you can aggregate. The window type is a product decision, not a tuning knob: tumbling windows answer per minute, hopping windows answer over the last five minutes updated every minute, session windows answer per visit, and a global window answers ever, until you say so. Triggers decide when a window emits, and the accumulation mode decides what each emission means downstream. Getting these wrong produces numbers that look plausible and are silently wrong.
Architecture. Architecture
Select a node to see why it exists, or an edge to see the protocol, direction, effect, and consequence.
Mermaid export
flowchart TB rec["1. Record at time t, key k"] kinds["2. Tumbling, hopping, session, global"] trig["3. Trigger decides when to emit"] mode["4. Accumulate, discard, or retract"] rec -->|1. Record at time t, key k| kinds kinds -->|2. Tumbling, hopping, session, global| trig trig -->|3. Trigger decides when to emit| mode
Assigning events
The Python is a mini version of Flink's window assigners, including session merging. Each event opens a proto-window from t to t plus the gap. Overlapping ones for the same key merge. Hopping doubles the per-event work at a slide of half the size.
Press Run. Snippets must be self-contained — no network, files, or native modules.
tumbling(10s): {(0, 10): 3, (10, 20): 2, (30, 40): 1, (40, 50): 1}
hopping(10s/5s): {(-5, 5): 2, (0, 10): 3, (5, 15): 2, (10, 20): 2, (15, 25): 1, (30, 40): 1, (35, 45): 2, (40, 50): 1}
session(gap=10s): {'u1': [(1, 24, 'n=3'), (39, 51, 'n=2')], 'u2': [(4, 14, 'n=1'), (16, 26, 'n=1')]}
global (one window, needs a custom trigger to ever fire): total = 7
hopping amplification: 2.0 window-updates per eventu1 produced two sessions because the 25 second gap between t=14 and t=39 exceeded the 10 second gap. u2's events at 4 and 16 are 12 seconds apart, so they are two one-event sessions. With a 15 second gap they would merge. Hopping wrote 2.0 window-updates per event. A one-hour window sliding every second would write 3600.
Triggers and accumulation
A trigger fires a window pane. A typical combo from Streaming 102: early firings on a processing-time timer for speculative dashboards, an on-time firing when the watermark passes the window end, and late firings for records that arrive within allowed lateness. Allowed lateness itself is the watermark lesson.
| Trigger | Fires when | Use |
|---|---|---|
| Event-time | Watermark passes the window end | Default for event-time windows |
| Processing-time | Wall-clock timer | Early or speculative results, and processing-time windows |
| Count | N elements in the pane | Global and count windows, batching |
| Purging or custom | Your logic, such as a punctuation record | End-of-session markers, an order-closed event |
| Accumulation mode | Each pane emits | Downstream must | Cost |
|---|---|---|---|
| Discarding | Only the delta since the last pane | Sum the deltas | Cheapest |
| Accumulating | The full value so far | Overwrite by key | Keeps state until the window is garbage-collected |
| Accumulating and retracting | Retract the old value, emit the new one | Apply minus-old and plus-new | Most expensive. This is what Flink SQL changelogs do |
Press Run. Snippets must be self-contained — no network, files, or native modules.
accumulating | EARLY: emit 15 ; ON-TIME: emit 22 ; LATE: emit 25
discarding | EARLY: emit 15 ; ON-TIME: emit 7 ; LATE: emit 3
accumulating+retracting | EARLY: emit 15 ; ON-TIME: retract 15 ; ON-TIME: emit 22 ; LATE: retract 22 ; LATE: emit 25If a downstream job sums window results per region, accumulating mode double-counts: 15 + 22 + 25 = 62. Discarding is fine only if every consumer sums. Retractions keep the sum of all emitted values at 25 no matter how downstream regroups. Flink SQL's update-before and update-after rows are this mechanism.
Engine mapping
| Concept | Flink DataStream / SQL | Spark Structured Streaming | Kafka Streams |
|---|---|---|---|
| Tumbling | TumblingEventTimeWindows, SQL TUMBLE | window(ts, "1 minute") | TimeWindows.ofSizeWithNoGrace |
| Hopping | SlidingEventTimeWindows, SQL HOP | window(ts, "10 min", "5 min") | advanceBy |
| Session | EventTimeSessionWindows, SQL SESSION | session_window | SessionWindows.ofInactivityGapWithNoGrace |
| Cumulative | SQL CUMULATE, daily totals updated every hour | Manual via state | Manual via state |
| Early firings | A custom trigger, or SQL early-fire config | Output mode update plus the trigger interval | No suppress means emit on every update |
A 24 hour rolling unique-user count updated every minute should not be a hopping window over raw events. Pre-aggregate per-minute HyperLogLog sketches in tumbling windows, then merge the last 1440 sketches. Mergeability is the point of approximate aggregations.
Pitfalls
In the session function, set the gap to 15 seconds and predict whether u2's events at 4 and 16 merge before you run it. Then explain which accumulation mode a region-level sum needs if those sessions can still merge late.
Interview Q&A
Tumbling versus hopping. When do you need hopping?
Answer
When the question is over the last N minutes and you need it refreshed more often than N, such as rolling averages and rate alerts. Tumbling gives one result per period and has boundary blind spots.
How are session windows implemented if events arrive out of order?
Answer
Each event creates a proto-window from t to t plus the gap. The engine merges overlapping windows for the same key and merges their state, so a late event can bridge two sessions into one. That requires mergeable accumulators and, in SQL, retractions for the replaced sessions.
What is a trigger versus a watermark?
Answer
The watermark is the engine's estimate of event-time completeness. A trigger is the policy that decides when to emit a pane, often when the watermark passes the window end, plus early and late firings.
Accumulating versus discarding. Which for a key-value sink?
Answer
Accumulating with an upsert sink keyed by the key and the window. Each pane overwrites the previous value. Discarding fits sinks that add deltas and breaks if anyone treats a pane as the full value.
How would you compute a 24 hour rolling unique-user count updated every minute?
Answer
Avoid a one-minute-slide hopping window over raw events. Pre-aggregate per-minute HyperLogLog sketches in tumbling windows, then merge the last 1440 sketches per output. Sketches are mergeable.
Why can a global window be useful?
Answer
For count-based or punctuation-based logic, flush every 100 events or when an order-closed record arrives, where time boundaries are meaningless. You must supply the trigger and an eviction policy.
Why are timezone-naive daily windows wrong?
Answer
Tumbling windows align to epoch UTC unless you set an offset. A civil day for America/Chicago needs that offset, and daylight saving time still moves the boundary. Without it, evening traffic lands in the next UTC day.