skip to content

Processor API and Punctuators

Dropping below the DSL to the Processor API: manual forwarding, attached state stores, and wall-clock or stream-time punctuators. Interviewers ask when they want to hear where the DSL stops being enough.

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

questions

6

What is the Processor API in Kafka Streams, and how does it differ from the high-level DSL?

level: juniorimportance: must knowfreq 60%

answer

  1. DSL = declarative, PAPI = imperative
  2. process(Record) once per record
  3. Topology.addProcessor / addStateStore
  4. DSL compiles down to PAPI nodes
  5. new api package: Processor<KIn,VIn,KOut,VOut>

basics

~20 s

The 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 s

Kafka 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

for a junior

Know PAPI is the low-level API and process() runs per record, vs. the declarative DSL.

for a middle

Explain how to build a PAPI topology and that the DSL compiles down to the same processor nodes.

for a senior

Discuss when to drop to PAPI, store connection, and mixing via process()/processValues().

for a principal

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.

context

open as a page

What does ProcessorContext.forward() do, and how do you control where a forwarded record goes?

level: middleimportance: must knowfreq 55%

basics

~20 s

context.forward(record) sends a record from your processor to its downstream child nodes. You can call it zero, one, or many times per input record. Passing forward(record, "childName") routes it to only one named child instead of all of them.

open as a page

Explain ProcessorContext.schedule() and the difference between PunctuationType.STREAM_TIME and WALL_CLOCK_TIME punctuators.

level: seniorimportance: must knowfreq 50%

basics

~20 s

schedule() registers a punctuator — a callback that fires periodically. STREAM_TIME advances by the timestamps of records flowing through, so it only fires when data arrives and progresses. WALL_CLOCK_TIME advances by the system clock, so it fires on a real-time interval even if no records arrive.

open as a page

How do you attach and access a state store from a custom Processor, including in the DSL via process()?

level: seniorimportance: must knowfreq 45%

basics

~20 s

Register the store with a StoreBuilder (addStateStore in PAPI, or builder.addStateStore in the DSL), connect it to the processor by name, then look it up in init() via context.getStateStore("name"). In the DSL, pass the store names as the second argument to process(supplier, "storeName").

open as a page

Walk through the lifecycle of a Processor: ProcessorSupplier.get(), init(), process(), and close(). Why is a supplier used instead of a single instance?

level: middleimportance: should knowfreq 35%

basics

~20 s

A ProcessorSupplier.get() returns a new Processor for each task/thread, so instances aren't shared across threads. init() runs once per task to grab the context and stores, process() runs per record, and close() runs once when the task shuts down to release resources.

open as a page

In the modern Processor API, how do KStream.process() and processValues() relate to the deprecated transform()/transformValues(), and when do they cause a repartition?

level: seniorimportance: should knowfreq 25%

basics

~20 s

process()/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.

open as a page