skip to content

Explain incremental cooperative rebalancing (KIP-415) in Kafka Connect: how does it differ from eager, and what is connect.protocol?

level: seniorimportance: must knowfreq 65%

answer

  1. KIP-415, Kafka 2.3
  2. revoke only the delta
  3. two rounds: revoke then assign
  4. connect.protocol: eager/compatible/sessioned
  5. compatible = default, falls back during upgrade
  6. sessioned = KIP-507 signed assignments

basics

~20 s

KIP-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 s

Incremental 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

for a junior

Know that KIP-415 means only tasks that must move are stopped, not the whole cluster.

for a middle

Explain the revoke-only-the-delta idea and name the connect.protocol values (eager/compatible/sessioned).

for a senior

Detail the two-round revoke-then-assign mechanics, the compatible-mode upgrade fallback, and how scheduled.rebalance.max.delay.ms pairs with it.

for a principal

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.

context