What is the co-partitioning requirement for Kafka Streams joins, and how do you satisfy or avoid it?
answer
- same key, same #partitions, same partitioner
- join runs per-task / per-partition
- repartition() or selectKey to fix
- GlobalKTable + FK join = exceptions
- TopologyException on partition mismatch
basics
~20 sCo-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.
solid answer
~40 sNon-global Kafka Streams joins are executed per partition by independent tasks, so a record can only join its counterpart if both land on the same partition number. That requires co-partitioning: (1) identical keys (and key serdes), (2) the same partition count on both topics, and (3) the same partitioner so equal keys hash to the same partition. If keys differ, you call `selectKey`/`map` and then `repartition()` (or `through`/`groupByKey`) to write to an internal repartition topic with matching partitioning. Kafka Streams validates partition counts at startup and throws a TopologyException on mismatch. The escape hatches are GlobalKTable joins (table fully replicated to every instance, so the stream needs no co-partitioning and joins via a key-mapper) and foreign-key KTable joins (KIP-213), which internally repartition by the foreign key.
go deeper
Know that joined topics must share key and partition count.
Enumerate the three conditions and fix mismatches with selectKey + repartition.
Explain task-per-partition execution as the root cause and the silent-failure mode of a partitioner mismatch.
Weigh repartition cost vs GlobalKTable replication vs FK-join machinery when designing a topology's partitioning strategy.
## Why co-partitioning exists Kafka Streams parallelizes by assigning each **partition** of the input topics to a **task**. A task only sees the data in its own partitions and maintains its own local state stores. For a join, the matching record from the other side must be in the *same task*—otherwise it simply isn't visible. Records are routed to partitions by `partition = hash(key) % numPartitions`. So two records with the same join key land together **only if** both topics agree on key, partition count, and partitioner. ## The three conditions 1. **Same key + key serde:** Records must be keyed by the join attribute, serialized identically. 2. **Same number of partitions:** topic A and topic B must have the same partition count. 3. **Same partitioner:** equal keys must hash to the same partition number (the default `DefaultPartitioner`/murmur2 satisfies this; a custom partitioner on one side breaks it). Kafka Streams checks partition counts when building the topology and fails fast with a `TopologyException` if they differ. It does **not** automatically detect a custom-partitioner mismatch—that produces silent missing joins. ## How to satisfy it If one side has the wrong key: ``` stream.selectKey((k, v) -> v.joinId()) .repartition() // writes to an internal repartition topic ``` The repartition topic is created with a partition count matching the co-partition group, restoring alignment. Aggregations (`groupBy`) also insert a repartition automatically when the key changes. ## How to avoid it - **GlobalKTable:** the entire table lives on every instance, so any stream record can be enriched locally. No co-partitioning; join via a `KeyValueMapper` that extracts the lookup key from the stream value. Cost: full replication + memory on every instance, eventual consistency, no time synchronization. - **Foreign-key KTable join (KIP-213):** joins a left table to a right table on a key derived from the left value. Kafka Streams handles the repartitioning by foreign key internally, including subscription/response topics, so you don't manually re-key. ## Edge cases & gotchas - A repartition adds an extra topic, network hop, and ordering boundary—use it deliberately. - Upstream operators that change the key set a 'repartition required' flag; a downstream join then forces a repartition automatically in newer versions, but explicit `repartition()` makes the topology clear. - Mismatched partition counts caught at startup; mismatched partitioners cause silent data loss in join output—test for it.
- What happens at runtime if two join inputs have a different number of partitions?Kafka Streams validates this when building the topology and throws a TopologyException at startup, refusing to run rather than producing wrong results.
- How does a GlobalKTable join sidestep co-partitioning, and what's the tradeoff?The table is fully replicated to every instance, so any stream record can look up any key locally via a key-mapper. The tradeoff is full data duplication per instance, higher memory, and no time-synchronized/strictly ordered semantics with the stream.
saying these in an interview costs you the question
- Saying co-partitioning is only about key equality (partition count and partitioner matter too).
- Claiming Kafka Streams auto-repartitions both sides of every join (it repartitions where it detects a key change / for FK joins, but you must align inputs).
- Believing a custom partitioner on one side is harmless—it silently breaks matching.