skip to content

Joins

Stream-stream, stream-table, table-table and foreign-key joins, and the co-partitioning they demand. Interviewers ask because join semantics plus co-partitioning is where most Streams designs go wrong.

part ofApache Kafkaoverview, primer and where to startread it →
on this pageshow

questions

5

What kinds of joins does Kafka Streams support, and how do they differ at a high level?

level: juniorimportance: must knowfreq 70%

answer

  1. 3 families: stream-stream, stream-table, table-table
  2. stream-stream needs a window
  3. stream-table = lookup current value
  4. GlobalKTable = no co-partitioning
  5. FK join relaxes same-key

basics

~20 s

Kafka Streams supports stream-stream joins (over a time window), stream-table joins (enrich a record with the latest table value), and table-table joins (combine two changelog tables). GlobalKTable joins let a stream join a fully-replicated table without matching partitioning.

solid answer

~40 s

Kafka Streams offers three join families plus a global variant. Stream-stream joins (KStream-KStream) match records from two event streams that fall within a time window (JoinWindows), since both sides are unbounded sequences. Stream-table joins (KStream-KTable) enrich each stream record with the current value of a table keyed identically, a lookup with no window. Table-table joins (KTable-KTable) combine two changelog-backed tables and emit a new changelog as either side updates; foreign-key table joins relax the same-key requirement. GlobalKTable joins replicate the whole table to every instance, so the stream side needs no co-partitioning and can join on a derived key. Each supports inner, left, and (where meaningful) outer semantics.

go deeper

for a junior

Know the three families and that stream-stream needs a window while stream-table is a lookup.

for a middle

Explain asymmetry of stream-table joins and which support outer semantics.

for a senior

Tie each join's behavior back to stream-vs-changelog semantics and co-partitioning needs.

for a principal

Frame join selection as a modeling decision driven by event-vs-state semantics, partitioning topology, and state-store cost.

## The mental model Kafka Streams has two core abstractions: - A **KStream** is an unbounded sequence of independent events (an append-only log), e.g. clicks or payments. Each record is a fact that happened. - A **KTable** is a changelog: a continuously-updated table where each key maps to its latest value. Later records with the same key overwrite earlier ones (an UPSERT). A **GlobalKTable** is a KTable replicated *in full* to every application instance rather than partitioned across them. Joins combine two of these by **key**. Which combinations exist, and how they behave, follows directly from stream-vs-table semantics. ### 1. Stream-Stream (KStream-KStream) Two unbounded streams have no "current value" to look up, so you join over a **time window** defined by `JoinWindows.ofTimeDifferenceAndGrace(...)`. A record on side A joins every record on side B whose timestamp falls within `[tA - before, tA + after]`. Both sides are buffered in **window stores** (RocksDB state). Supports inner, left, outer. ### 2. Stream-Table (KStream-KTable) For each stream record, look up the table's **current** value for that key and emit the combined result. It is *asymmetric*: only the stream side triggers output; a table update does not re-emit past stream records. No window. Supports inner and left (not outer, because there is no stream record to emit on a table-only change). ### 3. Table-Table (KTable-KTable) Both sides are changelogs. An update on either side looks up the other side's current value and emits an updated result (itself a changelog). Supports inner, left, outer. A **foreign-key (FK) join** (KIP-213) lets the left table join the right table on a key *derived* from the left value, instead of requiring identical keys. ### 4. GlobalKTable joins A `GlobalKTable` holds the entire table on every instance, so the stream side does **not** need to be co-partitioned and can join via a **key-mapper** function (join on any field of the stream value). Only `KStream-GlobalKTable` (inner/left) exists. ### Co-partitioning All non-global joins require the two inputs to be **co-partitioned**: same number of partitions and same partitioning so matching keys land on the same task. GlobalKTable and FK joins are the exceptions that relax this. ### Edge cases - Stream-table joins are sensitive to timing/ordering since they read whatever the table currently holds. - Outer/left joins on streams emit null-paired results only after the window closes (plus grace). - GlobalKTables are bootstrapped fully at startup and are not strictly time-synchronized with the stream side.

  • Why does a stream-stream join need a window but a stream-table join does not?
    A KTable has a defined 'current value' per key to look up instantly, so no window is needed. Two streams have no current value—only events over time—so you must bound the match with a time window.
  • Which join types are asymmetric (only one side triggers output)?
    Stream-table and stream-globaltable joins: only a record on the stream side produces output; a table update never re-emits past stream records.

saying these in an interview costs you the question

  • Saying stream-table joins use a time window (they don't—they're an instantaneous lookup).
  • Claiming all joins require co-partitioning (GlobalKTable and FK joins are exceptions).
  • Saying a table update re-triggers a stream-table join output.

context

open as a page

What is the co-partitioning requirement for Kafka Streams joins, and how do you satisfy or avoid it?

level: middleimportance: must knowfreq 65%

basics

~20 s

Co-partitioning means both join inputs must use the same key, the same number of partitions, and the same partitioning strategy, so matching keys land on the same task. You satisfy it by re-keying and repartitioning; GlobalKTable joins avoid it entirely.

open as a page

Explain KStream-KTable join semantics: what triggers output, what inner vs left produce, and why a GlobalKTable join differs.

level: middleimportance: must knowfreq 60%

basics

~20 s

A stream-table join enriches each stream record with the table's current value for that key. Only stream records trigger output (table updates don't). Inner emits only on a match; left emits the stream record paired with null when no match. GlobalKTable joins skip co-partitioning and join on any value-derived key.

open as a page

How does a foreign-key KTable-KTable join (KIP-213) work, and how is it different from an equi-key table join?

level: seniorimportance: should knowfreq 35%

basics

~20 s

A foreign-key join lets a left KTable join a right KTable on a key extracted from the left value, not the left's own key. Kafka Streams repartitions the left table by the foreign key internally, so the two tables don't need to be co-partitioned, and updates on either side correctly propagate.

open as a page

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

level: seniorimportance: should knowfreq 50%

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.

open as a page