Walk through what happens to state stores during a rebalance, and explain how standby replicas and state.dir affect restoration time.
answer
- Task moves → RESTORING → replay changelog → RUNNING
- Restore time ∝ changelog size
- num.standby.replicas = warm copies, fast failover
- Durable state.dir + checkpoint → replay only tail
- Cooperative-sticky / probing rebalance (KIP-441)
basics
~20 sWhen tasks move between instances during a rebalance, the new owner must rebuild each store by replaying its changelog topic before it can process records. This restore can be slow. Standby replicas keep warm copies on other instances to avoid or shorten it, and a persistent state.dir lets an instance reuse local state instead of restoring from scratch.
solid answer
~50 sOn a rebalance, the group coordinator reassigns tasks. If a stateful task lands on an instance that has no local copy of its store, that instance enters **restoration**: a restore consumer replays the store's changelog partition into the local store up to the latest offset, and only then does the task transition to RUNNING. Restore time scales with changelog size, so a cold task can stall processing. Two mitigations: (1) **`num.standby.replicas`** maintains warm shadow copies of stores on other instances that continuously consume the changelog, so on failover the standby is already (nearly) up to date and restoration is near-instant; the cooperative-sticky assignor also prefers placing tasks where state already exists. (2) A durable **`state.dir`** that survives restarts lets the same instance reuse its on-disk RocksDB files plus a checkpoint file, replaying only the changelog tail. Since 2.6+, **probing rebalances** move tasks back to instances with warm state to minimize long restores. With EOS, on restore Streams discards uncommitted state via the checkpoint to avoid replaying aborted writes.
go deeper
Know that after a rebalance a new instance rebuilds the store from the changelog before processing.
Describe the RESTORING→RUNNING flow and that restore time grows with changelog size.
Explain standby replicas, durable state.dir + checkpoints, and cooperative-sticky/probing rebalances as restore mitigations.
Set recovery-time SLAs by tuning standby count, storage durability, and assignor behavior against resource cost and EOS semantics.
## Why rebalances threaten state Kafka Streams instances form a consumer group. When membership changes — an instance joins, leaves, or crashes, or partitions are added — the group **rebalances** and tasks are reassigned. A **stateful task** carries a state store. If the task moves to an instance that doesn't already have that store on local disk, the store must be **restored** before processing can resume, because correctness requires the full prior state (you can't resume a running count from empty). ## The restoration sequence 1. After assignment, a task with missing/stale local state goes into the **RESTORING** state. 2. A dedicated **restore consumer** reads the store's **changelog** partition from the last checkpointed offset (or from the beginning if there's no local state) up to the high-water mark. 3. Each changelog record is written back into the local store (RocksDB or in-memory), rebuilding it. 4. Only when restore completes does the task become **RUNNING** and start processing new input. `restore.consumer.*` and `max.poll.records` influence throughput here; the `StateRestoreListener` lets you observe progress. Restoration time is **proportional to changelog size** (number of records to replay). For a large store this can mean minutes of unavailability for the affected partitions — a real production concern. ## Mitigation 1: standby replicas Setting **`num.standby.replicas=N`** tells Streams to maintain *N* extra warm copies of each store on other instances. A standby task continuously tails the changelog into a local store but does no processing. When the active instance fails, the coordinator promotes a standby whose state is already (nearly) current, so restore is near-instant instead of a full replay. Standbys cost extra disk, memory, and changelog read bandwidth, and add a small replication lag window, but they are the primary lever for fast failover and high availability. ## Mitigation 2: durable state.dir and stickiness If **`state.dir`** points at storage that **survives instance restarts** (e.g., a mounted volume, not ephemeral `/tmp`), a restarting instance finds its RocksDB files and a **checkpoint file** recording how far the store was restored. It then replays only the **tail** of the changelog (records after the checkpoint), which is fast. The **cooperative-sticky** assignment strategy and **probing rebalances** (KIP-441, Streams 2.6+) actively try to keep/return each task on the instance that already holds its state, avoiding cold restores. Streams reports a per-task **lag** so the assignor can move work to the warmest instance. ## Exactly-once interaction Under **exactly-once (EOS)**, a crash may have left uncommitted writes in the local store. On recovery Streams uses the checkpoint to discard state past the last committed offset and restores forward from the changelog, so aborted/duplicate writes don't survive. This is why EOS may wipe and rebuild local state more aggressively than at-least-once. ## Putting it together (capacity/SLA view) - Big stores + no standbys + ephemeral state.dir = worst case: full changelog replay on every failover. - Standbys + durable state.dir + cooperative-sticky = near-zero restore in the common case. - The trade-off is resource cost (standby disk/mem/bandwidth) vs. recovery-time objective.
- Why is a persistent state.dir on a mounted volume better than the default /tmp for restoration?Surviving RocksDB files plus a checkpoint file let a restarted instance replay only the changelog tail since the checkpoint, not the whole topic. With ephemeral /tmp the local state is gone on reboot, forcing a full cold restore.
- What does num.standby.replicas trade off?Faster failover and higher availability versus extra disk, memory, and changelog-read bandwidth on standby instances, plus a small replication-lag window. It's the main lever for shrinking restore-driven downtime.
- What did KIP-441 (probing rebalances) improve?It avoids long stop-the-world restores by assigning a moved task temporarily to the warm instance while a new target catches up in the background, then handing over once its state lag is small — minimizing unavailability during scaling.
saying these in an interview costs you the question
- Saying restoration replays the original input topics rather than the changelog
- Claiming standby replicas serve reads or do processing — they only keep warm state
- Believing the default state.dir (/tmp) is fine for production restore behavior
- Thinking a rebalance never requires restoration as long as the changelog exists — cold instances still must replay it