skip to content

How do you wire a materialized KTable state store in Spring and expose it for interactive queries?

level: principalimportance: nice to knowfreq 22%

answer

  1. Materialized.as("store") -> named RocksDB store + compacted changelog
  2. getKafkaStreams().store(StoreQueryParameters...)
  3. KafkaStreamsInteractiveQueryService (3.2+) wraps retries + host lookup
  4. APPLICATION_SERVER_CONFIG for cross-instance metadataForKey
  5. InvalidStateStoreException during REBALANCING / before RUNNING

basics

~10 s

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

A 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 lines
java
import 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

for a junior

Know a KTable can be materialized into a named store you can read back.

for a middle

Explain Materialized.as, the RocksDB store + changelog, and reading via getKafkaStreams().store().

for a senior

Cover KafkaStreamsInteractiveQueryService, InvalidStateStoreException timing, and store-type matching.

for a principal

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)

context