Walk through the lifecycle of a Processor: ProcessorSupplier.get(), init(), process(), and close(). Why is a supplier used instead of a single instance?
answer
- Supplier.get() -> new instance per task (no sharing)
- init() once: cache context, stores, schedule punctuators
- process() per record: logic + forward
- close() once: release your resources (Streams owns stores)
- stores() override = DSL auto-register
- init re-runs after rebalance
basics
~20 sA 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.
solid answer
~50 sKafka Streams runs many tasks across stream threads, and each must have its own Processor instance because a Processor holds mutable, non-thread-safe state (cached store references, context, buffers). So you supply a ProcessorSupplier whose get() returns a fresh Processor each time Streams needs one — never share one instance. Lifecycle per instance: init(ProcessorContext) is called once when the task starts; here you cache context, look up state stores via context.getStateStore, and register punctuators via context.schedule. process(Record) is called once per input record for the transform/forward logic. close() is called once on task shutdown, rebalance, or migration — release external resources, but note Streams manages store flushing/closing for you. The supplier may also override stores() to declare its StoreBuilders for DSL auto-registration. get() must return a new instance each call; returning a shared instance is a classic bug.
go deeper
Know the order: get -> init -> process (per record) -> close.
Explain why a supplier exists (per-task instances) and what belongs in init vs. process vs. close.
Discuss rebalance re-init, stores() auto-registration, and not managing store lifecycle yourself.
Reason about thread-affinity guarantees, crash semantics vs. close(), and per-task resource initialization tradeoffs.
**Why a supplier, not an instance?** Kafka Streams parallelizes work into **tasks** (one per input partition group), executed by a pool of **stream threads** (`num.stream.threads`). Each task that uses your processor needs its **own** Processor object, because a Processor typically holds **mutable, non-thread-safe** state: a cached `ProcessorContext`, cached state-store handles, in-flight buffers, counters. If two threads shared one instance, those fields would be corrupted. To guarantee isolation, you don't hand Streams an instance — you hand it a **`ProcessorSupplier<KIn,VIn,KOut,VOut>`**, a factory whose **`get()`** returns a **brand-new Processor each time** Streams asks. > Classic bug: implementing `get()` to return the *same* cached instance (or a lambda capturing one). Under multi-threading this shares mutable state across tasks/threads and causes data races. `get()` must `return new MyProcessor(...)` every call. **Lifecycle of one Processor instance:** 1. **`get()`** (on the supplier) — Streams calls it to obtain the instance for a task. Construction should be cheap; don't open stores here (the context/stores aren't available yet). 2. **`init(ProcessorContext<KOut,VOut> context)`** — called **once**, right after the task is assigned/started and stores are ready. This is where you: - cache the `context` (you need it later to forward), - resolve and cache **state stores**: `this.store = context.getStateStore("name")`, - register **punctuators**: `context.schedule(interval, type, punctuator)` (keep the `Cancellable`), - initialize any per-task resources. After a rebalance moves the task elsewhere, the new instance gets its own `init()` once restoration completes. 3. **`process(Record<KIn,VIn> record)`** — called **once per input record**, on the owning stream thread. Core logic: read/update stores, decide output, call `context.forward(...)` zero/one/many times. This is the hot path; keep it fast and non-blocking. 4. **`close()`** — called **once** when the task is shut down, reassigned during a rebalance, or the application stops. Release **your** external resources (clients, files, cancel custom timers). You generally do **not** close state stores or flush changelogs yourself — Streams owns store lifecycle, commit, and flushing. After `close()`, the instance is discarded; a fresh one is created if the task restarts. **ProcessorSupplier.stores().** The supplier can override `stores()` to return the set of `StoreBuilder`s the processor needs. When you inject the supplier via `KStream.process(supplier)` in the DSL, Streams **auto-registers and connects** those stores — so the processor is self-contained without a separate `addStateStore`/store-name list. This pairs naturally with the lifecycle: `stores()` declares, `init()` resolves. **Edge cases / gotchas.** - *Don't forward in init()/close():* forwarding belongs in `process()` or a punctuator; the topology/commit context around init/close differs. - *init() runs again after rebalance:* re-register punctuators every init; don't assume one-time-per-JVM. - *Heavy construction in get():* avoid — get() may be called for many tasks; do per-task setup in init(). - *close() not guaranteed on hard crash:* a kill -9 skips close(); rely on the changelog/transactions for correctness, not on close() side effects. - *Thread affinity:* all four methods for a given instance run on the same stream thread, so the instance's fields need no synchronization — but never touch them from another thread.
- What is the bug if ProcessorSupplier.get() returns a shared singleton instance?Multiple tasks/threads share one Processor's mutable state (context, store handles, buffers), causing data races and corruption. get() must return a new instance per call.
- Should you flush or close state stores in close()?No. Streams owns store lifecycle, commit, and flushing. In close() you only release your own external resources (clients, files, custom timers).
saying these in an interview costs you the question
- Returning a shared instance from get().
- Looking up state stores or scheduling punctuators in the constructor instead of init().
- Assuming init() runs only once per JVM — it runs per task start, including after rebalances.
- Manually closing/flushing state stores in close().
- Relying on close() running on hard crashes.