skip to content

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

level: seniorimportance: must knowfreq 45%

answer

  1. StoreBuilder via Stores.*
  2. addStateStore + connectProcessorAndStateStores
  3. context.getStateStore in init(), cache it
  4. DSL: process(supplier, "storeName") or override stores()
  5. changelog topic = restore on rebalance

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").

solid answer

~40 s

A state store gives a processor fault-tolerant local state. In pure PAPI: build a StoreBuilder (e.g., Stores.keyValueStoreBuilder(Stores.persistentKeyValueStore("s"), keySerde, valSerde)), register it with topology.addStateStore(builder, "processorName") (or addStateStore + connectProcessorAndStateStores). In init(ProcessorContext), call context.getStateStore("s") once and cache the reference. In the DSL, register the StoreBuilder with streamsBuilder.addStateStore(builder), then pass the store name(s) to KStream.process(supplier, "s") (or processValues) — that connects the store to the injected processor. The ProcessorSupplier can also override stores() to declare its StoreBuilders so the DSL auto-registers them. By default a persistent store is backed by a compacted changelog topic for fault tolerance; on failure/rebalance the store is restored from it. Stores are per-task and accessed only from the owning stream thread.

go deeper

for a junior

Know a processor uses a state store via context.getStateStore(name).

for a middle

Build a StoreBuilder, register and connect it, and look it up in init().

for a senior

Explain DSL connection options (process() names vs. supplier.stores()), changelog restoration, and caching effects.

for a principal

Reason about fault tolerance, standby replicas, store sizing/RocksDB tuning, and when withLoggingDisabled/caching is appropriate.

**What a state store is.** A **state store** is a local, embedded key-value (or window/session) store that a processor uses to remember things across records — counts, last-seen values, buffered records, dedup sets. Kafka Streams ships RocksDB-backed **persistent** stores and **in-memory** stores. Each store is **per-task** (per partition group): a task owns its store partition, and only that task's stream thread touches it. Persistent stores are fault-tolerant via a **changelog topic** — every write is also appended to a compacted Kafka topic so the store can be **restored** after a crash or rebalance onto another instance. **Step 1 — define the store.** You don't just `new` a store; you describe it with a **StoreBuilder**: ``` StoreBuilder<KeyValueStore<String,Long>> builder = Stores.keyValueStoreBuilder( Stores.persistentKeyValueStore("counts"), // the supplier Serdes.String(), Serdes.Long()); // key/value serdes ``` Variants: `inMemoryKeyValueStore`, `persistentWindowStore`, `persistentSessionStore`. You can also `.withLoggingDisabled()` (no changelog — not fault tolerant) or `.withCachingEnabled()` (record cache for dedup/throttling of downstream emits). **Step 2 — register and connect (pure PAPI).** With an explicit `Topology`: ``` topology.addStateStore(builder, "MyProcessor"); ``` The second varargs argument is the **processor name(s)** that may access the store. Equivalently `topology.addStateStore(builder)` then `topology.connectProcessorAndStateStores("MyProcessor", "counts")`. Connection is what grants a processor access **and** tells Streams which store partitions to co-locate with which task. **Step 3 — look it up at runtime.** In the processor's `init(ProcessorContext<KOut,VOut> context)` (called once when the task starts), fetch and cache the store: ``` this.store = context.getStateStore("counts"); ``` Then use it in `process()`: `store.get(key)`, `store.put(key, val)`, `store.delete(key)`, range/all iterators (remember to close iterators!). Do **not** call getStateStore on every record — cache it in init. **Connecting a store in the DSL via process().** When you splice a custom processor into a DSL pipeline you must connect the store yourself. Two ways: 1. **Pass store names to process():** ``` builder.addStateStore(storeBuilder); stream.process(processorSupplier, "counts"); // "counts" connects the store ``` 2. **Let the supplier declare stores:** override `ProcessorSupplier.stores()` to return the set of `StoreBuilder`s; then the DSL **auto-registers and connects** them when you call `process(supplier)` — you don't need a separate `addStateStore` or to list names. This keeps the processor self-contained. Use `processValues(...)` instead of `process(...)` when your logic only transforms values and preserves the key — it preserves partitioning info and avoids a possible repartition. Both accept store-name varargs. **Fault tolerance & restoration.** On rebalance, the task (and its store) may move to another instance; Streams replays the changelog to rebuild the store before processing resumes (standby replicas can keep warm copies to shorten this). This is why writes must go through the store API, not a side database, if you want the same guarantees. **Edge cases / gotchas.** - **Unconnected store:** calling `getStateStore("name")` for a store you never connected to that processor throws — connection is mandatory and name-scoped. - **Threading:** never access the store from another thread; it is single-threaded per task. - **Iterators leak:** `range`/`all` return `KeyValueIterator` that holds RocksDB resources — always close (try-with-resources). - **Caching delays emits:** `.withCachingEnabled()` and the global cache can buffer/forward fewer downstream records (only the latest per key on flush); this is good for dedup but surprising if you expect every put to emit. - **Logging disabled = data loss risk:** `.withLoggingDisabled()` removes the changelog; the store can't be restored after failure — only use for re-derivable state.

  • How does a persistent state store survive an instance crash or rebalance?
    Each write is also appended to a compacted changelog topic. When the task moves, Streams restores the store by replaying that changelog; standby replicas can keep warm copies to shorten restoration.
  • What is the difference between passing store names to process() and overriding ProcessorSupplier.stores()?
    Passing names connects already-registered stores by name. Overriding stores() lets the supplier declare its StoreBuilders so the DSL auto-registers and connects them — the processor becomes self-contained with no separate addStateStore call.
  • What happens if you disable changelogging on a store?
    withLoggingDisabled() drops the changelog, so the store can't be restored after failure. Only use it for state you can re-derive; otherwise you risk silent data loss on rebalance.

saying these in an interview costs you the question

  • Calling context.getStateStore() on every record instead of once in init().
  • Thinking a store is global/shared across tasks — stores are per-task, single-threaded.
  • Forgetting to connect the store, then expecting getStateStore to work.
  • Writing to an external DB and assuming the same exactly-once/restore guarantees as a changelog-backed store.
  • Not closing KeyValueIterators from range/all.

context