skip to content

What are Interactive Queries in Kafka Streams, and how do you read the value for a key from a local state store?

level: juniorimportance: must knowfreq 70%

answer

  1. Materialized.as("name") → queryable
  2. streams.store(StoreQueryParameters.fromNameAndType)
  3. ReadOnlyKeyValueStore: get/range/all, no put
  4. only local partitions visible
  5. InvalidStateStoreException before RUNNING / during rebalance

basics

~20 s

Interactive Queries let your app read directly from the state stores Kafka Streams already maintains, instead of querying an external database. You call KafkaStreams.store(...) to get a read-only store and look up a key with get().

solid answer

~40 s

Interactive Queries (IQ) expose the materialized state stores that Kafka Streams keeps for aggregations, KTables, and windowed operations so an application can serve reads from them directly rather than from an external datastore. You first materialize a store with a name (e.g. Materialized.as("counts")), then after the app reaches RUNNING you fetch a handle via streams.store(StoreQueryParameters.fromNameAndType("counts", QueryableStoreTypes.keyValueStore())). That returns a ReadOnlyKeyValueStore<K,V> whose get(key), range(from,to), and all() methods read the local RocksDB/in-memory state. The store is read-only by design — IQ never mutates state; writes only happen through the topology. A store handle only sees the partitions assigned to this instance, so a single instance answers only for the keys it owns.

code

java · 4 lines
java
ReadOnlyKeyValueStore<String, Long> store =
    streams.store(StoreQueryParameters.fromNameAndType(
        "word-counts", QueryableStoreTypes.keyValueStore()));
Long count = store.get("hello");   // null if this instance does not own the key's partition

go deeper

for a junior

Know IQ reads existing state stores, you name the store via Materialized.as, and the handle is read-only (get/range/all).

for a middle

Explain the StoreQueryParameters.fromNameAndType + QueryableStoreTypes pairing and that handles only see local partitions.

for a senior

Discuss lifecycle (RUNNING required), InvalidStateStoreException on rebalance, dynamic rebinding, and why locality forces a distributed design.

for a principal

Frame IQ as removing an external read store from the architecture and the tradeoffs that creates (operational coupling of reads to the streams cluster, rebalance-time availability).

**State stores.** Kafka Streams maintains local state to do stateful operations — counts, aggregations, joins, KTable materializations. That state lives in a *state store*, by default a RocksDB instance on local disk (or an in-memory store), and is backed for fault tolerance by a compacted *changelog topic* in Kafka. Normally this state is internal plumbing the topology reads and writes as records flow through. **Interactive Queries (IQ)** make that internal state *queryable from outside the processing loop*. Instead of writing your aggregation result back to a Kafka topic and loading it into Postgres/Redis to serve lookups, you query the store the stream already maintains. This removes a whole external system from the read path. **Making a store queryable.** You must give the store a name so IQ can find it. For a DSL aggregation: `.count(Materialized.as("word-counts"))`. For the Processor API you register a `StoreBuilder` with a name. Only named, materialized stores are queryable. **Getting a handle.** After `streams.start()` and once the instance is in the `RUNNING` state, call: ``` ReadOnlyKeyValueStore<String, Long> store = streams.store(StoreQueryParameters.fromNameAndType( "word-counts", QueryableStoreTypes.keyValueStore())); ``` `StoreQueryParameters.fromNameAndType` pairs the store name with a `QueryableStoreType` describing the shape: `keyValueStore()`, `windowStore()`, `sessionStore()`, `timestampedKeyValueStore()`, etc. The returned object is **read-only** — `ReadOnlyKeyValueStore` exposes `get(key)`, `range(from, to)`, `reverseRange`, `all()`, `reverseAll()`, and `approximateNumEntries()`, but no `put`/`delete`. That is intentional: IQ must never corrupt the state the topology owns; the changelog and processing order are the only source of writes. **Locality — the key edge case.** A store handle reflects **only the partitions assigned to this instance**. If your input topic has 6 partitions spread across 3 instances, each instance materializes ~2 partitions of the store. `get(key)` on instance A returns a value only if `key` hashes to a partition A owns; otherwise it returns `null` even though the key exists elsewhere. So IQ on a multi-instance app is inherently a *distributed* problem — locating which instance owns a key (covered by the metadata APIs) is the other half of the feature. **Lifecycle / exceptions.** Calling `store(...)` before the app is `RUNNING`, or during a rebalance when the store is being (re)assigned, throws `InvalidStateStoreException`. Robust code retries on that exception. The handle is also *dynamic*: it transparently rebinds to the underlying store across rebalances, but a query landing mid-rebalance can still throw and should be retried.

  • Why is the returned store type read-only?
    IQ must never mutate state the topology owns; all writes go through the processing pipeline and changelog, which preserves ordering and fault tolerance. A read-only handle prevents queries from corrupting state or diverging from the changelog.
  • Why might get(key) return null for a key you know exists?
    Because a local store handle only sees partitions assigned to this instance. If the key hashes to a partition owned by another instance, the local lookup returns null — you must locate the owning instance via the metadata APIs and route the query there.

saying these in an interview costs you the question

  • Claiming you can put/delete through an IQ store handle
  • Thinking a single instance can answer for every key in the topic
  • Querying the store before the app reaches RUNNING and not handling InvalidStateStoreException
  • Confusing IQ (reading internal state) with consuming the changelog topic

context