In Kafka Streams, what is the difference between groupByKey() and groupBy(), and why does one of them trigger a repartition?
answer
- groupByKey = same key, no shuffle
- groupBy = new key -> repartition topic
- <app-id>-<name>-repartition
- all same-key records -> one task
- Grouped.with(serdes) for re-key
basics
~20 sgroupByKey() groups by the existing record key with no repartition. groupBy() picks a NEW key, so Streams must repartition (re-shuffle records to topic partitions) so all records with the same new key land on the same task.
solid answer
~40 sAggregations in Kafka Streams operate per key, and Streams guarantees that all records for a key go to one task. groupByKey() keeps the current key, so the data is already correctly partitioned and no shuffle is needed. groupBy(KeyValueMapper) lets you derive a brand-new key; because records with the same new key may currently live on different partitions, Streams marks the stream for repartition and writes to an internal repartition topic (named <app-id>-<name>-repartition) before the aggregation. That extra topic costs network and disk and breaks ordering guarantees relative to the source. Both return a KGroupedStream. Prefer groupByKey() when the key already fits your aggregation; only use groupBy() when you genuinely need a different grouping key, and provide explicit Serdes via Grouped.with(...) since the repartition topic must serialize the new key/value.
go deeper
Know that groupByKey keeps the key and groupBy makes a new one, and that new keys mean reshuffling data.
Explain the repartition topic, its naming, and when to supply Grouped Serdes.
Reason about the cost (extra topic, latency, ordering) and that any upstream key change forces repartition.
Design topologies to minimize repartition topics, choose keys deliberately, and explain the partitioning invariant that makes it mandatory.
## What grouping is Kafka Streams is a library for processing records (key/value pairs) read from Kafka topics. An **aggregation** (like count or sum) combines many records that share the same key into one result. Before you can aggregate, you must declare *how records are grouped* — that produces a `KGroupedStream` (from a `KStream`) or `KGroupedTable` (from a `KTable`). ## The partitioning invariant A Kafka topic is split into **partitions**. Streams assigns partitions to **tasks**, and each task owns its own local state store. The core invariant: **all records with the same key must be processed by the same task**, otherwise a per-key aggregate would be split across machines and be wrong. Kafka's default producer partitioner sends a record to `hash(key) % numPartitions`, so records with equal keys already share a partition — *as long as the key doesn't change*. ## groupByKey() `groupByKey()` groups by the record's **existing** key. Since the upstream data was already partitioned by that key, no data movement is needed. It is cheap. You can pass `Grouped.with(keySerde, valueSerde)` to control serialization and the internal name. ## groupBy() `groupBy((key, value) -> newKey)` lets you **choose a new key** derived from the key and/or value. Now records that should be grouped together may sit on different partitions. Streams therefore sets a **repartition** flag: it writes the re-keyed records to an internal **repartition topic** named like `<application.id>-<operator-name>-repartition`, partitioned by the new key, and reads them back. This re-shuffle (a) costs an extra topic (network + broker disk), (b) adds latency, and (c) means the new stream's ordering is only guaranteed per new-key, not globally. ## Edge cases - Mark a stream as already-partitioned with `KStream.repartition()` only when needed; you cannot skip the repartition for `groupBy` because Streams can't know your data is pre-partitioned by the new key. - Always supply Serdes via `Grouped.with(...)` for `groupBy`; otherwise Streams falls back to the configured `default.key.serde`/`default.value.serde`, and a mismatch causes ClassCastExceptions at the repartition topic. - `selectKey()` followed by `groupByKey()` is equivalent to `groupBy()` and also triggers repartition — changing the key anywhere upstream sets the repartition-required flag.
- What is the name pattern of the internal topic created by groupBy(), and how can you control it?It is <application.id>-<operatorName>-repartition. You influence the operator name (and Serdes) by passing Grouped.as("my-name") or Grouped.with(keySerde, valueSerde) to groupBy().
- Does selectKey().groupByKey() avoid the repartition that groupBy() incurs?No. Changing the key with selectKey() (or map) sets the repartition-required flag, so the subsequent groupByKey() still creates a repartition topic. It is functionally the same as groupBy().
saying these in an interview costs you the question
- Claiming groupBy() never repartitions or is as cheap as groupByKey().
- Saying groupByKey() can re-key records (it cannot change the key).
- Thinking the repartition is optional and Streams figures out data is already partitioned by the new key.