Command Palette

Search for a command to run...

Hectal
PHASE 11Advanced ~9 min· topic 5 of 5

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.

0/5 · 0%

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

  1. 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.

  2. 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.

  3. 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.

  4. 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.

  5. 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.

  6. 06

    Repartitioning: after selectKey, map or groupBy with a new key, Streams writes records to an internal -repartition topic 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

OrderEnrichment.javajava
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)
join.mermaiddiagram
Rendering diagram…

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

01

What is co-partitioning and why do non-global joins require it?

Practice

01

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.