What is the Processor API in Kafka Streams, and how does it differ from the high-level DSL?
answer
- DSL = declarative, PAPI = imperative
- process(Record) once per record
- Topology.addProcessor / addStateStore
- DSL compiles down to PAPI nodes
- new api package: Processor<KIn,VIn,KOut,VOut>
basics
~20 sThe Processor API (PAPI) is the low-level Kafka Streams API. You write a Processor that handles each record one at a time, can attach state stores, and forward results downstream. The DSL is a higher-level, declarative layer (map, filter, join) built on top of it.
solid answer
~40 sKafka Streams offers two layers. The DSL is declarative — you compose operators like map, filter, groupBy, and join, and the library builds the topology for you. The Processor API (PAPI) is imperative and low-level: you implement Processor<KIn,VIn,KOut,VOut> (or the older Processor<K,V>), whose process(Record) method is invoked once per input record. You get a ProcessorContext to forward records downstream, look up the record's metadata, access attached state stores, and schedule punctuators. The DSL is actually compiled down into PAPI processor nodes, so anything the DSL does, PAPI can do — but PAPI gives you fine control over state, custom forwarding, and time-driven logic (punctuation) at the cost of more boilerplate. You build a PAPI topology with Topology.addSource/addProcessor/addStateStore/addSink, or mix PAPI into a DSL pipeline via KStream.process()/processValues().
go deeper
Know PAPI is the low-level API and process() runs per record, vs. the declarative DSL.
Explain how to build a PAPI topology and that the DSL compiles down to the same processor nodes.
Discuss when to drop to PAPI, store connection, and mixing via process()/processValues().
Reason about the old vs. api package migration, topology design tradeoffs, and where imperative control pays off architecturally.
**Kafka Streams** is a Java/Scala library for building stream-processing applications on top of Apache Kafka — it reads records from input topics, transforms them, and writes results to output topics, with built-in fault tolerance and local state. It exposes two programming layers: **1. The DSL (Domain-Specific Language).** A declarative, functional API on the `KStream`/`KTable`/`GlobalKTable` types. You chain operators: `stream.filter(...).mapValues(...).groupByKey().count()`. You describe *what* transformation you want; the library figures out *how* — which processor nodes, repartition topics, and state stores to create. This is what most applications use because it is concise and safe. **2. The Processor API (PAPI).** A lower-level, imperative API. The unit of work is a **Processor**: a class implementing `org.apache.kafka.streams.processor.api.Processor<KIn, VIn, KOut, VOut>`. Its key method is `process(Record<KIn, VIn> record)`, invoked exactly once for each incoming record. Inside, you can do arbitrary logic, read/update **state stores**, and call `context.forward(...)` to emit zero, one, or many output records to downstream nodes. A Processor is created by a **ProcessorSupplier**, whose `get()` returns a fresh Processor instance per stream thread/task. **Key differences:** - *Control vs. convenience:* PAPI lets you forward an arbitrary number of records, branch dynamically, maintain custom state layouts, and run time-driven logic via **punctuators** (`context.schedule(...)`). The DSL hides all of this. - *Topology construction:* With PAPI you build the topology explicitly — `Topology.addSource(...)`, `addProcessor(...)`, `connectProcessorAndStateStores(...)`, `addSink(...)` — naming each node and wiring parent/child relationships and store access by name. The DSL builds this graph for you via `StreamsBuilder`. - *Relationship:* The DSL is **not** a separate engine — `StreamsBuilder.build()` produces a `Topology` made of the same processor nodes PAPI uses. So PAPI is strictly more general. **When to drop to PAPI:** custom state access patterns, fine-grained forwarding control, time/wall-clock-driven emission (punctuation), or behavior the DSL operators cannot express. You can also **mix** the two: stay in the DSL and inject a custom processor with `KStream.process(...)` (terminal, can forward to other processors) or `KStream.processValues(...)` (keeps the key, value-only) — getting DSL convenience plus a hand-written node where needed. **Edge case / modern note:** the original interfaces lived in `org.apache.kafka.streams.processor` (`Processor<K,V>`, `Transformer`, `ValueTransformer`). Since Kafka 2.7+/3.0 these are superseded by the typed `org.apache.kafka.streams.processor.api` package (`Processor<KIn,VIn,KOut,VOut>`, `ProcessorContext<KOut,VOut>`, `Record<K,V>`), which carry output types in the signature and pass a whole `Record` (with headers + timestamp) instead of separate key/value. Prefer the `api` package in new code.
- Can you use both the DSL and the Processor API in the same application?Yes. Stay in the DSL and splice in a custom processor with KStream.process() or processValues(); the DSL node and the PAPI node share one underlying topology. You can also connect state stores so the processor reads/writes the same store the DSL uses.
- Why might you choose PAPI over the DSL?For behavior the DSL can't express: time/wall-clock-driven emission via punctuators, custom state layouts, dynamic branching, controlled multi-record forwarding, or fine control over downstream routing.
saying these in an interview costs you the question
- Saying the DSL and PAPI are separate runtimes — the DSL compiles to a PAPI topology.
- Claiming PAPI can't use state stores — state stores are a core PAPI feature.
- Thinking process() is called once per partition or per poll — it is called once per record.