skip to content

How do session windows work, including session merging and the inactivity gap?

level: seniorimportance: should knowfreq 50%

answer

  1. Inactivity gap defines session boundaries
  2. Variable length, per key, [minTs, maxTs]
  3. Out-of-order → sessions MERGE
  4. aggregate needs a merger function
  5. SessionStore keyed by (key,start,end)

basics

~20 s

A session window groups records for a key that arrive within an inactivity gap of each other. The window grows with each new record and closes when no record arrives for longer than the gap. Out-of-order records can merge two adjacent sessions into one.

solid answer

~40 s

Session windows (`SessionWindows.ofInactivityGapAndGrace(gap, grace)`) are **data-driven, variable-length** windows scoped per key. A session is a contiguous run of records where consecutive timestamps differ by no more than the inactivity `gap`. Each new in-window record extends the session's end. If a record arrives whose timestamp falls within `gap` of two existing sessions (common with out-of-order data), those sessions **merge** into one spanning session — Streams removes the old session keys and writes the merged session. Sessions are stored in a `SessionStore`, keyed by (key, sessionStart, sessionEnd). Aggregation uses `aggregate(initializer, aggregator, merger, ...)` — the extra **merger** function combines the aggregate values of two sessions being merged. Grace governs how long late records can still extend or merge sessions before the result is final; retention must cover gap + grace. Typical use: clickstream / user-activity sessionization.

go deeper

for a junior

Know a session groups activity separated by idle gaps and closes after inactivity.

for a middle

Know the inactivity gap defines boundaries, sessions are per key and variable length, stored in a SessionStore.

for a senior

Explain merging on out-of-order data and why aggregate() needs a merger; tie grace/retention to gap.

for a principal

Reason about store churn/changelog cost from frequent merges under skewed/out-of-order load and design sessionization SLAs with suppress + grace.

## Concept A **session window** captures a burst of activity. Instead of fixed clock boundaries, a session is defined by **periods of activity separated by periods of inactivity**. This matches real user behavior: a user clicks several times, goes idle, comes back later — each burst is a session. `SessionWindows.ofInactivityGapAndGrace(Duration inactivityGap, Duration grace)`. ## The inactivity gap Two records (for the same key) belong to the same session if the time between them is `<= inactivityGap`. If the gap between a new record and the nearest existing session boundary exceeds `inactivityGap`, a **new** session starts. So a session's `[start, end]` is `[min timestamp, max timestamp]` of the records it contains, and its effective duration is variable. ## Growth When an in-gap record arrives, the session's `end` (and possibly `start`, for out-of-order earlier records) is extended to include it, and the aggregate is updated. The old session entry (with its previous boundaries) is deleted and a new one with the new boundaries is written — sessions are keyed by their boundaries, so changing boundaries changes the key. ## Merging — the defining feature With **out-of-order** data, a record can land **between** two existing sessions and be within `inactivityGap` of both. Those two sessions must then **merge** into a single session covering both. Kafka Streams: 1. Finds all sessions within `gap` of the new record. 2. Removes them from the store. 3. Creates one merged session whose boundaries span all of them and whose aggregate is the combination of all their values. This is why session aggregation requires a **merger** in addition to the aggregator: ``` .aggregate(initializer, aggregator, merger, Materialized...) ``` The `merger: (key, aggLeft, aggRight) -> aggMerged` tells Streams how to fold two sessions' partial aggregates together (e.g. sum the counts). ## Storage Sessions live in a **SessionStore** (a specialized RocksDB store, segment-based like window stores) keyed by `(key, sessionStartTime, sessionEndTime)`. Because boundaries are part of the key, merges manifest as delete-old + put-new operations, reflected in the changelog topic. ## Grace and retention `grace` defines how long after a session's end late records may still extend or merge it. Once `streamTime > sessionEnd + grace`, the session is final; later records for it are dropped. The store retention must be at least `inactivityGap + grace`. ## suppress for final-only output Like other windows, intermediate session updates stream out as a changelog. Pair with `suppress(Suppressed.untilWindowCloses(...))` to emit one final record per session after it closes. ## Edge cases - A single isolated record is a valid session of zero duration (`start == end`). - Heavy out-of-order data causes frequent merges → extra store churn and changelog traffic. - Different keys have completely independent session boundaries.

  • Why does session-window aggregation require a merger function in addition to the aggregator?
    Because out-of-order records can merge two existing sessions; the merger combines their two partial aggregate values into the single aggregate for the merged session. Other window types never merge, so they don't need it.
  • Can a session window have zero duration?
    Yes. A single record with no neighbors within the inactivity gap forms a session where start == end (zero-length), which is perfectly valid.

saying these in an interview costs you the question

  • Describing session windows as fixed-size
  • Forgetting the merger function in aggregate()
  • Claiming sessions never merge / ignoring out-of-order merging
  • Saying sessions are global rather than per key

context