How do windowed aggregations work in Kafka Streams? Cover the window types, the resulting key type, and how late records and grace periods are handled.
answer
- windowedBy() between group and aggregate
- tumbling / hopping / sliding / session
- result key = Windowed<K> (key + start/end)
- grace period: late-within-grace folds, past-grace drops
- retention >= size + grace; WindowStore/SessionStore
basics
~20 sYou call windowedBy(...) on a KGroupedStream before count/reduce/aggregate, so records are bucketed by time. The result is a KTable keyed by Windowed<K> (key + window bounds). Window types include tumbling, hopping, sliding, and session windows. Late records that arrive within the grace period still update their window; beyond grace they are dropped.
solid answer
~50 sWindowing partitions a key's records into time buckets so each bucket aggregates independently. You insert windowedBy() between groupByKey()/groupBy() and the aggregator. Time windows (TimeWindows) cover tumbling (size only, non-overlapping) and hopping (size + advance, overlapping). SlidingWindows produce a window per distinct pair of in-range records. SessionWindows group activity separated by an inactivity gap and can merge windows as records fill the gap. The output KTable is keyed by Windowed<K>, which carries the original key plus window start/end, and is materialized in a WindowStore (or SessionStore) with bounded retention. Event-time semantics use a TimestampExtractor. Late records are folded into their window as long as they arrive within the window's grace period (Duration ofSizeAndGrace / ofSizeWithNoGrace, or windows with grace). After grace expires the window is closed: late records are dropped and counted in the late-record sensor. Retention must be >= window size + grace or old windows are purged early.
go deeper
Know windowing buckets records by time before aggregating.
Name tumbling/hopping/session and that the key becomes Windowed<K>.
Explain grace period, stream time, dropped late records, and retention >= size + grace.
Choose window types per use case, reason about overlap/merging semantics and retention/storage trade-offs at scale.
## Why windows An unwindowed aggregation grows forever — a count per key over all time. Often you want 'per key, per time bucket': clicks per user per minute. **Windowing** buckets a key's records by time so each bucket aggregates separately. ## Where it goes in the topology `stream.groupByKey().windowedBy(window).count()`. `windowedBy()` turns a `KGroupedStream` into a `TimeWindowedKStream` (or `SessionWindowedKStream`), and the subsequent count/reduce/aggregate runs **per (key, window)**. ## Window types - **Tumbling** (`TimeWindows.ofSizeWithNoGrace(Duration)`): fixed size, non-overlapping, gap-free. Each record belongs to exactly one window. - **Hopping** (`TimeWindows.of(...).advanceBy(Duration)`): fixed size with an advance smaller than the size, so windows **overlap** and a record can fall into multiple windows. - **Sliding** (`SlidingWindows.ofTimeDifferenceAndGrace(...)`): a window is created around the timeframe between records that are within the time difference — efficient for 'within N seconds of each other' aggregates. - **Session** (`SessionWindows.ofInactivityGapAndGrace(...)`): data-driven, variable-length windows defined by a period of inactivity (the gap). New records can **merge** two adjacent sessions, which emits a tombstone for the merged-away session key plus the new merged value. ## Result key: Windowed<K> The aggregate KTable is keyed by **`Windowed<K>`** = original key + `Window` (start/end timestamps). When you sink it you typically use `WindowedSerdes` or extract `windowed.key()` / `windowed.window().start()`. ## Time and lateness Windows use **event time** by default, derived from a `TimestampExtractor` (default reads the record timestamp). Streams tracks **stream time** = the max observed timestamp. A record is **late** if its timestamp is older than the stream time minus window boundaries. - **Grace period**: each window has a grace `Duration`. Late records within grace are still folded into their window (updating the aggregate). Past grace the window is **closed**; later records for it are **dropped** and surfaced via the `dropped-records` / late-record metrics. - Modern API forces an explicit choice: `ofSizeAndGrace(size, grace)` or `ofSizeWithNoGrace(size)` (grace 0). ## Retention The backing `WindowStore`/`SessionStore` keeps windows for a **retention** period (default 1 day, set via `Materialized.withRetention(...)`). Retention **must be >= window size + grace**, else Streams throws because windows would be deleted before they can legally accept late data. Expired windows are compacted away in the changelog. ## Edge cases - A hopping window's overlap means a single record contributes to multiple aggregates — counts across windows are not additive. - Without `suppress()`, each windowed aggregate still emits intermediate updates as records arrive; suppress(untilWindowCloses) holds emission until window+grace passes. - Session-window merges produce tombstones; downstream must handle null values. - Stream time only advances with new records; with no traffic, a window may stay 'open' indefinitely and never close (relevant for suppress).
- What is the relationship between window size, grace period, and store retention?Retention must be >= window size + grace period. The store has to keep a window alive for its whole duration plus the grace window during which late records may still arrive; otherwise Streams rejects the configuration and would lose data.
- How does a session window differ from a tumbling window, and what extra event can it produce?A session window is data-driven and variable-length, defined by an inactivity gap rather than a fixed clock boundary. A new record bridging two sessions merges them, emitting a tombstone for the removed session key plus the merged aggregate.
- What happens to a record whose event timestamp falls before stream-time minus the grace period?Its window is already closed, so the record is dropped from the aggregation and counted by the dropped-records/late-record metric. It does not update any aggregate.
saying these in an interview costs you the question
- Saying the windowed result is keyed by the plain K instead of Windowed<K>.
- Claiming late records always update their window regardless of grace.
- Setting retention shorter than window size + grace.
- Confusing hopping (overlapping, fixed) with sliding or session windows.
- Assuming windowing uses processing time by default rather than event time via the TimestampExtractor.