How are windowed aggregation results keyed and stored — explain Windowed keys, windowed serdes, and the segmented store layout?
answer
- windowedBy → KTable<Windowed<K>, V>
- WindowedSerdes: time needs size; session stores start+end
- WindowKeySchema: key ‖ start ‖ seqnum
- Segmented store: expire = drop whole segment (O(1))
- retention >= size + grace; longer for IQ
basics
~20 sA windowed aggregation produces a KTable keyed by Windowed<K> (the original key plus the window's start/end). It is stored in a segmented WindowStore (RocksDB split into time segments) and serialized with a WindowedSerdes that encodes key + window timestamp. Old segments are dropped wholesale when they fall outside retention.
solid answer
~50 sWindowing turns `KGroupedStream<K,V>` into `KTable<Windowed<K>, V>`. `Windowed<K>` wraps the raw key with the `Window` (start/end). The state lives in a **WindowStore**, which RocksDB-internally is **segmented**: time is divided into segments (segment interval ≈ retention / number of segments, min 2), and each window's data goes into the segment covering its start time. The storage key is a composite: serialized record key + window-start timestamp + a sequence number, produced by `WindowKeySchema`. For output and Interactive Queries you use `WindowedSerdes.timeWindowedSerdeFrom(...)` / `sessionWindowedSerdeFrom(...)`, and on the topic the windowed key is serialized with the window size so it can be reconstructed. The big operational win of segmentation: **expiring old data is dropping whole segments**, an O(1) file deletion, instead of scanning and deleting individual keys. Retention determines how many segments are kept; grace must be <= retention. This layout also makes range scans by time efficient.
go deeper
Know windowed aggregations are keyed by the key plus the window (Windowed<K>).
Know windowed serdes exist and that time-windowed deserializers need the window size.
Explain the WindowStore is segmented and retention must be >= size + grace, longer for IQ.
Explain segment intervals, WindowKeySchema composite keys, O(1) segment-drop expiry, changelog/restore implications, and tuning segment count vs retention for disk/IQ trade-offs.
## From grouped stream to windowed table A non-windowed aggregation yields `KTable<K, V>`. Adding `.windowedBy(TimeWindows...)` changes the result type to `KTable<Windowed<K>, V>`. The key is now a **`Windowed<K>`** object: the original key **plus** the `Window` (which carries `start()` and `end()` event-time bounds). So the same logical key K appears once per window it participated in. ## Windowed serdes Because the key is now compound, it needs special serialization: - **`WindowedSerdes.timeWindowedSerdeFrom(innerClass, windowSize)`** (and the `TimeWindowedSerializer/Deserializer`) encodes the inner key bytes followed by the window-start timestamp. The **window size** must be supplied on the deserializer side so the window `end` can be reconstructed (the serialized form stores start, and size lets you derive end). - **`WindowedSerdes.sessionWindowedSerdeFrom(innerClass)`** for session windows encodes both start and end (sessions are variable length, so end can't be derived from a fixed size). A classic bug: deserializing a time-windowed topic with the wrong window size yields wrong window ends. Since around KIP-659, helpers were added to ease constructing these serdes correctly. ## Segmented store layout The `WindowStore` is backed by RocksDB but **partitioned into time segments** to make retention cheap: - Time is divided into **segments**. The **segment interval** ≈ `retentionPeriod / (numSegments - 1)`, with a minimum of 2 segments and a 60s floor on the interval. - A window is placed in the segment whose time range contains the window's **start** timestamp. - The **physical key** within a segment is built by `WindowKeySchema`: `serialized-key ‖ windowStartTimestamp (8 bytes) ‖ sequenceNumber (4 bytes)`. The sequence number disambiguates records that share key+timestamp. - For session stores, `SessionKeySchema` encodes `key ‖ end ‖ start` so range queries by end time are efficient. ## Why segmentation matters (the principal-level point) With windowed data you must **expire** old windows once they fall outside **retention** (= at least size + grace, often longer for Interactive Queries). If each expired window were deleted individually, you'd pay a scan + per-key delete cost and generate RocksDB tombstones/compaction churn. Instead, when stream time advances past a segment's whole range, Streams **drops the entire segment** (deletes the underlying RocksDB column family / files) in essentially O(1). Segmentation thus turns retention enforcement into bulk file deletion. ## Retention, grace, and IQ - **Constraint:** `grace <= retention`, and retention must be `>= windowSize` (+ grace). The builder validates this. - Retention can be set **longer** than grace via `Materialized.withRetention(...)` when **Interactive Queries** need to serve historical windows after they've stopped updating. - Reducing the number of segments lowers open-file overhead but makes retention coarser (you keep up to one extra segment of stale data). ## Changelog The window store's changelog topic is a compacted+deleted topic; record keys are the composite windowed keys, and segment expiry corresponds to tombstones/retention on the changelog so restoration rebuilds only live segments. ## Edge cases - Interactive Queries (`ReadOnlyWindowStore.fetch(key, from, to)`) translate the time range into a scan across the relevant segments. - Choosing too small a retention silently drops windows you might still want for IQ; too large inflates disk and restore time. - The window-size mismatch on the windowed deserializer is a silent correctness bug, not an error.
- Why must you pass the window size to a TimeWindowedDeserializer?The serialized time-windowed key stores only the window start; the deserializer reconstructs window end = start + size. A wrong size silently produces wrong window ends without throwing.
- What is the operational advantage of a segmented window store over a flat key-value store for windowed data?Retention enforcement becomes bulk segment deletion: when stream time passes a segment's range, the whole segment's RocksDB files are dropped in O(1), avoiding per-key scans, tombstones, and compaction churn.
- How do session-windowed serdes differ from time-windowed ones?Time-windowed serdes store only the start and derive end from a fixed window size; session serdes must store both start and end because session windows are variable length and end cannot be derived.
saying these in an interview costs you the question
- Saying the windowed key stores both start and end for time windows (it stores start; end derives from size)
- Claiming expired windows are deleted key-by-key (they're dropped per-segment)
- Thinking retention can be shorter than grace
- Forgetting that a wrong deserializer window size is a silent correctness bug