skip to content

Interactive Queries

Serving reads directly from a Streams app's local state and routing a query to whichever instance owns the key. A strong design question about turning a stream processor into a queryable service.

part ofApache Kafkaoverview, primer and where to startread it →
on this pageshow

questions

5

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

open as a page

In a multi-instance Kafka Streams app, how do you find which instance can answer an Interactive Query for a given key?

level: middleimportance: must knowfreq 60%

basics

~20 s

A local store only holds the partitions assigned to its instance. To find the owner of a key, call streams.queryMetadataForKey(storeName, key, serializer). It returns the host:port (from application.server) of the instance that owns that key's partition, so you can route the query there.

open as a page

How would you build the RPC layer that turns local Interactive Queries into a cluster-wide queryable service?

level: seniorimportance: should knowfreq 45%

basics

~20 s

Run an HTTP (or gRPC) server on each instance with a query endpoint. On a request, use queryMetadataForKey to find the owning host. If it's local, read the local store and return the value; if remote, forward the request to that host's endpoint and relay the result.

open as a page

How do standby replicas (num.standby.replicas) change Interactive Query availability and routing, and what are the consistency tradeoffs?

level: seniorimportance: should knowfreq 35%

basics

~20 s

Setting num.standby.replicas > 0 makes other instances keep warm copies of a store's state. If the active instance fails, queryMetadataForKey still lists standby hosts, so you can serve reads from a standby for higher availability — but a standby may be slightly behind, so its data can be stale.

open as a page

What is the IQv2 API (KIP-796), and how do you interactively query windowed/session stores versus key-value stores?

level: principalimportance: nice to knowfreq 25%

basics

~20 s

IQv2 is a newer, type-safe query API (StateQueryRequest + Query types like KeyQuery/RangeQuery) returning a StateQueryResult, designed to be extensible and partition-aware. Windowed stores use a different queryable type (windowStore/sessionStore) and you query by key plus a time range, not just a key.

open as a page