How do you wire a materialized KTable state store in Spring and expose it for interactive queries?
answer
- Materialized.as("store") -> named RocksDB store + compacted changelog
- getKafkaStreams().store(StoreQueryParameters...)
- KafkaStreamsInteractiveQueryService (3.2+) wraps retries + host lookup
- APPLICATION_SERVER_CONFIG for cross-instance metadataForKey
- InvalidStateStoreException during REBALANCING / before RUNNING
basics
~10 sBuild a KTable with Materialized.as("store-name") so Kafka Streams keeps a named state store. Then reach the running client via StreamsBuilderFactoryBean.getKafkaStreams() (or KafkaStreamsInteractiveQueryService) and call store(...) to read current values by key.
solid answer
~40 sA KTable is Kafka Streams' materialized 'latest value per key' view; when you name its store with Materialized.as("store"), Kafka Streams maintains a local state store (RocksDB by default) backed by a compacted changelog topic for fault tolerance. In Spring you define this in a topology @Bean on the injected StreamsBuilder. To serve reads, you need the running client: inject the StreamsBuilderFactoryBean and call getKafkaStreams().store(StoreQueryParameters.fromNameAndType("store", QueryableStoreTypes.keyValueStore())). Spring for Kafka 3.2+ adds KafkaStreamsInteractiveQueryService, which wraps this and also resolves the host that owns a key across a multi-instance app (via metadataForKey), so you can proxy remote reads. Queries must run only when the client is RUNNING; querying during REBALANCING throws InvalidStateStoreException, so guard with retries or the service's helpers.
code
java · 41 linesimport org.apache.kafka.streams.StreamsBuilder;
import org.apache.kafka.streams.kstream.KTable;
import org.apache.kafka.streams.kstream.Materialized;
import org.apache.kafka.streams.state.QueryableStoreTypes;
import org.apache.kafka.streams.state.ReadOnlyKeyValueStore;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.kafka.annotation.EnableKafkaStreams;
import org.springframework.kafka.streams.KafkaStreamsInteractiveQueryService;
import org.springframework.stereotype.Service;
@Configuration
@EnableKafkaStreams
class CountsTopology {
@Bean
KTable<String, Long> counts(StreamsBuilder builder) {
return builder.<String, String>stream("events")
.groupByKey()
.count(Materialized.as("counts-store"));
}
// Spring for Kafka 3.2+: register the query service (or inject the factory bean directly).
@Bean
KafkaStreamsInteractiveQueryService iqService(
org.springframework.kafka.config.StreamsBuilderFactoryBean fb) {
return new KafkaStreamsInteractiveQueryService(fb);
}
}
@Service
class CountQuery {
private final KafkaStreamsInteractiveQueryService iq;
CountQuery(KafkaStreamsInteractiveQueryService iq) { this.iq = iq; }
Long countFor(String key) {
ReadOnlyKeyValueStore<String, Long> store =
iq.retrieveQueryableStore("counts-store", QueryableStoreTypes.keyValueStore());
return store.get(key); // may be owned by another instance in a scaled-out app
}
}go deeper
Know a KTable can be materialized into a named store you can read back.
Explain Materialized.as, the RocksDB store + changelog, and reading via getKafkaStreams().store().
Cover KafkaStreamsInteractiveQueryService, InvalidStateStoreException timing, and store-type matching.
Design multi-instance interactive queries: APPLICATION_SERVER_CONFIG, metadataForKey proxying, standby replicas, query-availability SLAs.
**KTable and state stores (wiring-level, since library internals are a sibling topic).** A `KTable<K,V>` represents the current value per key. When you materialize it — `builder.table("topic", Materialized.as("my-store"))` or `.aggregate(..., Materialized.as("my-store"))` — Kafka Streams keeps a **local state store** (default: RocksDB on local disk under `state.dir`) holding those values, and backs it with a **compacted changelog topic** so the store can be rebuilt after a crash or rebalance. The store name you pass is the handle you'll query. **Defining it in Spring.** It's still a normal topology bean: ```java @Bean public KTable<String, Long> counts(StreamsBuilder builder) { return builder.stream("events") .groupByKey() .count(Materialized.as("counts-store")); } ``` The `StreamsBuilderFactoryBean` builds and starts the client as usual; the store comes online once the client reaches `RUNNING`. **Interactive queries — the classic way.** 'Interactive query' = reading a state store directly from your service (no extra Kafka consumer). You need the live `KafkaStreams`: ```java ReadOnlyKeyValueStore<String, Long> store = factoryBean.getKafkaStreams().store( StoreQueryParameters.fromNameAndType("counts-store", QueryableStoreTypes.keyValueStore())); Long v = store.get("key-1"); ``` Inject the `StreamsBuilderFactoryBean` (by type, or the reserved name `defaultKafkaStreamsBuilder`) to reach `getKafkaStreams()`. **Interactive queries — the Spring helper (3.2+).** `KafkaStreamsInteractiveQueryService` wraps the boilerplate: `retrieveQueryableStore(name, type)` returns the store with built-in retry while the client is still stabilizing, and it exposes `getKafkaStreamsApplicationMetadata` / host-info lookups. In a **multi-instance** deployment each instance owns only *some* partitions' keys; the service uses Kafka Streams' `metadataForKey(store, key, serializer)` to tell you *which host* holds a given key so you can HTTP-proxy the read to the owning instance. To enable cross-host discovery you must set `StreamsConfig.APPLICATION_SERVER_CONFIG` (host:port) so instances advertise themselves. **Lifecycle & timing gotchas.** 1. **`InvalidStateStoreException`.** Querying before the client is `RUNNING`, or during a `REBALANCING`, throws this. Guard with the service's retry, a `StateListener` gating readiness, or catch-and-retry. 2. **`getKafkaStreams()` returns null before start.** Don't query during context init. 3. **Store type must match.** Key-value vs. windowed vs. session stores use different `QueryableStoreTypes`; a mismatch fails. 4. **Locality.** A single instance can only answer for the keys/partitions it currently owns. Ignoring `metadataForKey` in a scaled-out app yields wrong 'not found' answers for keys owned elsewhere. 5. **Changelog cost.** Every materialized store creates an internal compacted changelog topic; naming stores (vs. anonymous) keeps those topic names stable across restarts and avoids orphaned topics. 6. **Standby replicas.** `NUM_STANDBY_REPLICAS_CONFIG > 0` lets other instances serve stale reads during rebalance — relevant for query availability SLAs. **When to use.** Interactive queries suit low-latency point lookups of aggregated/materialized state (dashboards, enrichment lookups) served directly from the streams app, avoiding a separate database. If you need rich ad-hoc querying, push results to an external store via `.to(...)` instead. **Multi-topology / multi-app note.** Serving several independent applications means several named `StreamsBuilderFactoryBean`s, each with its own `application.id` and `application.server`; you'd inject the specific factory bean (by qualifier) to query the right one.
- In a 3-instance deployment, a key returns null on one instance but exists. Why, and how do you fix it?Each instance only holds the partitions (and thus keys) it owns. Set APPLICATION_SERVER_CONFIG so instances advertise host:port, then use metadataForKey (or KafkaStreamsInteractiveQueryService host lookup) to find and proxy the read to the owning instance.
- Why might a query throw InvalidStateStoreException right after startup?The store isn't queryable until the KafkaStreams client reaches RUNNING and finishes restoring/rebalancing. Guard with retries (the interactive query service does this) or gate on the RUNNING state via a StateListener.
saying these in an interview costs you the question
- Assuming a single instance can answer queries for all keys in a scaled-out app
- Querying getKafkaStreams() during context init before the client exists/RUNNING
- Thinking a materialized store needs no changelog topic (it always creates one)