How do you create and configure a state store in Kafka Streams? Contrast Materialized (DSL) with Stores/StoreBuilder (Processor API).
answer
- DSL → Materialized.as(name).withKeySerde/...
- PAPI → Stores factory + StoreBuilder + addStateStore
- Materialized: logging/caching/serdes/retention
- Name it to make it queryable
- context.getStateStore(name) inside processor
basics
~20 sIn the DSL you configure stores with Materialized — passing it to operators like count() or aggregate() to set the store name, serdes, and persistence. In the lower-level Processor API you build a store explicitly with the Stores factory and a StoreBuilder, then register it on the topology.
solid answer
~40 sTwo paths. In the **DSL**, stateful operators accept a `Materialized<K, V, S>` object: `count(Materialized.as("counts-store").withKeySerde(...).withValueSerde(...))`. Materialized also exposes `withLoggingDisabled()`, `withLoggingEnabled(configs)`, `withCachingDisabled()`, and `withRetention(...)` for windowed stores; the store type (persistent RocksDB vs in-memory) is chosen via the supplier, e.g. `Materialized.as(Stores.inMemoryKeyValueStore("x"))`. In the **Processor API**, you build the store yourself: pick a supplier from the `Stores` factory (`Stores.persistentKeyValueStore`, `inMemoryKeyValueStore`, `persistentWindowStore`, `persistentSessionStore`), wrap it in a `StoreBuilder` (`Stores.keyValueStoreBuilder(supplier, keySerde, valueSerde)`), then attach it with `topology.addStateStore(builder, ...)` or `streamsBuilder.addStateStore(builder)` and connect it to a processor by name. The `StoreBuilder` is where you call `.withLoggingEnabled/Disabled` and `.withCachingEnabled/Disabled`. The DSL path is concise; the Processor API path gives explicit lifecycle and naming control.
go deeper
Recognize Materialized.as(name) configures a DSL store and Stores is the factory for explicit ones.
Use Materialized to set serde/logging/caching, and build a Processor-API store with Stores + StoreBuilder + addStateStore.
Choose store types deliberately, connect stores to processors, and know which knobs live on Materialized vs StoreBuilder.
Design store topology for shared access, queryability, and the DSL↔PAPI interop boundary across a team's codebase.
## Two abstraction levels Kafka Streams offers two APIs. The high-level **DSL** (`KStream`, `KTable`, `count`, `aggregate`, `join`) auto-creates state stores for you and lets you tune them with a **`Materialized`** object. The low-level **Processor API** (custom `Processor`/`Transformer`) requires you to build and register stores explicitly via the **`Stores`** factory and a **`StoreBuilder`**. ## DSL: Materialized `Materialized<K, V, S extends StateStore>` is the configuration handle passed to stateful DSL operators: - **Naming & queryability:** `Materialized.as("my-store")` names the store so it can be looked up for **Interactive Queries**. Without a name, Streams generates an internal name and the store is not queryable. - **Serdes:** `.withKeySerde(...)`, `.withValueSerde(...)` — required when the store's types differ from the stream's defaults, because the store and its changelog must serialize entries. - **Fault tolerance:** `.withLoggingEnabled(Map)` / `.withLoggingDisabled()` controls the changelog. - **Caching:** `.withCachingEnabled()` / `.withCachingDisabled()` controls the record cache (dedup of updates before forwarding). - **Store type:** by default persistent RocksDB; pass a supplier to change it, e.g. `Materialized.as(Stores.inMemoryKeyValueStore("x"))`. - **Windowed retention:** `.withRetention(Duration)` for window/session stores. ## Processor API: Stores + StoreBuilder Here you assemble the store yourself, in three steps: 1. **Choose a store supplier** from the `Stores` factory: - `Stores.persistentKeyValueStore(name)` — RocksDB key-value. - `Stores.inMemoryKeyValueStore(name)` — heap key-value. - `Stores.persistentWindowStore(name, retention, windowSize, retainDuplicates)` — windowed. - `Stores.persistentSessionStore(name, retention)` — session. - `Stores.persistentTimestampedKeyValueStore(name)` — value plus timestamp. 2. **Wrap it in a StoreBuilder** with serdes: `Stores.keyValueStoreBuilder(supplier, keySerde, valueSerde)` (or `windowStoreBuilder`, `sessionStoreBuilder`). On the builder you call `.withLoggingEnabled(configs)` / `.withLoggingDisabled()` and `.withCachingEnabled()` / `.withCachingDisabled()`. 3. **Register and connect:** `topology.addStateStore(builder, "my-processor")` (or `streamsBuilder.addStateStore(builder)` then reference the store name in `process(...)`). Inside the processor, retrieve it via `context.getStateStore("name")`. ## Why two APIs The DSL trades control for brevity — it picks store types, names, and changelog behavior with sensible defaults. The Processor API is for custom processing logic that the DSL can't express, where you need explicit control over which processors share which stores, the exact store type, and lifecycle. They interoperate: you can drop into the Processor API from the DSL with `process`/`transform` and reference DSL-or-PAPI-registered stores by name. ## Common mistakes - Forgetting `.withKeySerde/.withValueSerde` when store types differ from defaults → `ClassCastException`/serde errors at runtime. - Not naming a Materialized store you intend to query interactively. - Registering a `StoreBuilder` but never connecting it to a processor (the store exists but no processor can access it).
- What does naming a store via Materialized.as("...") enable beyond readability?A named store is exposed for Interactive Queries — you can look it up with KafkaStreams.store(...) and read state directly from the running app. Unnamed stores get an internal generated name and are not (reliably) queryable.
- How do you switch a DSL aggregation from RocksDB to an in-memory store?Pass an in-memory supplier into Materialized, e.g. count(Materialized.as(Stores.inMemoryKeyValueStore("counts"))). The aggregation then materializes into a heap-backed store instead of RocksDB.
saying these in an interview costs you the question
- Thinking Materialized and StoreBuilder are interchangeable across both APIs — Materialized is DSL, StoreBuilder is Processor API
- Forgetting that store types differing from stream defaults require explicit serdes
- Believing every Materialized store is automatically interactive-queryable even without a name
- Registering a StoreBuilder but never connecting it to a processor