Explain exactly when and how Kafka Streams creates a repartition topic from a stateless key-changing operator, and how to control or avoid it.
answer
- selectKey/map/flatMap set the flag
- Lazy: topic only on downstream stateful op
- Name: app-id + node + -repartition
- groupBy = selectKey + groupByKey
- repartition(Repartitioned.as(...)) to control
basics
~20 sKey-changing stateless ops (selectKey, map, flatMap) set a repartition-required flag. When a stateful op follows, Streams injects an internal repartition topic (named app-id + node + '-repartition') to re-shuffle by the new key. Use mapValues to avoid it, or repartition() to control it.
solid answer
~40 sStreams marks a stream as 'needs repartition' when a key-changing operator runs: selectKey, map, flatMap, or transform/process that may emit a new key. The flag is lazy — nothing happens until a downstream operator requires co-partitioning (groupBy/groupByKey + aggregation, joins, windowed ops). At that point Streams inserts an internal repartition topic, typically `<application.id>-<generated-or-named-node>-repartition`, with the same partition count as the source by default, written and re-consumed so records land on the partition matching the new key's hash. This costs an extra topic, broker storage, a network round-trip, and ordering only within the new key. To avoid it: don't change the key (mapValues/flatMapValues/filter). To control it: call `KStream.repartition(Repartitioned.as("name").withNumberOfPartitions(n).withStreamPartitioner(...))` to fix the topic name, partition count, and partitioner explicitly instead of relying on the auto-generated one. groupBy is essentially selectKey + repartition.
go deeper
Know that changing the key may cause re-shuffling.
Explain the flag and that mapValues avoids it.
Explain laziness, the topic naming, default partition count, and explicit repartition() control.
Reason about stable internal-topic names across upgrades, partition-count budgets, and sharing one repartition across multiple stateful consumers.
## The problem repartitioning solves Kafka spreads a topic's records across **partitions**; the producer chooses a partition by hashing the record **key** (`hash(key) % numPartitions`). Kafka Streams assigns partitions to **tasks**, and stateful operators assume **co-partitioning**: all records sharing a key are in the same partition, hence the same task, hence the same state store. If you change the key mid-topology, that assumption breaks — records with the new key are scattered across partitions by their OLD key. They must be physically re-shuffled. ## When the flag is set Streams keeps a `repartitionRequired` boolean on each stream node. It is set true by **key-changing** stateless operators: - `selectKey` - `map` - `flatMap` - `transform` / `process` / `flatTransform` (which may emit new keys) It is NOT set by key-preserving ops: `mapValues`, `flatMapValues`, `filter`, `filterNot`, `peek`, `foreach`, `merge`, `branch`/`split`. The flag is set **conservatively** — even `map` that returns the same key sets it, because Streams cannot inspect your lambda. ## When the topic is actually created (laziness) The flag alone produces nothing. A repartition **topic** materializes only when a downstream operator needs correct partitioning: - `groupByKey()` / `groupBy()` followed by `count`/`reduce`/`aggregate` - `join` / `leftJoin` / `outerJoin` (stream-stream or stream-table) - windowed aggregations So `source.selectKey(...).to("out")` creates **no** repartition topic — you just wrote rekeyed records straight out. But `source.selectKey(...).groupByKey().count()` creates one. ## The topic itself Naming: `<application.id>-<node-name>-repartition` (the node name is auto-generated, e.g. `KSTREAM-KEY-SELECT-0000000003`, unless you name it). It's an internal topic: default partition count = number of partitions of the source it derives from, `cleanup.policy=delete`, short retention managed by Streams. Records are produced to it and re-consumed, giving a full broker round-trip (durability + latency cost). Ordering is preserved only **per new key**. ## Controlling it - **Avoid:** keep the key — `mapValues`/`flatMapValues`/`filter`. - **Name and size it:** `stream.repartition(Repartitioned.as("my-repart").withNumberOfPartitions(12).withKeySerde(...).withValueSerde(...).withStreamPartitioner(...))`. Explicit `repartition()` gives a stable topic name (important across topology changes / upgrades) and lets you scale partitions or supply a custom `StreamPartitioner`. - **Reuse:** if multiple stateful ops follow one rekey, insert one explicit `repartition()` so Streams doesn't create several. - `groupBy(keySelector)` is sugar for `selectKey(...).groupByKey()` and triggers the same repartition. ## Operational notes Repartition topics count toward broker partition limits and disk. Renaming nodes (by adding/removing operators) changes auto-generated topic names, which can orphan old internal topics on upgrade — another reason to name them explicitly with `Repartitioned.as(...)` in long-lived apps.
- Why might you call repartition() explicitly even though Streams would auto-create the topic?To pin a stable topic name (auto-generated names shift when the topology changes, orphaning old internal topics on upgrade), to set a specific partition count or serdes, to supply a custom StreamPartitioner, or to share one repartition across several downstream stateful ops instead of creating several.
- Does selectKey followed only by to("output-topic") create a repartition topic?No. The flag is set but never consumed by a co-partitioning operator. The rekeyed records are written straight to the output topic. Repartition topics appear only before stateful ops that require co-partitioning.
saying these in an interview costs you the question
- Saying selectKey immediately creates a repartition topic regardless of what follows
- Claiming repartition topics use a state store / changelog
- Thinking mapValues can trigger a repartition
- Forgetting that groupBy is selectKey + repartition