skip to content

In Kafka Streams, what is the difference between groupByKey() and groupBy(), and why does one of them trigger a repartition?

level: juniorimportance: must knowfreq 70%

answer

  1. groupByKey = same key, no shuffle
  2. groupBy = new key -> repartition topic
  3. <app-id>-<name>-repartition
  4. all same-key records -> one task
  5. Grouped.with(serdes) for re-key

basics

~20 s

groupByKey() 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 s

Aggregations 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

for a junior

Know that groupByKey keeps the key and groupBy makes a new one, and that new keys mean reshuffling data.

for a middle

Explain the repartition topic, its naming, and when to supply Grouped Serdes.

for a senior

Reason about the cost (extra topic, latency, ordering) and that any upstream key change forces repartition.

for a principal

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.

context