skip to content

Aggregations

groupByKey and groupBy with count, reduce and aggregate, and suppressing intermediate results until a window closes. Interviewers ask because KTable update streams surprise people expecting one output per window.

part ofApache Kafkaoverview, primer and where to startread it →
on this pageshow

questions

6

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

open as a page

Compare count(), reduce(), and aggregate() on a KGroupedStream. When must you use aggregate() with an Initializer and Aggregator?

level: middleimportance: must knowfreq 65%

basics

~20 s

count() returns how many records per key. reduce() combines two values of the SAME type into one (e.g. sum of longs). aggregate() is the general form: an Initializer creates a starting value of ANY type, and an Aggregator folds each record into it — use it when the result type differs from the input.

open as a page

What problem does suppress(Suppressed.untilWindowCloses(...)) solve in windowed aggregations, and what are its operational requirements and caveats?

level: seniorimportance: must knowfreq 55%

basics

~20 s

By default a windowed aggregation emits an update on every record, so downstream sees many intermediate counts per window. suppress(Suppressed.untilWindowCloses(...)) buffers updates and emits only the FINAL result for each window after the window plus grace period has closed. It needs a buffer config (e.g. unbounded or a bytes/records limit) and only works on windowed KTables.

open as a page

How do windowed aggregations work in Kafka Streams? Cover the window types, the resulting key type, and how late records and grace periods are handled.

level: seniorimportance: must knowfreq 60%

basics

~20 s

You call windowedBy(...) on a KGroupedStream before count/reduce/aggregate, so records are bucketed by time. The result is a KTable keyed by Windowed<K> (key + window bounds). Window types include tumbling, hopping, sliding, and session windows. Late records that arrive within the grace period still update their window; beyond grace they are dropped.

open as a page

An aggregation produces a KTable. Explain what the resulting KTable's update stream emits, and why a downstream consumer may see many intermediate values per key.

level: middleimportance: should knowfreq 50%

basics

~20 s

A KTable is a changelog: each key maps to its latest aggregate. After every input record, the aggregate for that key changes, so the KTable emits a new (key, newAggregate) update. Downstream therefore sees the running aggregate update again and again, not just the final value.

open as a page

When aggregating a KTable (KGroupedTable) rather than a KStream, why do reduce()/aggregate() require both an adder and a subtractor? What goes wrong if the subtractor is incorrect?

level: principalimportance: should knowfreq 35%

basics

~20 s

A KTable is an update stream: a key's value can change or be deleted. When you re-aggregate it, an update means the OLD value must be removed from the aggregate and the NEW value added. So you pass a subtractor (remove old) and an adder (add new). A wrong subtractor leaves stale contributions, so the aggregate drifts and never self-corrects.

open as a page