skip to content

What are the join types and co-partitioning requirements in ksqlDB, and how do stream-stream, stream-table, and table-table joins differ?

level: seniorimportance: should knowfreq 50%

answer

  1. Stream-stream = WITHIN window → STREAM
  2. Stream-table = enrichment lookup, fires on stream only
  3. Table-table = materialized, re-emits on either change
  4. Co-partition: same key + same partition count + same hash
  5. Non-key join → auto repartition topic

basics

~20 s

Stream-stream joins are windowed (you join events within a time window) and produce a stream. Stream-table joins are non-windowed enrichment lookups against the table's current state. Table-table joins maintain a continuously-updated joined table. All require co-partitioned inputs (same key, same partition count).

solid answer

~60 s

ksqlDB supports INNER, LEFT, and (for stream-stream and table-table) FULL OUTER joins, but semantics depend on the input types. A **stream-stream** join must specify a `WITHIN` window because both sides are unbounded event flows — it matches events whose timestamps fall within that interval and emits a stream of joined events; it keeps both sides' recent records in windowed state. A **stream-table** join is **not** windowed: each stream event is enriched by looking up the **current** value in the table — a temporal/'as-of' lookup that only fires on the stream side. A **table-table** join continuously maintains a materialized joined table, re-emitting when either side changes. A hard prerequisite for all of them is **co-partitioning**: both inputs must be keyed on the join column, have the **same number of partitions**, and use the same partitioning strategy — otherwise matching records may sit on different partitions/nodes. ksqlDB will auto-**repartition** a side (creating an internal repartition topic) when you join on a non-key column. Foreign-key table-table joins relax the same-key requirement in newer versions.

go deeper

for a junior

Know the three join shapes exist and that stream-stream joins use a time window.

for a middle

Explain stream-table enrichment vs windowed stream-stream, and which produce a stream vs a table.

for a senior

Cover co-partitioning rules, auto-repartitioning on non-key joins, fires-on semantics, and outer-join support per type.

for a principal

Reason about timing/ordering pitfalls in enrichment, foreign-key table-table joins, repartition-topic cost, and designing topic keys/partition counts for joinability at scale.

## Joining streams of data A **join** combines records from two collections on a join key. Because ksqlDB operates on unbounded data, *what 'match' means* depends on whether each side is a STREAM (events) or a TABLE (current state). ### Co-partitioning: the universal prerequisite Kafka data is split into **partitions**, and a given key always hashes to the same partition. For a join to find matching records on the **same node/partition**, both inputs must be **co-partitioned**: 1. **Same join key** — joined on the column that is each side's key. 2. **Same number of partitions** on both topics. 3. **Same partitioning strategy** (same partitioner/hash) so equal keys land on equal partition numbers. If you join on a **non-key** column, ksqlDB transparently inserts a **repartition** (an internal topic re-keyed by the join column) so co-partitioning holds. If partition counts differ, you must repartition one side yourself. Violating co-partitioning silently produces missing matches — a classic gotcha. ### Stream-stream join — windowed Both sides are unbounded event streams, so 'match' must be bounded in **time**. You supply a `WITHIN` clause: ```sql SELECT * FROM orders o JOIN shipments s WITHIN 1 HOUR ON o.order_id = s.order_id EMIT CHANGES; ``` An order matches a shipment if their **timestamps are within 1 hour** of each other. ksqlDB buffers each side's recent records in **windowed state stores**. Output is a **STREAM** of joined events. Supports INNER/LEFT/FULL OUTER. Use for correlating two event flows (e.g., request/response, click/conversion). ### Stream-table join — enrichment (not windowed) Here one side is the **current state**. Each **stream** event triggers a lookup of the **latest** value for that key in the **table**: ```sql SELECT * FROM clicks c JOIN users u ON c.user_id = u.user_id EMIT CHANGES; ``` This is a **temporal/as-of** join: the click is enriched with the user's value *as it was at processing time*. It fires only on the **stream** side — a table update does **not** retroactively re-emit past stream events. Supports INNER and LEFT (no FULL OUTER). This is the bread-and-butter pattern for enriching events with reference/dimension data. ### Table-table join — materialized Both sides are evolving state. ksqlDB maintains a **continuously-updated joined TABLE**; a change on **either** side re-evaluates and re-emits the joined row: ```sql SELECT * FROM accounts a JOIN profiles p ON a.id = p.id EMIT CHANGES; ``` Supports INNER/LEFT/FULL OUTER. The result is itself materialized and can serve pull queries. Newer ksqlDB also supports **foreign-key (n:1) table-table joins**, which relax the requirement that the join column be the right side's primary key (and don't require identical partition counts the same way). ### Summary of differences | Join | Windowed? | Fires on | Output | Outer support | |------|-----------|----------|--------|---------------| | stream-stream | yes (`WITHIN`) | either side within window | STREAM | INNER/LEFT/FULL | | stream-table | no | stream side only | STREAM | INNER/LEFT | | table-table | no | either side | TABLE | INNER/LEFT/FULL | ### Edge cases - A stream-table join with a **null** table value (or no match) yields no row for INNER, or null-filled columns for LEFT. - Stream-stream windowed state respects a **grace/retention** bound; very old events expire and won't match. - The stream-table join uses the table value **current at the time the stream record is processed**, which depends on relative arrival/timing of the two topics — a common source of 'why didn't it enrich' confusion (table not yet loaded). - Foreign-key table-table joins re-key internally and have their own subscription/response topics.

  • Why does a stream-stream join require a WITHIN clause but a stream-table join does not?
    Both sides of a stream-stream join are unbounded event flows, so 'match' must be bounded in time (WITHIN) and both sides buffered in windowed state. A stream-table join looks up the table's current state per stream event — there's a single current value, so no time window is needed.
  • What is co-partitioning and what breaks if it's violated?
    Inputs must share the join key, the same partition count, and the same partitioning strategy so equal keys land on the same partition/node. If violated, matching records sit on different partitions and the join silently misses them. ksqlDB repartitions automatically when you join on a non-key column.
  • In a stream-table join, does a later update to the table re-emit past stream events?
    No. The join fires only on the stream side; each stream event is enriched with the table's value at processing time. Updating the table afterward does not retroactively re-emit earlier joined rows.

saying these in an interview costs you the question

  • Saying stream-stream joins are non-windowed
  • Claiming a stream-table join re-emits when the table changes (it fires only on the stream side)
  • Forgetting co-partitioning / same-partition-count requirement
  • Saying stream-table joins support FULL OUTER (they support INNER/LEFT only)

context