skip to content

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%

answer

  1. flatMap=0..n + may rekey; flatMapValues keeps key
  2. Empty iterable = drop record
  3. peek = pass-through side effect (non-terminal)
  4. foreach = terminal, returns void
  5. Side effects not transactional -> can duplicate

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.

solid answer

~50 s

flatMap takes a mapper returning an Iterable<KeyValue>, emitting 0..n records and possibly new keys (sets the repartition flag); flatMapValues returns an Iterable of values, keeps the key, and does NOT flag repartition. Use flatMapValues when fanning out values without rekeying. peek((k,v)->...) is a non-terminal pass-through for observability — logging, metrics, debugging — records continue downstream unchanged. foreach((k,v)->...) is terminal: it consumes the stream, returns void, and you cannot chain further. The pitfall with peek/foreach (and any side effect) is that Kafka Streams gives at-least-once processing by default, so on rebalance/retry a record can be reprocessed and the side effect runs again; with exactly-once (processing.guarantee=exactly_once_v2) the Kafka writes are transactional but EXTERNAL side effects in peek/foreach are NOT rolled back. So never put non-idempotent external writes (DB inserts, emails, payments) in peek/foreach — they can duplicate.

go deeper

for a junior

Know flatMap = one-to-many, peek = log-and-pass, foreach = terminal action.

for a middle

Use flatMapValues to avoid rekeying, know empty-iterable-drops, and that foreach is terminal.

for a senior

Explain at-least-once reprocessing and that exactly-once doesn't cover external side effects.

for a principal

Design idempotent external sinks; decide foreach vs sink connector vs to(); reason about duplicate-safety end to end.

## flatMap and flatMapValues (record fan-out) These turn ONE input record into ZERO, ONE, or MANY output records. - `flatMapValues(value -> Iterable<newValue>)` — emit several values, **key preserved** on each. Example: split a sentence record into one record per word, keeping the original key. Key-preserving, so **no repartition flag**. - `flatMap((key, value) -> Iterable<KeyValue<newKey,newValue>>)` — can emit new keys too, so it **sets the repartition flag** (like map). Use only when you genuinely need to rekey while fanning out. - Returning an empty Iterable drops the record entirely (a fan-out of zero) — a way to filter-and-expand in one step. `flatMapValues` should be preferred over `flatMap` whenever the key stays the same, for the same repartition-avoidance reason as mapValues over map. ## peek (non-terminal side effect) `peek((key, value) -> { log/metric })` runs an action for each record and then passes the record downstream **unchanged**. It returns a `KStream`, so you keep chaining. Ideal for logging, counters, tracing, debugging a topology mid-stream without altering data. ## foreach (terminal side effect) `foreach((key, value) -> { ... })` runs an action per record and **terminates** the branch — it returns `void`, so nothing can follow. Use it as a sink when you want to push records into something the DSL has no first-class sink for. The DSL alternative for writing back to Kafka is `to(...)`; foreach is for non-Kafka sinks or custom handling. ## The big pitfall: side-effect failure semantics Kafka Streams' default `processing.guarantee` is **at_least_once**. On a failure, rebalance, or crash, records after the last committed offset are **reprocessed**, so any side effect in peek/foreach **runs again** — duplicates. Even with **exactly_once_v2** (the Kafka transactional guarantee), the atomic unit is the Kafka write + offset commit + state-store changelog — all wrapped in a transaction. **External** side effects in peek/foreach (writing to a database, sending an email, calling a payment API) are NOT part of that transaction and are NOT rolled back if the transaction aborts. So exactly-once protects Kafka-internal effects, not your foreach's DB insert. Consequences / best practice: - Keep peek strictly to **idempotent, observability-only** actions (logging, metrics). Don't mutate external systems in it. - For external sinks, prefer a properly designed **sink connector** or make the write **idempotent** (upsert by key, dedup key) so reprocessing is safe. - Don't rely on peek/foreach ordering or exactly-once for external systems. ## Statelessness summary flatMapValues, peek, foreach are key-preserving and stateless (no store, no changelog). flatMap is stateless but key-changing (repartition flag).

  • Why is putting a database insert inside foreach risky even with exactly_once_v2 enabled?
    exactly_once_v2 only makes Kafka-internal effects (output records, offset commits, changelog writes) atomic. An external DB insert in foreach isn't part of that Kafka transaction and isn't rolled back on abort/retry, so it can be executed more than once. Make the external write idempotent instead.
  • How can flatMap or flatMapValues be used to drop a record entirely?
    Return an empty Iterable for that record — it produces zero output records, effectively filtering it out while still allowing fan-out for others.

saying these in an interview costs you the question

  • Saying exactly-once makes external side effects in foreach safe
  • Believing peek can modify or filter records (it's pass-through only)
  • Thinking you can chain operators after foreach (it's terminal)
  • Forgetting flatMapValues keeps the key while flatMap can rekey

context