skip to content

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