Topic 11.5
Joins, Co-Partitioning, and Repartitioning
In one line
Streams can join stream–stream (within a time window), stream–table (enrich events with latest state), table–table (combine latest states), and stream–GlobalKTable (enrich without co-partitioning). Non-global joins require co-partitioning: same key, same partition count, same partitioner. Changing keys triggers repartition topics, which cost network, storage and latency.
Think of it like this
Matching delivery notes with purchase orders. It's easy if both piles are sorted into the same numbered trays by order number (co-partitioned). If one pile is sorted by customer instead, someone must re-sort it first (repartitioning).
Key ideas
- 01
Stream–table join (
orders.join(customers, ...)): for each order event, look up the customer's current state in the local table store. Most common enrichment. The table must be keyed like the stream and co-partitioned. - 02
Stream–stream join: both sides are events; they must occur within a join window (
JoinWindows.ofTimeDifferenceWithNoGrace(10m)), for example matching an order with its payment within 10 minutes. Both sides are buffered in window stores. - 03
Table–table join: produces a KTable updated when either side changes; foreign-key joins (
KTable.join(otherTable, foreignKeyExtractor, ...)) join on a field other than the key. - 04
GlobalKTable join: every instance holds the full table, so the stream needn't be co-partitioned and can join on any derived key; fine for small reference data (countries, products), expensive for large tables.
- 05
Co-partitioning: records with the same key must be in the same partition number on both topics, which requires the same partition count and partitioning strategy. Kafka Streams checks partition counts and throws if they differ.
- 06
Repartitioning: after
selectKey,maporgroupBywith a new key, Streams writes records to an internal-repartitiontopic keyed by the new key, then reads them back. It's automatic but adds a network round trip, broker storage and latency; minimise re-keying and reuse repartitioned streams.
Code & diagrams
KStream<String, Order> orders = b.stream("commerce.orders"); // key = orderId
KTable<String, Customer> customers = b.table("customers.state"); // key = customerId
orders
.selectKey((orderId, o) -> o.customerId()) // re-key -> automatic repartition topic
.join(customers, // stream-table join, co-partitioned by customerId
(order, customer) -> EnrichedOrder.of(order, customer),
Joined.with(Serdes.String(), orderSerde, customerSerde))
.selectKey((customerId, e) -> e.orderId()) // back to orderId for downstream ordering
.to("commerce.orders.enriched");
// Internal topics created:
// <app-id>-KSTREAM-KEY-SELECT-0000000001-repartition
// <app-id>-KSTREAM-KEY-SELECT-0000000005-repartition (on write-out if later grouped)Interview problem
The problem
Join orders with customer data
Topics orders (key orderId) and customers (key customerId, latest profile). Produce orders enriched with customer information. Design the join, covering co-partitioning, keys, state stores and repartitioning.
When it breaks
customers has 12 partitions and orders.by-customer has 24
What you see
The join can't be co-partitioned; Kafka Streams fails at startup with a TopologyException about partition counts (or, with manual consumers, joins silently miss records).
Fix & prevent
Match partition counts, let Streams repartition via an internal topic, or use a GlobalKTable.
Explain it without notes
What is co-partitioning and why do non-global joins require it?
Practice
Implement the enrichment with TopologyTestDriver, including an order that arrives before its customer (inner vs left join).
Trade-offs
- ↔
GlobalKTables avoid repartitioning but replicate the whole table everywhere; re-keying costs a repartition topic.
Done when you can
I can design stream-table, stream-stream and table-table joins with correct co-partitioning.