skip to content

How do workers in a distributed Connect cluster coordinate work assignment, and how did incremental cooperative rebalancing (KIP-415) improve the original protocol?

level: principalimportance: should knowfreq 35%

answer

  1. group coordinator + leader assigns
  2. eager = stop-the-world
  3. KIP-415 incremental cooperative (2.3)
  4. only moved tasks revoked
  5. scheduled.rebalance.max.delay.ms

basics

~20 s

Workers sharing a group.id use Kafka's group-membership protocol to elect a leader that assigns connectors and tasks. The original protocol stopped all work on every change (stop-the-world); incremental cooperative rebalancing (KIP-415) only reassigns the affected tasks, avoiding global pauses.

solid answer

~40 s

Distributed workers join a group via their shared `group.id` and run Kafka's rebalance protocol: members heartbeat to the group coordinator (a broker), one worker is elected **leader**, and the leader computes the assignment of connectors and tasks across all workers, distributed through the config topic. Originally Connect used **eager rebalancing**: any membership or config change revoked *all* assignments cluster-wide and reassigned from scratch — a 'stop-the-world' pause where every connector/task briefly halted, painful at scale. **KIP-415 (Incremental Cooperative Rebalancing**, Kafka 2.3) changed this: rebalances happen in cooperative rounds where only the tasks that actually need to move are revoked, unaffected tasks keep running, and a configurable `scheduled.rebalance.max.delay.ms` lets the cluster wait briefly for a bounced worker to return before redistributing its load — avoiding needless churn during rolling restarts.

go deeper

for a junior

Know that workers coordinate and a leader assigns tasks; rebalances happen on changes.

for a middle

Explain leader election via group.id and what triggers a rebalance.

for a senior

Contrast eager stop-the-world vs. incremental cooperative rebalancing and its benefit.

for a principal

Tune rebalance behavior for rolling upgrades/failover and reason about cluster-wide stability trade-offs.

## Coordination and rebalancing in distributed Connect ### How assignment works Workers sharing a `group.id` form a group managed by a Kafka broker acting as the **group coordinator** — the same membership machinery that backs consumer groups. Members send **heartbeats**; if one misses its `session.timeout.ms`, it's considered dead. One worker is the **leader**; the leader runs the assignment algorithm that distributes **connectors** (the management/coordination objects) and **tasks** (the data-moving units) across the live workers. The resulting assignment and all configs propagate via the internal `config.storage.topic`, so every worker can compute and apply its share. A **rebalance** is triggered when: a worker joins or leaves, a connector is created/deleted/reconfigured, or `tasks.max` changes. ### The original problem: eager (stop-the-world) rebalancing In the first protocol, *any* trigger caused **all** workers to revoke **all** their assignments and rejoin, after which the leader recomputed everything from scratch. Consequences: - Every connector and task in the cluster **paused** during the rebalance, even ones unaffected by the change. - During a **rolling restart** of N workers you'd incur N full stop-the-world rebalances. - At scale (hundreds of tasks), these pauses caused noticeable data-flow gaps and latency spikes. ### KIP-415: Incremental Cooperative Rebalancing (Kafka 2.3+) The redesign borrows the cooperative model: - **Incremental**: only the assignments that must change are revoked. Tasks already on the correct worker **keep running** through the rebalance — no global pause. - **Cooperative rounds**: rebalancing can take multiple short rounds; the leader signals which tasks to revoke, those are released, then a follow-up round assigns them — converging without halting everything. - **Delayed reassignment**: `scheduled.rebalance.max.delay.ms` (default 5 minutes) lets the cluster **wait** for a temporarily-absent worker (e.g., one being restarted) to rejoin before redistributing its tasks. This prevents a flurry of pointless reassignments during rolling restarts and brief network blips. ### Why it matters architecturally - **Rolling upgrades** of a Connect cluster become far cheaper — most tasks stay running while one worker bounces. - **Operational stability**: transient worker loss doesn't immediately reshuffle the whole cluster. - **Tuning**: operators balance `scheduled.rebalance.max.delay.ms` (longer = more tolerant of restarts but slower failover of genuinely dead workers) against failover speed. ### Edge cases - A genuinely crashed worker still has its tasks held idle until `scheduled.rebalance.max.delay.ms` elapses, then they're reassigned — so very long delays slow recovery from real failures. - Mixed-version clusters: cooperative protocol negotiates; ensure all workers support it to avoid falling back to eager behavior. - Connector imbalance can still occur; later improvements refined the assignor, but the cooperative core is KIP-415.

  • What does scheduled.rebalance.max.delay.ms control, and what's the trade-off in tuning it?
    How long the cluster waits for an absent worker to rejoin before reassigning its tasks. Longer values tolerate rolling restarts without churn but slow recovery when a worker has actually died; shorter values fail over faster but reshuffle on transient blips.
  • During a rolling restart of a large cluster, why is cooperative rebalancing dramatically better than eager?
    Eager pauses every task on each worker bounce (N stop-the-world rounds). Cooperative keeps unaffected tasks running and the delay timer often lets a bounced worker rejoin before any reassignment, so most data flow never stops.

saying these in an interview costs you the question

  • Saying ZooKeeper coordinates worker assignment (it's Kafka's group protocol via a broker coordinator).
  • Claiming eager rebalancing only paused the affected connector (it paused the whole cluster).
  • Not knowing scheduled.rebalance.max.delay.ms or attributing cooperative rebalancing to the wrong mechanism.

context