In the modern Processor API, how do KStream.process() and processValues() relate to the deprecated transform()/transformValues(), and when do they cause a repartition?
answer
- process() = new Processor; processValues() = FixedKeyProcessor
- replaces deprecated transform()/transformValues()
- process() may re-key -> repartition flag
- processValues() key-immutable -> no repartition
- repartition only materializes before a stateful op
- FixedKeyRecord key is read-only
basics
~20 sprocess()/processValues() are the modern replacements for the deprecated transform()/transformValues(). process() can change the key, so Streams may mark the stream for repartitioning if a downstream operation needs it. processValues() keeps the key, so it avoids that repartition.
solid answer
~40 sSince Kafka 3.x the typed Processor API (org.apache.kafka.streams.processor.api) replaced Transformer/ValueTransformer. In the DSL you now splice custom logic with KStream.process(ProcessorSupplier, storeNames...) — terminal-style, returns void in newer signatures and forwards via context — or KStream.processValues(FixedKeyProcessorSupplier, storeNames...), which keeps the key and uses a FixedKeyRecord/FixedKeyProcessor so you can only change the value. The old transform/transformValues/flatTransform are deprecated in favor of these. Repartition implications mirror the old rules: process() (like the old transform/map/selectKey) may change the key, so Streams sets a 'requiresRepartition' flag; if a downstream stateful operation (join, aggregation, groupByKey) follows, it inserts a repartition topic to re-shuffle by the new key. processValues() (like mapValues/transformValues) guarantees the key is unchanged, so no repartition is triggered — preferred when you only transform values. Choosing processValues avoids needless network/topic overhead.
go deeper
Know process() can change the key and processValues() keeps it.
Map transform()->process() and transformValues()->processValues(), and that key changes can cause repartitioning.
Explain the repartition flag, when a topic is actually materialized, and FixedKeyProcessor's key immutability.
Reason about topology cost/latency tradeoffs, when to deliberately re-key vs. preserve partitioning, and migration strategy off the deprecated API.
**Background — the API migration.** Early Kafka Streams mixed custom logic into the DSL with `transform()`, `transformValues()`, `flatTransform()`, and `flatTransformValues()`, backed by `Transformer`/`ValueTransformer` from `org.apache.kafka.streams.kstream`. These were **untyped on output** and returned values you then had to forward awkwardly. Kafka 2.7 introduced, and 3.x standardized, the **typed Processor API** under `org.apache.kafka.streams.processor.api`: `Processor<KIn,VIn,KOut,VOut>`, `ProcessorContext<KOut,VOut>`, and `Record<K,V>`. The DSL hooks for it are: - **`KStream.process(ProcessorSupplier<KIn,VIn,KOut,VOut>, String... storeNames)`** — inject a full processor; it forwards via `context.forward(...)` and can change the key and emit 0..N records. - **`KStream.processValues(FixedKeyProcessorSupplier<KIn,VIn,VOut>, String... storeNames)`** — a **key-preserving** variant. It uses a **`FixedKeyProcessor`** receiving a **`FixedKeyRecord`** whose key is read-only; you may change only the value (and forward via a `FixedKeyProcessorContext`). `transform`/`transformValues`/`flatTransform`/`flatTransformValues` are **deprecated** in favor of `process`/`processValues`. Migration: move logic from `Transformer.transform(k,v)` returning a value into `Processor.process(Record)` calling `context.forward(...)`; move `ValueTransformerWithKey` logic into a `FixedKeyProcessor`. **Repartitioning — the key question.** Kafka Streams partitions data by key. **Stateful** operations (`groupByKey`, aggregations, `join`) require that all records for a key live in the same partition/task. If an upstream operator **might change the key**, Streams can no longer trust the existing partitioning, so it marks the stream as **requiring repartition** (an internal `repartitionRequired` flag). When a downstream **key-based** operation follows such a marked stream, Streams transparently inserts an **internal repartition topic**: it writes records out keyed by the new key and reads them back, re-shuffling so each key's records co-locate. This costs an extra topic, network round-trip, and latency. - **`process()` may change the key** (the output `Record` can have any `KOut`), so — like `map`, `selectKey`, `flatMap`, and the old `transform` — it **sets the repartition flag**. If you then `groupByKey()`/`join()`/aggregate, a repartition topic is created. - **`processValues()` guarantees the key is unchanged** (the `FixedKeyRecord` key is immutable) — like `mapValues`/`transformValues`. So it does **not** set the flag and a following stateful op needs **no** repartition. **Practical guidance.** If your custom node only needs to transform/enrich the **value** (and maybe maintain per-key state), use **`processValues()`** — it's cheaper and keeps partitioning intact. Use **`process()`** only when you genuinely need to re-key, fan out with different keys, or route to specific children. Misusing `process()` for value-only work can silently add a repartition topic when a join/aggregation follows, inflating cost and end-to-end latency. **Edge cases / nuances.** - The repartition is only **materialized** if a downstream operation actually needs the new partitioning; a `process()` followed only by a `to()`/sink won't force one. - You can suppress unnecessary repartitions in the DSL elsewhere with explicit `repartition()`/naming, but the cleanest fix is choosing `processValues()` when the key is stable. - `FixedKeyProcessor` cannot call the regular `context.forward(Record)`; it forwards a `FixedKeyRecord` via its `FixedKeyProcessorContext`, enforcing key immutability at the type level. - State store connection works the same for both: pass `storeNames` or override `stores()` on the supplier. - The deprecated transform methods still work for now but should be migrated; they share the same repartition semantics (transform = key may change, transformValues = key fixed).
- You only enrich values but used process() before a groupByKey(). What's the hidden cost?process() marks the stream as requiring repartition (it could change the key), so the groupByKey triggers an internal repartition topic — extra topic, network shuffle, and latency. Switching to processValues() keeps the key fixed and avoids it.
- How does FixedKeyProcessor enforce that you can't change the key?It receives a FixedKeyRecord whose key is read-only and forwards via FixedKeyProcessorContext using FixedKeyRecord — there's no API to set a new key, so key immutability is guaranteed at compile time.
saying these in an interview costs you the question
- Claiming process() and processValues() are interchangeable with no partitioning difference.
- Saying processValues() can change the key — it cannot; the key is fixed.
- Thinking transform()/transformValues() are the current recommended API — they're deprecated in favor of process/processValues.
- Assuming every process() always creates a repartition topic — it only materializes when a downstream key-based op needs it.
- Believing the repartition is free.