skip to content

How do JoinWindows work for stream-stream joins, including grace period and inner vs outer/left emission timing?

level: seniorimportance: should knowfreq 50%

answer

  1. |tsA - tsB| <= timeDifference
  2. before()/after() = asymmetric window
  3. grace replaced until()/retention (KIP-633)
  4. inner emits on match; left/outer null after window close
  5. spurious-result fix delays null rows

basics

~20 s

A JoinWindows defines a time difference: record A joins record B if their timestamps are within the window. The grace period allows late records before the window closes. Inner results emit immediately on a match; left/outer null-paired results emit only after the window plus grace expires.

solid answer

~50 s

Stream-stream joins use `JoinWindows.ofTimeDifferenceAndGrace(timeDifference, grace)`. Two records join if |tsA - tsB| <= timeDifference (you can make it asymmetric with `before`/`after`). Both sides are buffered in window stores keyed by record key + window, and joins are evaluated on event time. The grace period (KIP-633 replaced the deprecated `until`/`retention`) defines how long after a window ends late-arriving records are still accepted; records later than grace are dropped. For inner joins, a result is emitted as soon as a matching pair is found. For left and outer joins, a spurious-result fix (KIP-633 era) means the null-paired result (A with null, or null with B) is only emitted once the window has closed past the grace period and no match arrived—avoiding the old behavior of emitting false negatives that were later corrected. Window store retention must be at least window size + grace.

go deeper

for a junior

Know a window means records must be close together in time to join.

for a middle

Use ofTimeDifferenceAndGrace, know event-time semantics and that late records are dropped.

for a senior

Explain asymmetric windows, the grace period, and delayed left/outer emission (spurious-result fix).

for a principal

Reason about correctness-vs-latency tradeoffs of grace, store retention sizing, and state cost across high-throughput joins.

## What a JoinWindows is In a KStream-KStream join there is no 'current value', so you bound matches in **event time**. `JoinWindows.ofTimeDifferenceAndGrace(Duration timeDifference, Duration grace)` says: a record `a` (timestamp `tA`) joins a record `b` (timestamp `tB`) when ``` tB ∈ [tA - timeDifference, tA + timeDifference] ``` You can make it asymmetric: `JoinWindows.ofTimeDifferenceWithNoGrace(d).before(x).after(y)` lets you say 'b must be no earlier than x before a and no later than y after a'—useful for causal joins (e.g. an ad-click must follow an impression). ## State buffering Each side is stored in a **window store** (RocksDB-backed, with a changelog topic for fault tolerance). When a record arrives on side A, Streams scans side B's window store for records in the matching time range, and vice-versa. This is why both sides retain data for the window duration. ## Event time, not wall-clock Window membership is decided by **record timestamps** (the `TimestampExtractor`), not processing time. Stream time advances as records are observed. A record whose timestamp is older than the current stream time by more than (window + grace) is **late** and dropped. ## Grace period (KIP-633) Older APIs used `JoinWindows.of(d).until(retention)`; these are deprecated. The modern API forces an explicit **grace**: the extra time a window stays open after its end to admit out-of-order/late records. Larger grace = more correctness for late data, at the cost of holding state longer and delaying null-paired emission. ## Emission timing: inner vs left/outer - **Inner:** emit `(a, b)` the moment a matching pair is found. Straightforward. - **Left:** for each A, emit `(a, null)` if no B matched—but *when*? Naively emitting on arrival produces **spurious results**: you'd emit `(a, null)`, then later a matching B arrives and you emit `(a, b)`, leaving a wrong `(a, null)` downstream. KIP-633 fixed this: the null-paired result is emitted only **after the window closes (end + grace)** and confirms no match ever arrived. This delays left/outer null rows but eliminates false negatives. - **Outer:** same logic applied symmetrically—`(a, null)` and `(null, b)` emitted only after their windows close with no partner. ## Edge cases - **Retention:** the window store retention must be >= window size + grace; otherwise records are evicted before the window logically closes. - **Duplicate emission:** if multiple B records fall in A's window, you get one output per pair (a fan-out), not a single combined row. - **Out-of-order within grace:** still joined correctly; beyond grace, silently dropped (monitored via the dropped-records metric / `late-record-drop` sensors). - **Self-join / same topic:** allowed; be careful about a record matching itself depending on keys.

  • Why were left/outer null-paired results historically 'spurious', and how was it fixed?
    Older Streams emitted (a,null) immediately, then if a match arrived later emitted (a,b)—leaving a wrong null row downstream. KIP-633 delays the null result until the window fully closes past grace, so it's only emitted when no match can ever arrive.
  • What must the window store retention be relative to the JoinWindows?
    At least the window size plus the grace period; otherwise records are evicted before the window logically closes and valid late joins are missed.

saying these in an interview costs you the question

  • Saying window membership is by processing/wall-clock time (it's event time via the TimestampExtractor).
  • Claiming left-join null results emit immediately when no match is found (they wait until window close + grace).
  • Confusing grace with retention—grace admits late records; retention is the physical store TTL and must cover window+grace.

context