skip to content

What does num.standby.replicas do in Kafka Streams, and how does it affect failover and availability?

level: seniorimportance: must knowfreq 55%

answer

  1. num.standby.replicas default 0
  2. hot copy fed by changelog topic
  3. fast failover, not throughput
  4. needs N+1 distinct instances
  5. warmup (scaling) vs standby (failover)

basics

~20 s

num.standby.replicas keeps shadow copies of a stateful task's state store on other instances, continuously fed from the changelog topic. On failover, an up-to-date standby takes over almost immediately instead of rebuilding state from scratch, cutting recovery time.

solid answer

~50 s

num.standby.replicas (default 0) sets how many hot standby copies Kafka Streams maintains for each stateful task. A standby replica lives on a different instance and continuously consumes the task's changelog topic, keeping a near-real-time copy of the local state store (RocksDB). If the instance running the active task dies, the assignor promotes a standby — which is already nearly caught up — so the new active task only needs to replay the small remaining changelog tail instead of restoring the entire store from the beginning. This drastically shortens the recovery window and improves availability for stateful applications. The cost is extra storage, memory, and changelog-consumption bandwidth on the standby instances, and you need at least num.standby.replicas+1 instances (or threads spread across instances) for standbys to actually be placed. Standbys are about fast failover; they do not add processing throughput.

go deeper

for a junior

Know num.standby.replicas keeps a backup copy of state for faster recovery; default 0.

for a middle

Explain that the standby follows the changelog and shortens restore time on failover, at storage cost.

for a senior

Distinguish standby vs warmup replicas, placement requirements, and the throughput-vs-availability separation.

for a principal

Design availability targets: rack-aware placement, recovery-lag tuning, and the resource budget for N copies of state.

## The problem standbys solve A **stateful** Streams task (aggregation, join, windowed op) keeps its working state in a local **state store**, typically **RocksDB** on the instance's disk. Every update is also written to a Kafka **changelog topic** (a compacted topic that is the durable source of truth for the store). If the instance hosting that task crashes, the task must be reassigned. Without standbys, the new host has an **empty** store and must **restore** it by replaying the entire changelog from the beginning (or last checkpoint). For a large store this can take minutes — during which the task (and the keys it owns) is unavailable. ## What a standby replica is `num.standby.replicas = N` tells Streams to keep **N additional hot copies** of each stateful task's store on **other** instances. A standby: - Does **not** process input or produce output. - **Continuously consumes the changelog topic** for that task, applying updates to its own local RocksDB copy. - Stays **nearly in sync** with the active task's store. ## Failover behavior When the active instance fails, the partition assignor prefers to **promote a standby** that already holds the state. Because the standby is almost caught up, the new active task only needs to consume the **small tail** of changelog records produced since the standby last polled — recovery goes from minutes to (often) sub-second to seconds. This is the 'hot failover' the coverage area refers to. ## Interaction with warmup replicas (KIP-441) Distinct concept: **warmup replicas** are temporary copies created during **scaling/rebalancing** to pre-restore a task that is about to move, so the move doesn't stall processing. Standby replicas are **permanent** copies for **failover**. They share the same restoration machinery but serve different goals; `max.warmup.replicas` caps concurrent warmups. ## Costs and constraints - **Resources**: each standby costs disk + memory (another RocksDB instance) and network/CPU to follow the changelog. N=1 roughly doubles state footprint. - **Placement requirement**: you need enough distinct instances. With N=1 you need ≥2 instances for the standby to be placed on a *different* host; otherwise it can't help with machine failure. - **Not throughput**: standbys never run active tasks, so they add zero processing capacity — purely availability/recovery. - **Rack awareness** (KIP-708): `rack.aware.assignment.tags` can force standbys onto different racks/AZs so an AZ outage still leaves a usable copy. ## Edge cases - A standby can lag if it can't keep up with changelog throughput; on promotion it still replays the lag tail, so recovery isn't strictly instant. - Stateless tasks have no state store and thus no standby cost — standbys only matter for stateful sub-topologies. - Combined with `acceptable.recovery.lag`, the assignor decides whether a standby is 'caught up enough' to be promoted directly.

  • Do standby replicas increase processing throughput?
    No. Standbys never run active tasks or produce output; they only maintain a hot state copy for fast failover. Throughput is bounded by the active task count (input partitions).
  • How is a standby replica different from a warmup replica?
    A standby is a permanent failover copy controlled by num.standby.replicas. A warmup replica (KIP-441) is a temporary copy created during scaling/rebalancing to pre-restore state before a task moves, avoiding a processing stall; it's capped by max.warmup.replicas.
  • How would you ensure a standby survives an availability-zone outage?
    Use rack-aware assignment (KIP-708) via rack.aware.assignment.tags so Streams places the standby on a different rack/AZ than the active task.

saying these in an interview costs you the question

  • Saying standby replicas add throughput / process records
  • Setting num.standby.replicas>0 with a single instance and expecting failover protection
  • Confusing standby replicas with Kafka broker topic replication (replication.factor)

context