What problem did the original eager rebalancing protocol in Kafka Connect cause, and why was it painful at scale?
answer
- eager = stop-the-world
- revoke ALL then reassign
- every trigger pauses whole cluster
- rolling restart = N pauses
- KIP-415 fixed it
basics
~20 sThe old 'eager' protocol was stop-the-world: on any membership change every worker revoked ALL its connectors and tasks, then everything was reassigned. So one worker joining or a brief blip paused the entire cluster's data flow.
solid answer
~40 sThe original Connect rebalancing was the **eager** protocol (inherited from the early consumer-group design). On any rebalance trigger — a worker joining, leaving, crashing, or even a connector config change — the leader revoked the **entire** assignment from **every** worker. All connectors and tasks across the whole cluster stopped, then a single new global assignment was computed and everyone restarted their work. This is 'stop-the-world': total downtime proportional to nothing useful, paid on every membership change. It hurt badly at scale because large clusters rebalance often (rolling restarts, autoscaling, transient network issues), and each event halted all pipelines — even tasks that didn't need to move. It also amplified instability: a flapping worker could repeatedly pause the cluster. KIP-415 introduced incremental cooperative rebalancing to fix exactly this.
go deeper
Know the phrase 'stop-the-world': the old protocol paused the whole cluster on any change.
Explain the revoke-all-then-reassign sequence and why rolling restarts and flapping workers were so costly.
Connect the failure mode to availability SLAs and rebalance storms, and frame KIP-415 as the direct remedy.
Discuss the systemic fragility eager rebalancing created under churn and how it shaped operational practices before cooperative rebalancing.
## Background: what a rebalance is A **rebalance** in Kafka Connect is the process where the cluster recomputes which worker runs which connectors and tasks. It is triggered by: - a worker **joining** (scale-out, restart) - a worker **leaving** or **crashing** (scale-in, failure) - a connector being **created, deleted, or reconfigured** - a connector changing how many tasks it wants ## The eager protocol (the old way) Connect originally used the **eager** assignment strategy, borrowed from how Kafka consumer groups first worked. Its defining rule: **before** any new assignment can be computed, every member must give up everything it currently holds. Concretely, on each rebalance: 1. Every worker **revokes (stops) all** of its connectors and tasks. 2. All workers rejoin the group and report they hold nothing. 3. The leader computes a fresh global assignment from scratch. 4. Every worker starts the tasks it was just assigned. Because step 1 stops *all* work *everywhere*, this is called **stop-the-world**. During the rebalance window no data flows through *any* pipeline in the cluster. ## Why this hurt at scale The cost is paid **on every trigger**, regardless of how small the actual change is. Examples of the pain: - **Adding one worker** to a 20-worker cluster pauses all 20 workers, even though you only wanted to offload a few tasks onto the newcomer. - **A rolling restart** (restart workers one at a time to upgrade) triggers a stop-the-world rebalance *per* worker restarted — N pauses for N workers. - **A flapping worker** (one repeatedly dropping and rejoining due to GC pauses or network blips) can pause the entire cluster over and over. - **A trivial config edit** to one connector stops every unrelated connector too. The downtime is also proportional to how long it takes the slowest worker to revoke and restart its tasks, which grows with cluster size and task count. ## Secondary problems - **Reduced availability** directly: SLA-sensitive sink connectors (e.g. feeding a dashboard) show gaps on every rebalance. - **Rebalance storms**: because each event is expensive and membership changes can cascade, eager rebalancing made clusters feel fragile under churn. - **Wasted movement**: tasks that were on a perfectly healthy worker and didn't need to move were still stopped and (often) reassigned to the same worker — pure overhead. ## The fix **KIP-415 — Incremental Cooperative Rebalancing** replaced eager assignment so that, instead of revoking everything, the cluster revokes **only the tasks that actually need to move** and leaves the rest running. That eliminates the stop-the-world pause for the common case. Understanding the eager protocol's failure mode is the motivation for everything cooperative rebalancing does.
- Why was a rolling restart especially painful under eager rebalancing?Each worker restart is a membership change, and each membership change triggers a full stop-the-world rebalance. Restarting N workers one at a time causes roughly N cluster-wide pauses.
- Did eager rebalancing move tasks that didn't need to move?Yes. Every worker revoked everything, so even tasks staying on a healthy worker were stopped and re-started, which is pure wasted overhead.
saying these in an interview costs you the question
- Saying only the joining/leaving worker stopped its tasks — eager stopped ALL workers.
- Confusing eager rebalancing with the incremental cooperative protocol that replaced it.
- Claiming rebalances were cheap or rare; at scale they were frequent and expensive.
- Thinking a config change to one connector only affected that connector under eager.