How does cooperative rebalancing work in Kafka Streams, and why was it introduced over the older eager protocol?
answer
- KIP-429, default since 2.4
- eager = stop-the-world, revoke all
- cooperative = revoke only what moves
- two-phase, sticky, keeps RocksDB local
- warmup replicas (KIP-441) build on it
basics
~20 sCooperative rebalancing (the default since 2.4) lets instances keep tasks they already own during a rebalance instead of giving everything up. Only the tasks that must move are revoked, avoiding a global stop-the-world pause and unnecessary state restoration.
solid answer
~40 sBefore 2.4, Kafka Streams used the eager protocol: on any rebalance every member revoked all its partitions/tasks, then the assignor redistributed from scratch — a 'stop-the-world' pause where stateful tasks might land elsewhere and have to restore state from changelog topics. Cooperative rebalancing (KIP-429), the default via StreamsPartitionAssignor / CooperativeStickyAssignor, splits the rebalance into two phases. In phase one, the assignor computes the target assignment; members only revoke tasks that are actually being reassigned to someone else and keep processing the rest. In phase two, a follow-up rebalance hands the freed tasks to their new owners. The benefits: most tasks never pause, stateful tasks stay put (stickiness preserves local RocksDB state, avoiding costly changelog restore), and adding/removing one instance causes minimal disruption. This is what makes rolling restarts and incremental scaling cheap.
go deeper
Know the term: cooperative rebalancing avoids stopping the whole app on every rebalance.
Contrast eager (revoke all) vs cooperative (revoke only what moves) and why it's faster.
Explain the two-phase sticky protocol, state locality, and how warmup replicas extend it for scaling.
Own the upgrade path (upgrade.from), tune probing/recovery-lag configs, and reason about availability during fleet changes.
## Background: what a rebalance is A **rebalance** is the Kafka consumer-group process of (re)assigning partitions/tasks to members when membership changes (an instance joins, leaves, or crashes) or topic metadata changes. Kafka Streams piggybacks task assignment on the consumer group protocol via the **StreamsPartitionAssignor**. ## The old way: eager rebalancing Under the **eager** protocol, the rebalance is **stop-the-world**: 1. Every member **revokes ALL** of its partitions (stops processing entirely). 2. The leader computes a brand-new assignment. 3. Every member gets its (possibly different) partitions and resumes. Problems: - **Global pause**: all processing halts during the rebalance, even for tasks that won't move. - **State restoration cost**: a stateful task is backed by a local **state store** (RocksDB) whose updates are also written to a **changelog topic**. If the task moves to a new instance, that instance must **restore** the store by replaying the changelog from the start (or from a checkpoint) before it can process — potentially minutes of downtime for large stores. ## The new way: cooperative (incremental) rebalancing — KIP-429 Default since Kafka 2.4. The key idea: **only revoke what actually moves**, in two passes. 1. **First rebalance**: the leader computes the desired final assignment. Members compare it to what they currently own and **revoke only the tasks being reassigned away**. Crucially, they **keep processing all tasks they get to retain** — no global pause. 2. **Second (follow-up) rebalance**: the now-free partitions are handed to their new owners, who begin restoring/processing them. The assignor is **sticky**: it tries to keep each stateful task on the instance that already has its local state, so it avoids changelog restoration whenever possible. ## Why it matters for scaling - **Rolling restarts / deploys**: bouncing one instance only moves that instance's tasks; the rest of the fleet keeps running. - **Incremental scale-out**: adding one instance moves just enough tasks to balance, instead of churning the whole assignment. - **Warmup replicas (KIP-441)**: builds on cooperative rebalancing — when a task must move, Streams first spins up a **warmup replica** that restores state in the background, and only transfers the active task once it's caught up (governed by `acceptable.recovery.lag` and probing rebalances every `probing.rebalance.interval.ms`, default 10 min). This keeps availability high during scaling. ## Edge cases & gotchas - **Upgrades**: moving from eager to cooperative requires a careful two-rolling-bounce upgrade path (set `upgrade.from` on the first bounce) because the two protocols can't be mixed naively. - A misbehaving member that fails to rejoin within `session.timeout.ms` / `max.poll.interval.ms` still triggers reassignment of its tasks. - Cooperative rebalancing reduces but does not eliminate restoration: a task genuinely moving to a fresh instance still restores (warmup replicas hide most of that latency).
- Why does cooperative rebalancing reduce changelog/state restoration?Because it's sticky: tasks that don't need to move stay on the instance that already holds their local RocksDB state, so no changelog replay is needed. Only genuinely relocated tasks restore — and warmup replicas (KIP-441) restore those in the background first.
- What is special about upgrading an app from eager to cooperative rebalancing?It needs a two-step rolling bounce: on the first bounce you set upgrade.from to the prior version so members still speak the old protocol, then a second bounce switches everyone to cooperative. You can't safely mix the two protocols in one group.
saying these in an interview costs you the question
- Saying cooperative rebalancing means no rebalance ever happens
- Claiming every rebalance is stop-the-world in modern Streams (eager is no longer the default)
- Confusing cooperative rebalancing with standby replicas (different mechanisms)