How would you design a rolling restart / upgrade of a large Connect cluster to minimize rebalancing disruption?
answer
- cooperative on (compatible/sessioned)
- delay > single-worker restart time
- one worker at a time, settle between
- mixed-version = eager fallback window
- freeze connector configs
- wait for tasks RUNNING, not just process up
basics
~20 sUse cooperative rebalancing (connect.protocol=compatible/sessioned) so only moving tasks pause. Set scheduled.rebalance.max.delay.ms above your per-worker restart time so restarting workers reclaim their own tasks. Restart one worker at a time and wait for the cluster to settle between each.
solid answer
~50 sThree levers. First, ensure **cooperative rebalancing** is active (`connect.protocol=compatible` or `sessioned`) so each restart only moves tasks that must move instead of stopping the cluster; note that during a mixed-version upgrade `compatible` may sit in eager fallback until the last old worker is upgraded, so expect heavier rebalances until then. Second, set `scheduled.rebalance.max.delay.ms` comfortably above a single worker's restart time (the default 5 min usually suffices) so a restarting worker rejoins and **reclaims its own tasks** without a move-out/move-back. Third, restart **one worker at a time**, waiting for the cluster to fully settle (assignments stable, tasks RUNNING) between each — restarting several at once defeats the delay logic and can orphan many tasks. Drain in-flight work where the connector supports it, monitor rebalance counts and task status via the REST API, and avoid editing connector configs mid-rollout (each edit triggers its own rebalance).
go deeper
Know to restart workers one at a time and rely on cooperative rebalancing, not all at once.
Explain using the scheduled delay so a restarting worker reclaims its tasks, and settling between restarts.
Account for the mixed-version eager-fallback window, config-edit freezes, and monitoring task status during the rollout.
Design the full choreography: protocol settings, delay sizing, sequencing automation with readiness gates, capacity headroom, and exactly-once source considerations.
## The goal Upgrading or restarting a large Connect cluster (new image, JVM, config) means each worker leaves and rejoins the group — every one of those is a potential rebalance. The aim is to keep data flowing and avoid **rebalance storms** while every worker cycles. ## Lever 1 — cooperative rebalancing must be on Confirm `connect.protocol` is `compatible` (default ≥ 2.3) or `sessioned`, **not** `eager`. Under cooperative rebalancing, restarting a worker revokes only the **delta** of tasks, leaving healthy pipelines running. **Mixed-version caveat:** if the upgrade changes the Connect version, `compatible` mode **falls back to eager** for as long as any worker still runs the old version. So during the rollout you may get stop-the-world behavior until the **last** old worker is upgraded — plan for heavier rebalances in that window, and do the rollout promptly rather than leaving the cluster half-upgraded for long. ## Lever 2 — the scheduled delay buys free restarts Set `scheduled.rebalance.max.delay.ms` **above** how long one worker takes to restart and rejoin (default 300000 ms / 5 min is usually plenty for a process bounce; size up for slow container scheduling). Then, when you stop a worker: 1. Its tasks become orphaned but are **held unassigned** for the delay. 2. The worker restarts and rejoins **within** the delay. 3. It **reclaims its own tasks** — no migration, minimal disruption. If the restart overruns the delay, the tasks get reassigned to survivors and then likely move back when the worker returns — so the delay must exceed your real restart time. ## Lever 3 — one worker at a time, settle between Restart workers **sequentially**: - Stop worker, wait for it to come back and for the cluster to **settle**: assignments stable, all tasks `RUNNING` (check via `GET /connectors/{name}/status` and the worker logs / metrics for rebalance activity). - Only then move to the next worker. Restarting **several at once** orphans many tasks simultaneously, can exceed the delay's protection, and forces large reassignments to the shrinking set of survivors — exactly the storm you're avoiding. ## Supporting practices - **Freeze connector configs during the rollout.** Editing a connector (including `tasks.max`) triggers its own rebalance; stacking that on top of restart rebalances multiplies churn. - **Watch the right signals:** rebalance rate/latency metrics, task status via REST, consumer lag on sink connectors (to spot stalled data), and the leader's logs for assignment changes. - **Mind exactly-once source connectors** (KIP-618): they have stricter restart/fencing semantics; ensure transactional state settles before proceeding. - **Capacity headroom:** during each single-worker restart the rest of the fleet briefly carries (or holds) that worker's share; ensure survivors can absorb it if the delay expires. - **Order/automation:** orchestrators (Kubernetes rolling updates, Ansible) should be configured with `maxUnavailable: 1`-style one-at-a-time semantics and a readiness gate that waits for tasks RUNNING, not just the process up. ## Why this works Cooperative rebalancing minimizes *what* moves; the scheduled delay minimizes *whether* anything moves for a quick restart; one-at-a-time sequencing minimizes *how many* tasks are in flux at once. Together they turn a fleet upgrade from a series of cluster-wide pauses into a near-transparent rollout — except for the unavoidable eager-fallback window during a version change. ## Common mistakes - Rolling all workers at once and overwhelming the delay/cooperative protections. - Setting the delay **below** restart time, causing move-out-then-move-back churn. - Forgetting the mixed-version eager fallback and being surprised by stop-the-world pauses mid-upgrade. - Editing connectors during the rollout. - Treating 'process up' as 'settled' instead of waiting for tasks RUNNING.
- During a version upgrade, why might you still see stop-the-world rebalances even with connect.protocol=compatible?Compatible mode falls back to eager while any worker is still on the old (cooperative-unaware) version. Until the last old worker is upgraded, the cluster behaves eagerly, so plan for heavier rebalances in that window.
- What's the danger of restarting multiple workers simultaneously during the rollout?Many tasks are orphaned at once, which can exceed the scheduled-delay protection and force large reassignments onto a shrunken set of survivors — a rebalance storm and a capacity spike. Restart one at a time and settle between.
- Why should the scheduled.rebalance.max.delay.ms exceed a single worker's restart time?So the restarting worker rejoins within the delay and reclaims its own tasks. If the restart takes longer than the delay, its tasks get reassigned to survivors and then likely move back when it returns — extra churn.
saying these in an interview costs you the question
- Restarting the whole cluster at once and expecting no disruption.
- Ignoring the mixed-version eager fallback during a version upgrade.
- Setting the scheduled delay below the real worker restart time.
- Editing connector configs in the middle of the rollout.
- Treating 'process started' as 'cluster settled' without checking tasks are RUNNING.