Command Palette

Search for a command to run...

Hectal
PHASE 11Advanced ~8 min· topic 3 of 5

Topic 11.3

Windowing: Tumbling, Hopping, Sliding, Session

In one line

Windows group events by time for aggregation. Tumbling windows are fixed and non-overlapping, hopping windows are fixed and overlapping, sliding windows are defined by the time difference between records, and session windows are separated by inactivity gaps. Each answers a different business question.

0/5 · 0%

Think of it like this

Counting cars on a road. Every 5 minutes on the dot (tumbling). A 5-minute window recalculated every minute (hopping). Any 5-minute stretch in which cars came close together (sliding). Each traffic burst separated by at least 2 quiet minutes (session).

Key ideas

  1. 01

    Tumbling: TimeWindows.ofSizeWithNoGrace(Duration.ofMinutes(5)): 10:00–10:05, 10:05–10:10; each event belongs to exactly one window. Use for periodic reports.

  2. 02

    Hopping: TimeWindows.ofSizeAndGrace(5m, grace).advanceBy(1m): overlapping windows; each event belongs to several windows. Use for moving averages.

  3. 03

    Sliding: SlidingWindows.ofTimeDifferenceWithNoGrace(5m): windows defined relative to records' timestamps, capturing "within 5 minutes of each other" patterns (fraud: 3 purchases within 5 minutes).

  4. 04

    Session: SessionWindows.ofInactivityGapWithNoGrace(30m): a session continues while events keep arriving within the gap; windows merge; variable length. Use for user sessions.

  5. 05

    Output: windowed aggregations emit updates as data arrives; suppress(Suppressed.untilWindowCloses(...)) emits one final result per window after it closes (plus grace), trading latency for a single final answer.

Code & diagrams

PurchaseCounts.javajava
KStream<String, Purchase> purchases = b.stream("purchases");   // key = userId

// Count purchases per user every 5 minutes, allow 2 minutes of late data, emit once per window
purchases.groupByKey()
    .windowedBy(TimeWindows.ofSizeAndGrace(Duration.ofMinutes(5), Duration.ofMinutes(2)))
    .count(Materialized.as("purchases-5m"))
    .suppress(Suppressed.untilWindowCloses(Suppressed.BufferConfig.unbounded()))
    .toStream((windowedKey, count) -> windowedKey.key() + "@" + windowedKey.window().startTime())
    .to("purchases.per-user.5m", Produced.with(Serdes.String(), Serdes.Long()));
windows.txttext
events (minutes): 1   2        6    7                  40  41
tumbling 5m:     [0-5): 2    [5-10): 2               [40-45): 2
hopping 5m/1m:   each event counted in 5 overlapping windows
sliding 5m:      one window per group of events within 5 min of each other (inclusive): {1,2,6}, {6,7}, {40,41}
session gap 10m: session A = 1..7 (4 events), session B = 40..41 (2 events)

Interview problem

The problem

Count purchases per user every 5 minutes

"Count purchases per user every 5 minutes." Compare tumbling, hopping, sliding and session windows and explain which semantics each provides for this requirement.

Explain it without notes

01

Which window type would you use for fraud detection of rapid purchases, and why?

Practice

01

Use TopologyTestDriver to feed timestamped events and assert tumbling window counts, advancing stream time to close windows.

Trade-offs

  • ↔

    Overlapping and session windows answer richer questions at the cost of more state and output volume.

Done when you can

  • I can choose and implement tumbling, hopping, sliding and session windows.