Explain incremental cooperative rebalancing (KIP-415) in Kafka Connect: how does it differ from eager, and what is connect.protocol?
answer
- KIP-415, Kafka 2.3
- revoke only the delta
- two rounds: revoke then assign
- connect.protocol: eager/compatible/sessioned
- compatible = default, falls back during upgrade
- sessioned = KIP-507 signed assignments
basics
~20 sKIP-415 stops revoking everything. Only tasks that must move are revoked; the rest keep running, so there's no cluster-wide pause. It uses two rebalance rounds (revoke, then assign). The connect.protocol setting selects eager, compatible, or sessioned.
solid answer
~40 sIncremental cooperative rebalancing (KIP-415, Kafka 2.3) replaces stop-the-world. The leader computes the desired assignment and revokes **only** the connectors/tasks that need to relocate, leaving everything else running. Revocation and assignment happen in **two rebalance rounds**: round one revokes the to-move tasks, round two assigns them to their new homes. Unrelated tasks never stop, so a worker joining or a single config change no longer pauses healthy pipelines. The behavior is selected by the worker-level `connect.protocol` config: `eager` (legacy stop-the-world), `compatible` (default since 2.3 — supports cooperative but can fall back to eager during a mixed-version upgrade), and `sessioned` (KIP-507, adds signed assignments so workers can't be impersonated). Cooperative also pairs with `scheduled.rebalance.max.delay.ms` to defer reassigning a briefly-absent worker's tasks, avoiding churn from transient blips.
go deeper
Know that KIP-415 means only tasks that must move are stopped, not the whole cluster.
Explain the revoke-only-the-delta idea and name the connect.protocol values (eager/compatible/sessioned).
Detail the two-round revoke-then-assign mechanics, the compatible-mode upgrade fallback, and how scheduled.rebalance.max.delay.ms pairs with it.
Reason about upgrade choreography across mixed-version clusters, the security model of sessioned/KIP-507, and tuning the delay knob against availability SLAs.
## The problem it solves The legacy **eager** protocol was *stop-the-world*: on any membership or config change, every worker revoked **all** its connectors and tasks before a new assignment was computed, pausing every pipeline in the cluster. **KIP-415 — Incremental Cooperative Rebalancing** (shipped in Apache Kafka 2.3) fixes this. ## Core idea: revoke only what must move Instead of asking everyone to drop everything, the leader: 1. Looks at the **current** assignment (who runs what now). 2. Computes the **target** assignment (who should run what after the change). 3. Revokes **only the delta** — the specific connectors/tasks that need to relocate to reach balance. 4. Leaves every other task running undisturbed. So adding a worker means the leader picks a handful of tasks to move onto it; all other tasks keep flowing. There is no cluster-wide pause. ## Two-round (incremental) mechanics A single round of the group protocol can't both revoke from old owners and assign to new owners safely (a task must be fully stopped on worker A before it starts on worker B, or you'd run it twice). So cooperative rebalancing is **incremental**, using two rebalance rounds: - **Round 1 (revoke):** the leader tells the relevant workers to revoke the to-be-moved tasks. Those workers stop just those tasks and rejoin. - **Round 2 (assign):** now that the tasks are free, the leader assigns them to their new owners, which start them. The word *cooperative* captures this: members cooperate by giving up only their share of the delta rather than everything. ## The connect.protocol setting A **worker-level** config (`connect.protocol`) chooses the strategy. All workers in a cluster negotiate the highest protocol they all support: - **`eager`** — legacy stop-the-world. Use only to force old behavior. - **`compatible`** — the **default** since 2.3. Workers use cooperative rebalancing, but during a **mixed-version upgrade** (some workers still on an old image) the cluster automatically **falls back to eager** until every worker supports cooperative. This makes upgrades safe with no config gymnastics. - **`sessioned`** — adds **KIP-507** session protection: the leader signs assignments with a key so a rogue/stale worker cannot forge or replay an assignment. It builds on cooperative behavior and adds security/correctness for the internal Connect topics. ## Why two rounds isn't a regression You might think two rounds is slower than eager's one. But round 1 only stops the *moving* tasks, not the whole cluster, and round 2 only starts those same few. The total disruption is proportional to the size of the **delta**, not the size of the **cluster** — which is the whole point. ## Companion: scheduled.rebalance.max.delay.ms Cooperative rebalancing introduced a delay knob (KIP-415, `scheduled.rebalance.max.delay.ms`, default 300000 ms / 5 min). When a worker **leaves**, the leader can leave that worker's tasks **unassigned** for up to this delay, betting the worker will come back (e.g. a quick restart or GC pause) and reclaim its own tasks — avoiding a move-out-then-move-back churn. If the worker doesn't return within the delay, the tasks are reassigned to the survivors. ## Edge cases / gotchas - **Mixed versions**: while any worker doesn't support cooperative, `compatible` keeps the whole cluster on eager — you don't get the benefit until the *last* old worker is upgraded. - **Lost tasks during the delay window**: with a non-zero `scheduled.rebalance.max.delay.ms`, a genuinely dead worker's tasks stay **down** until the delay elapses, trading availability for stability. Tune it to your tolerance. - **Imbalance after delay**: after the delay reassigns tasks, a later rebalance may rebalance again when the worker returns. Cooperative rebalancing's incremental moves keep that cheap.
- Why does cooperative rebalancing need two rebalance rounds instead of one?A task must be fully stopped on its old worker before it starts on the new one, or it would run in two places at once. Round 1 revokes the moving tasks; round 2 assigns them once they're free.
- What does connect.protocol=compatible do during a mixed-version rolling upgrade?Workers use cooperative rebalancing only if all members support it; while any worker is still on an old (eager-only) version, the cluster falls back to eager automatically, so the upgrade is safe without manual config changes.
- What does the sessioned protocol (KIP-507) add over compatible?It adds signed (session-keyed) assignments so workers can verify the leader's assignment is authentic, preventing a stale or rogue worker from forging assignments that touch the internal config topic.
saying these in an interview costs you the question
- Saying cooperative rebalancing still revokes all tasks — it revokes only the delta.
- Claiming connect.protocol is a connector-level setting; it's worker-level.
- Forgetting that 'compatible' (the default) falls back to eager during mixed-version upgrades.
- Confusing sessioned (KIP-507, security) with the basic cooperative behavior (KIP-415).
- Thinking the two-round approach makes it slower than eager overall.