How do JoinWindows work in a KStream-KStream join, and what does the window size mean?
answer
- |tsA - tsB| <= diff (symmetric, width 2*diff)
- .before()/.after() = asymmetric
- Both sides buffered in window stores
- Left/outer non-matches emit after window close (no spurious)
- State ~ rate × (span + grace)
basics
~20 sJoinWindows define how close in event time two records from the two streams must be to join. With JoinWindows.ofTimeDifferenceAndGrace(d), records join if their timestamps differ by at most d (symmetric: a within [b-d, b+d]). Each stream is buffered in a window store for that span.
solid answer
~50 sA KStream-KStream join is windowed because two unbounded streams can't be joined unboundedly. `JoinWindows.ofTimeDifferenceAndGrace(diff, grace)` says: a record from stream A and a record from stream B join if `|tsA - tsB| <= diff` — a symmetric window of width 2*diff around each record. You can make it asymmetric with `.before(d1)` and `.after(d2)`, giving `tsB - d1 <= tsA <= tsB + d2`. Both sides are buffered in **window stores** so a record can match partners that arrive later within the window. Inner joins emit per matching pair; left/outer joins additionally emit null-padded results, and since KIP-633/improvements those non-matches are emitted when the window closes (after grace) rather than eagerly, so spurious left/outer results are avoided. Grace extends how long late records can still produce matches. State size is driven by diff + grace times the input rate on both sides.
go deeper
Know that stream-stream joins need a JoinWindows and join records close in time.
State the |tsA - tsB| <= diff rule and that both sides are buffered in window stores.
Explain before/after asymmetry, inner vs left/outer semantics, and the post-window-close emission of non-matches.
Size join state/changelog from rate × span × grace, reason about fan-out cardinality and producer clock skew, and choose window width against correctness/cost SLAs.
## Why a window at all? A **KStream-KStream join** correlates two record streams. Both are unbounded, so you can't keep every record from both forever waiting for a partner. The **JoinWindows** bounds *how far apart in event time* two records may be and still be considered "the same event" — and therefore how long each side must be buffered. ## The symmetric model `JoinWindows.ofTimeDifferenceAndGrace(Duration diff, Duration grace)` joins records A and B when: ``` |timestamp(A) - timestamp(B)| <= diff ``` This is a window of total width **2 * diff** centered on each record. If diff = 5 minutes, an A record joins any B record from 5 minutes before to 5 minutes after it (and vice versa). ## Asymmetric windows Use `.before(Duration)` and `.after(Duration)` to skew it: ``` timestamp(B) - before <= timestamp(A) <= timestamp(B) + after ``` Useful when one stream is expected to lead the other (e.g. an order then a payment that always comes *after*). ## How it executes — dual window stores Each side of the join is materialized in its own **window store** (a `WindowStore`, segment-based, see windowed-serde/segment questions). When a record arrives on stream A: 1. It is stored in A's window store. 2. Streams scans B's window store for all B records whose timestamps fall inside the join window relative to A, and emits a joined output per match. 3. The symmetric step happens when B records arrive too. Because both sides are stored, a match is produced regardless of which side arrives first, as long as both fall within the window before it expires. ## Inner vs left vs outer and the spurious-result fix - **Inner join**: emits only when a pair matches. - **Left join**: emits `(a, null)` if no B partner appears. - **Outer join**: emits `(a, null)` and `(null, b)` for unmatched records on either side. Historically left/outer joins emitted null-padded results **eagerly** (as soon as the first record arrived with no current match), producing **spurious** results that were later "corrected." Improvements around KIP-633 / the windowed-join rework made non-matches emit only **after the window closes** (i.e., after `diff + grace`), removing those spurious outputs. ## Grace and retention `grace` extends the time late records can still find partners. The underlying window stores must retain at least `windowSize + grace` (here windowSize relates to the join span); retention defaults derive from that, and you can't set retention below it. ## State and sizing Buffered state ≈ (input rate A + input rate B) × (window span + grace) × record size. Wide windows or long grace dramatically increase RocksDB state and changelog volume. This is the main scaling lever to reason about. ## Common edge cases - A record can join **multiple** partners (fan-out) within the window — joins are not 1:1. - Self-joins and high-rate streams with wide windows can explode output cardinality. - Time is event time via the TimestampExtractor; clock skew between the two producers directly affects which records match.
- How do you make a JoinWindow asymmetric, and when would you?Use .before(d1).after(d2). Use it when one stream reliably leads the other — e.g. an order always precedes its payment, so you allow more 'after' than 'before'.
- Why did older Kafka Streams versions emit spurious left/outer join results, and how was it fixed?They emitted null-padded results eagerly before the window closed, so a late-arriving partner made the earlier (a,null) result incorrect. The fix emits non-matches only after the join window (size + grace) closes.
saying these in an interview costs you the question
- Saying a KStream-KStream join is unwindowed or uses a single fixed bucket
- Claiming a record joins at most one partner (it can fan out to many)
- Thinking only one side is buffered (both sides have window stores)
- Believing left/outer joins still emit eagerly with spurious corrections in current versions