skip to content

Stateless Operations

The stateless operators such as map, filter, flatMap, branch and merge, and which of them change the key and force a repartition. Interviewers ask because an accidental selectKey adds a hidden internal topic.

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

questions

6

In Kafka Streams, what does it mean for a DSL operation to be 'stateless', and which common operators fall into that category?

level: juniorimportance: must knowfreq 70%

answer

  1. No memory of past records
  2. No state store, no changelog
  3. map/filter/flatMap/selectKey/branch/merge/peek/foreach
  4. Output type usually KStream
  5. Rekeying is the only catch

basics

~10 s

A stateless operation processes each record independently, keeping no memory of past records. Examples: map, mapValues, filter, filterNot, flatMap, flatMapValues, selectKey, branch/split, merge, peek, foreach.

solid answer

~40 s

Stateless operations transform a record using only that record itself — they hold no state store and never remember earlier records, so the output for a given input never depends on history. The result type is typically a KStream (an unbounded event stream). The core stateless transforms are map/mapValues (rekey/reshape values), filter/filterNot (predicate keep/drop), flatMap/flatMapValues (one record to zero-or-more), selectKey (set a new key), branch/split (route to multiple streams), merge (combine streams), and the side-effect operators peek and foreach. Because there is no state, these operators are cheap, need no RocksDB store or changelog topic, and recover trivially on failure. Aggregations, joins, windowing and reduce are the stateful counterparts.

go deeper

for a junior

Just know stateless = each record handled on its own, no memory, and recognize the operator list.

for a middle

Explain the no-store/no-changelog consequence and which operators rekey.

for a senior

Connect rekeying to repartition topics and recovery/scaling cost trade-offs.

for a principal

Reason about topology design: minimize rekeys, place stateless filters early, teach the state vs stateless boundary.

## What 'stateless' means Kafka Streams is a Java/Kotlin library for processing records (key/value pairs) flowing through Kafka topics. Its DSL (high-level API) builds a topology — a graph of processing nodes. Each node is either **stateless** or **stateful**. A **stateless** operation computes its output purely from the **current record**. It keeps no memory of records it has already seen. Process the same record twice and you get the same output; the order of unrelated records does not change any single record's result. Because nothing is remembered, there is: - **no state store** (no RocksDB on disk, no in-memory map), - **no changelog topic** (the internal Kafka topic that backs up state for fault tolerance), - **trivial recovery** — on restart there is nothing to restore. A **stateful** operation, by contrast, accumulates information across records: counts, sums, joins, windowed aggregates. It needs a state store and a changelog. ## The stateless operators - **mapValues** — transform the value, keep the key. Does NOT change the key, so it never forces repartitioning. - **map** — transform key AND value (returns a new `KeyValue`). Changing the key marks the stream for possible repartition. - **filter / filterNot** — keep records where a predicate is true (filter) or false (filterNot). Drops the rest. - **flatMap / flatMapValues** — turn one input record into zero, one, or many output records (flatMapValues keeps the key; flatMap may rekey). - **selectKey** — set a brand-new key from the record; value untouched. Always marks for repartition. - **branch / split** — route one stream into several streams by predicates (modern API: `KStream.split().branch(...).defaultBranch()` returning a `Map<String, KStream>`; the old `KStream.branch(...)` returning a `KStream[]` array is deprecated). - **merge** — interleave two streams of the same key/value type into one. - **peek** — run a side effect (e.g. logging, metrics) per record WITHOUT modifying the stream; records pass through unchanged. - **foreach** — terminal side effect; consumes the stream and returns `void` (no downstream node). ## Why it matters Stateless nodes are cheap, horizontally scalable, and fault-tolerant for free. The one subtlety is that some of them (`map`, `selectKey`, `flatMap`, `transform` with key change) **change the record key**, and a downstream stateful operation will then trigger an automatic **repartition topic** so records with the same new key land on the same partition. So even 'stateless' rekeying has a downstream cost.

  • Name two side-effect-only stateless operators and the difference between them.
    peek and foreach. peek runs a side effect and passes records through unchanged (non-terminal); foreach is terminal — it consumes the stream and returns void with no downstream node.
  • If stateless ops keep no state, why can map still cause an internal Kafka topic to appear?
    map can change the record key. A downstream stateful op (aggregation/join) then needs a repartition topic so all records with the same new key go to the same partition. The repartition topic is the cost, not a state store.

saying these in an interview costs you the question

  • Claiming mapValues changes the key (it does not — only map/selectKey/flatMap can)
  • Saying stateless ops use RocksDB or a changelog topic
  • Confusing peek (pass-through) with foreach (terminal)
  • Listing aggregate/reduce/join/count as stateless

context

open as a page

What is the practical difference between map and mapValues (and filter vs transformValues) in Kafka Streams, and why would you prefer mapValues?

level: middleimportance: must knowfreq 65%

basics

~10 s

map can change both key and value; mapValues changes only the value and keeps the key. Prefer mapValues because keeping the key avoids a repartition topic when a stateful operation follows.

open as a page

Explain exactly when and how Kafka Streams creates a repartition topic from a stateless key-changing operator, and how to control or avoid it.

level: seniorimportance: must knowfreq 60%

basics

~20 s

Key-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.

open as a page

How do you route a single KStream into multiple branches in modern Kafka Streams, and how do you recombine streams? Contrast split/branch with merge.

level: middleimportance: should knowfreq 45%

basics

~10 s

Use KStream.split() then chained .branch(predicate, Branched.as("name")) and optionally .defaultBranch(), which returns a Map<String, KStream> of routed sub-streams. merge(otherStream) does the opposite: it interleaves two same-typed streams into one. Both are stateless.

open as a page

Compare flatMap/flatMapValues with the side-effect operators peek and foreach. When is each appropriate, and what are the failure-semantics pitfalls?

level: middleimportance: should knowfreq 40%

basics

~20 s

flatMap/flatMapValues turn one input record into zero-or-many output records (flatMapValues keeps the key). peek runs a side effect per record and passes records through unchanged. foreach is terminal: it runs a side effect and ends the branch with no downstream. Side effects aren't transactional, so they can repeat on retries.

open as a page

What does selectKey do, how does it relate to groupBy and the deprecated through(), and what is the modern replacement?

level: seniorimportance: should knowfreq 35%

basics

~20 s

selectKey sets a new record key (value unchanged) and flags repartition. groupBy is essentially selectKey + groupByKey. The old through() (write-to-topic-then-reconsume) is deprecated; use repartition() instead for the implicit re-keying topic, or to()+stream() for an explicit user topic.

open as a page