In a multi-instance Kafka Streams app, how do you find which instance can answer an Interactive Query for a given key?
answer
- queryMetadataForKey(store, key, serializer) → KeyQueryMetadata
- activeHost / standbyHosts / partition
- application.server=host:port advertises your RPC address
- streamsMetadataForStore → all hosts (scatter-gather)
- replaced old metadataForKey
basics
~20 sA 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.
solid answer
~40 sBecause each instance materializes only its assigned partitions, a key's value may live on another instance. Kafka Streams gossips placement metadata via the consumer group, exposed through KafkaStreams.queryMetadataForKey(store, key, keySerializer), which returns a KeyQueryMetadata containing the active host (HostInfo) and standby hosts for that key's partition. To make HostInfo meaningful you must set application.server=host:port on each instance — that's the address peers advertise. For whole-store discovery you use streamsMetadataForStore(store), which returns one StreamsMetadata per instance hosting that store. The typical pattern: compute partition/owner with queryMetadataForKey; if it's the local host, read the local store; otherwise issue an HTTP/gRPC call to the remote host's IQ endpoint. queryMetadataForKey replaced the older metadataForKey API.
code
java · 9 linesKeyQueryMetadata md = streams.queryMetadataForKey(
"word-counts", "hello", Serdes.String().serializer());
HostInfo owner = md.activeHost();
if (owner.equals(thisHostInfo)) {
return streams.store(StoreQueryParameters.fromNameAndType(
"word-counts", QueryableStoreTypes.keyValueStore())).get("hello");
} else {
return rpcClient.get(owner.host(), owner.port(), "word-counts", "hello");
}go deeper
Know that a key may live on another instance and that an API tells you which host owns it.
Use queryMetadataForKey to get the active HostInfo, set application.server, and route local-vs-remote.
Contrast queryMetadataForKey vs streamsMetadataForStore, handle standbys, custom partitioners, and stale-metadata-during-rebalance.
Design the full routing/RPC layer and reason about availability windows during rebalances and metadata propagation delay.
**The problem.** With N instances and a P-partition input topic, the state store is sharded: instance i materializes whichever partitions it was assigned. A query for `key` can only be answered locally if `key`'s partition lives on this instance. So IQ across a cluster needs two steps: (1) *locate* the owning instance, (2) *route* the query there. This leaf is about step 1. **How placement is known.** Kafka Streams instances share metadata through the consumer group protocol — each instance knows the full assignment of partitions to members. That metadata is surfaced through the IQ metadata APIs so any instance can answer "who owns this key?" without an external coordinator. **`application.server`.** For the answer to be *actionable*, each instance must advertise a reachable address via the config `application.server=<host>:<port>` (e.g. `app-1.internal:8080`). This is the host/port of *your* RPC layer, not Kafka. Kafka Streams stores it in the metadata and hands it back as a `HostInfo`. If you don't set it, metadata still works but `HostInfo` values are meaningless (you can't route). **Locating a key — `queryMetadataForKey`:** ``` KeyQueryMetadata md = streams.queryMetadataForKey( "word-counts", "hello", Serdes.String().serializer()); HostInfo active = md.activeHost(); // owner of the key's partition Set<HostInfo> standbys = md.standbyHosts(); // standby replicas, if configured int partition = md.partition(); ``` `queryMetadataForKey` serializes the key with the provided serializer, computes its partition with the default partitioner (or a supplied `StreamPartitioner`), and returns a `KeyQueryMetadata` with the **active** host, any **standby** hosts, and the partition number. This is the modern API; the older `metadataForKey(...)` returned a single `StreamsMetadata` and predates standby-aware routing — interviewers like to hear you name the newer one. **Discovering all hosts — `streamsMetadataForStore`:** ``` Collection<StreamsMetadata> all = streams.streamsMetadataForStore("word-counts"); ``` Returns one `StreamsMetadata` per instance that hosts the store, each with its `HostInfo`, the topic-partitions it owns, and the store names. Use it for `range`/`all` scatter-gather queries (which span every partition) or to build a routing table. **Routing pattern.** The canonical flow: 1. `queryMetadataForKey(store, key, serializer)` → `HostInfo owner`. 2. If `owner.equals(thisHostInfo)` → read the **local** store via `streams.store(...)`. 3. Else → make an RPC (HTTP/gRPC) to `owner.host():owner.port()` hitting that instance's IQ endpoint, which does the local read and returns the value. This turns the cluster into a transparent distributed key-value store. **Edge cases.** - During a rebalance, ownership shifts; metadata can briefly be stale, and a routed query may hit an instance that no longer owns the partition → it should return a retriable error / `InvalidStateStoreException` so the caller retries. - `queryMetadataForKey` can return `KeyQueryMetadata.NOT_AVAILABLE` when metadata isn't ready (e.g. before first assignment). - A custom `StreamPartitioner` used in the topology must be passed to `queryMetadataForKey` too, or the computed partition won't match where data actually landed.
- What config must be set for the returned HostInfo to be useful, and what does it represent?application.server=host:port on each instance. It advertises the address of your own RPC/IQ endpoint (not a Kafka broker), which Streams stores in metadata and returns as HostInfo so peers can route queries to you.
- When would you use streamsMetadataForStore instead of queryMetadataForKey?For queries that aren't keyed to a single partition — range scans or all() — which must hit every instance hosting the store. streamsMetadataForStore lists all those hosts so you can scatter-gather and merge results, or to build a routing table at startup.
- Why must a custom StreamPartitioner be passed to queryMetadataForKey?Routing computes the key's partition; if the topology used a custom partitioner, the default one would compute a different partition than where the data actually landed, sending the query to the wrong instance.
saying these in an interview costs you the question
- Saying you can query any instance for any key without routing
- Forgetting application.server, so HostInfo is unusable
- Confusing application.server with bootstrap.servers / a broker address
- Not handling stale metadata during rebalances
- Using a default partition computation when the topology has a custom StreamPartitioner