skip to content

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

level: seniorimportance: should knowfreq 45%

answer

  1. RPC server bound to application.server host:port
  2. queryMetadataForKey → local read OR proxy
  3. ?local=true hop guard, no forwarding loops
  4. range/all → scatter-gather via streamsMetadataForStore
  5. 503 not-ready vs 404 not-found; retry on InvalidStateStoreException

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.

solid answer

~40 s

Each instance hosts an RPC server bound to the same host:port it advertises via application.server. A query endpoint (e.g. GET /store/{name}/{key}) calls queryMetadataForKey(store, key, serializer); if the active HostInfo equals the local host it reads the local ReadOnlyKeyValueStore and returns the value, otherwise it proxies the call to the owning host's endpoint (with a hop guard to prevent infinite forwarding). For range/all queries you can't pinpoint one partition, so you use streamsMetadataForStore to scatter the request to every host and merge results. The layer must: handle InvalidStateStoreException with retries (rebalances), distinguish 'key not found' from 'not ready', expose the app's RUNNING state for readiness, and optionally fall back to standby hosts when the active is unavailable. This is the pattern Confluent's kafka-streams-interactive-queries example implements.

code

java · 9 lines
java
// GET /store/{store}/{key}
KeyQueryMetadata md = streams.queryMetadataForKey(store, key, serializer);
if (md.equals(KeyQueryMetadata.NOT_AVAILABLE)) return status(503);
if (md.activeHost().equals(thisHost) || localOnly) {
    Object v = streams.store(StoreQueryParameters.fromNameAndType(
            store, QueryableStoreTypes.keyValueStore())).get(key);
    return v == null ? status(404) : ok(v);
}
return proxy(md.activeHost(), "/store/" + store + "/" + key + "?local=true");

go deeper

for a junior

Understand that an HTTP endpoint per instance plus 'local read or forward' makes IQ work cluster-wide.

for a middle

Implement the point-lookup route with queryMetadataForKey and local-vs-proxy logic.

for a senior

Add scatter-gather for range/all, hop guards, readiness, and correct 404-vs-503 semantics with retries.

for a principal

Reason about availability vs consistency (standby fallback), operational coupling of reads to streams health, and rebalance blast radius.

**Goal.** IQ gives you *local* reads; the metadata APIs tell you *who* owns a key. The RPC layer is the glue that makes the cluster behave like one distributed read-only KV store, so a client can hit *any* instance and get the right answer. **Topology of the service.** Each Streams instance also runs a lightweight web server — Spring MVC, Javalin, Vert.x, or gRPC — bound to exactly the `host:port` it advertises in `application.server`. That symmetry is essential: the address peers route to *is* the address the server listens on. **Point lookup endpoint.** ``` GET /store/{store}/{key} ``` Handler: 1. `KeyQueryMetadata md = streams.queryMetadataForKey(store, key, serializer);` 2. If `md == NOT_AVAILABLE` → 503 (metadata not ready, retry). 3. If `md.activeHost().equals(thisHost)` → local read: `streams.store(...).get(key)`; return 200 with value or 404 if null. 4. Else → **proxy** to `http://activeHost.host():activeHost.port()/store/{store}/{key}`. Add a query flag like `?local=true` (a *hop guard*) so the forwarded request is served locally and never re-forwarded — this prevents forwarding loops during stale-metadata windows. **Range / all (scatter-gather).** `range(from,to)` and `all()` span multiple partitions, so no single host has the full answer. Use `streamsMetadataForStore(store)` to enumerate every host, query each (locally on self, via RPC on others), and merge/sort the partial results. This is fan-out and far costlier than a point lookup. **Reliability concerns.** - **Rebalances:** during reassignment, `store(...)` and local reads can throw `InvalidStateStoreException`, and routed requests may hit a host that just lost the partition. Wrap reads in a bounded retry; return a retriable status (503) so clients back off. - **Readiness:** expose `streams.state() == RUNNING` (or `REBALANCING`) on a `/health` or readiness probe; don't route to instances that aren't ready. - **Standbys:** if active is unreachable, you may serve a (possibly slightly stale) read from `md.standbyHosts()` — a deliberate availability-over-consistency choice; gate it behind a flag. - **404 vs 503:** carefully separate "key genuinely absent" (404) from "store not currently available / metadata not ready" (503). Collapsing them makes clients retry on real misses or give up on transient errors. - **Serialization:** the RPC layer must agree on serdes; the key serializer passed to `queryMetadataForKey` must match the store's key serde so partition computation is correct. **Why this design.** It keeps reads co-located with state (no external DB), scales horizontally with the streams cluster, and tolerates the cluster's own rebalancing semantics. The cost is operational coupling: query availability is now tied to streams health and rebalance frequency, which informs whether you also keep standby replicas (see standby-store routing) and how aggressive your retry/back-off is. Confluent ships a canonical `kafka-streams-interactive-queries` reference app implementing exactly this skeleton.

  • How do you prevent infinite request forwarding between instances?
    Add a 'local-only' flag (e.g. ?local=true) on the proxied request so the receiving instance serves it from its local store and never re-routes, even if its metadata view is momentarily stale. Optionally cap a hop count too.
  • How do range() and all() queries differ from a point lookup in this layer?
    They span all partitions, so no single host answers them. You enumerate hosts with streamsMetadataForStore, fan the query out to every instance, and merge the partial results — a scatter-gather that's far more expensive than a single-key route.
  • How should the endpoint distinguish a missing key from an unavailable store?
    Return 404 when a local get() returns null (key genuinely absent) and 503/retriable when metadata is NOT_AVAILABLE or a read throws InvalidStateStoreException (transient, during rebalance). Conflating them breaks client retry logic.

saying these in an interview costs you the question

  • No hop guard, so proxied requests can loop during stale metadata
  • Treating range/all like a single-host point lookup
  • Returning the same status for missing-key and store-unavailable
  • Binding the RPC server to a port different from the advertised application.server
  • No retry/back-off around rebalance-time InvalidStateStoreException

context