skip to content

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

level: seniorimportance: should knowfreq 35%

answer

  1. num.standby.replicas=K → K warm changelog copies
  2. KeyQueryMetadata.standbyHosts() for fallback
  3. enableStaleStores() = opt-in to possibly-stale reads
  4. fast failover (skip cold restore) + availability vs consistency
  5. extra disk/bandwidth; KIP-441 warmup / acceptable.recovery.lag

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.

solid answer

~40 s

num.standby.replicas controls how many extra instances replicate each store's changelog into a local copy. With standbys, KeyQueryMetadata from queryMetadataForKey exposes both activeHost() and standbyHosts(). During an active instance's downtime or a long restore, your RPC layer can fall back to a standby and keep serving reads instead of returning errors — trading consistency for availability, because a standby's local state may lag the active by some changelog offsets. Standbys also dramatically shorten failover: the new active is promoted from an already-warm copy rather than restoring the whole changelog from scratch. You typically combine this with StoreQueryParameters.enableStaleStores() so a promoted-but-still-restoring or standby store can be queried while it's catching up. The decision is explicit: only enable standby reads where bounded staleness is acceptable.

code

java · 5 lines
java
StoreQueryParameters<ReadOnlyKeyValueStore<String, Long>> params =
    StoreQueryParameters
        .fromNameAndType("word-counts", QueryableStoreTypes.keyValueStore())
        .enableStaleStores();   // allow standby / restoring reads
ReadOnlyKeyValueStore<String, Long> store = streams.store(params);

go deeper

for a junior

Know standbys are warm backup copies that help availability if an instance fails.

for a middle

Know num.standby.replicas, that metadata exposes standbyHosts(), and that standby reads can be stale.

for a senior

Use enableStaleStores(), design active-then-standby fallback, and articulate the availability-vs-consistency tradeoff.

for a principal

Weigh standby count vs cost, KIP-441 warmup/acceptable.recovery.lag interactions, and per-query freshness policy across the platform.

**What a standby is.** By default Kafka Streams keeps one *active* copy of each store partition (on the instance assigned that partition). `num.standby.replicas=K` tells Streams to keep K additional instances continuously consuming the store's changelog topic into their own local RocksDB copy. These standbys never serve the topology's processing — they exist purely as warm replicas. **Why they matter for IQ.** Two benefits: 1. **Fast failover.** If the active instance dies, one standby is promoted to active. Because it already has a near-current copy, it skips the potentially huge changelog *restore* that a cold instance would need (which can take minutes for large state), so the partition becomes available again quickly. 2. **Read availability during gaps.** While the active is down/restoring, the active store is unqueryable. `queryMetadataForKey` returns a `KeyQueryMetadata` whose `standbyHosts()` lists the warm replicas, so your routing layer can serve the read from a standby instead of erroring. **The consistency catch.** A standby tails the changelog asynchronously, so its local state can lag the active by some offsets. Reading from it gives **eventually-consistent / possibly-stale** results. This is an availability-vs-consistency choice you make explicitly per query type: dashboards/approximate counts tolerate it; 'read your writes' semantics do not. **`enableStaleStores()`.** Normally `streams.store(...)` only returns a handle for *active, fully-restored* stores; querying a restoring or standby store throws. To deliberately allow reads from stores that may be stale (standbys, or an active mid-restore), build params with: ``` StoreQueryParameters.fromNameAndType(name, type).enableStaleStores(); ``` This is the explicit opt-in that says 'I accept possibly-stale data in exchange for availability.' Without it, you only ever read fully-caught-up active stores. **Routing logic with standbys.** A robust handler: 1. `queryMetadataForKey` → active + standby hosts. 2. Try the active first (fresh data). 3. If active is unreachable / `InvalidStateStoreException` / not RUNNING, optionally fall back to a `standbyHosts()` member, using `enableStaleStores()` on the remote read. 4. Surface in the response whether the answer came from active or standby (e.g. a header) so callers know the freshness. **Operational notes.** - Standbys cost extra disk and changelog-consumption bandwidth on the replica instances — `K` replicas means ~`K+1`× the state footprint. - `acceptable.recovery.lag` and warmup-replica settings (KIP-441 'smooth' rebalancing) govern when a warm replica is considered caught-up enough to take over, which interacts with how quickly standbys become useful after scaling. - Standbys don't help if *all* copies are simultaneously rebalancing; they reduce, not eliminate, IQ unavailability windows.

  • What does StoreQueryParameters.enableStaleStores() do and when do you need it?
    It allows querying stores that aren't fully caught up — standby replicas or an active that's still restoring. Without it, store(...) only returns active, fully-restored stores and otherwise throws. You enable it to favor availability over freshness.
  • Why can a standby return stale data?
    A standby asynchronously tails the store's changelog topic, so its local copy may lag the active by some offsets at any instant. Reads from it are eventually consistent and can miss the most recent updates.
  • Beyond IQ availability, what is the main operational benefit of standbys?
    Fast failover: a promoted standby already holds a near-current copy, so it avoids a full changelog restore (which can take minutes for large state), shrinking the window the partition is unavailable.

saying these in an interview costs you the question

  • Claiming standby reads are always consistent with the active
  • Forgetting enableStaleStores() is required to read standby/restoring stores
  • Thinking standbys serve the processing topology (they don't — they're warm replicas)
  • Ignoring the extra disk/changelog cost of K replicas
  • Assuming standbys eliminate, rather than reduce, IQ downtime

context