skip to content

What does scheduled.rebalance.max.delay.ms control, and how do you tune it for worker restarts vs. genuine failures?

level: seniorimportance: should knowfreq 45%

answer

  1. default 300000 ms (5 min)
  2. wait before reassigning orphaned tasks
  3. bet the worker comes back
  4. high = less churn, longer stall on real death
  5. 0 = reassign immediately
  6. only departed worker's tasks idle

basics

~20 s

It's how long (default 5 minutes) the leader waits before reassigning a departed worker's tasks, hoping the worker returns and reclaims them. It avoids churn from quick restarts but keeps those tasks down during the wait if the worker is really dead.

solid answer

~50 s

`scheduled.rebalance.max.delay.ms` (worker-level, default 300000 ms = 5 min) is a cooperative-rebalancing knob. When a worker leaves the group, the leader doesn't immediately move its tasks to survivors. Instead it leaves them **unassigned** for up to this delay, betting the worker is just restarting or paused (GC, brief network blip) and will rejoin to reclaim its own tasks — avoiding an expensive move-out-then-move-back. If the worker returns within the window, no task migration happens. If it doesn't, after the delay the tasks are reassigned to the remaining workers. The trade-off: a higher value smooths rolling restarts and transient flaps but leaves tasks **idle (data not flowing)** longer when a worker truly dies; a lower value (even 0) restores work faster but reintroduces churn. Set it near your typical worker restart time, balanced against how long you can tolerate those partitions stalled.

go deeper

for a junior

Know it's a ~5-minute wait before a gone worker's tasks move to others, so quick restarts don't cause churn.

for a middle

Explain the transient-vs-permanent trade-off and that only the departed worker's tasks are idle during the wait.

for a senior

Tune it against typical worker restart time and the tolerable data-stall window; distinguish it from session.timeout.ms.

for a principal

Frame it as an availability-vs-stability dial across deployment patterns (K8s rollouts, latency-sensitive sinks) and reason about its interaction with the full membership-timeout chain.

## What it is `scheduled.rebalance.max.delay.ms` is a **worker-level** Kafka Connect configuration introduced with **incremental cooperative rebalancing (KIP-415)**. Default: **300000 ms (5 minutes)**. It answers one question: *when a worker disappears, how long should the cluster wait before giving its work to someone else?* ## Why a delay exists at all Workers leave the group for two very different reasons: - **Transient**: a rolling restart, a JVM GC pause, a momentary network hiccup, a redeploy in Kubernetes. The worker will be back in seconds. - **Permanent**: the host died, the process crashed and won't restart, the node was terminated. Without a delay, both look identical, and the cluster would immediately move the departed worker's tasks to survivors. For the **transient** case that's wasteful: the worker comes back moments later and the tasks get moved **again** — two rebalances and two task migrations for a blip. With many workers churning (e.g. a rolling restart of a big cluster), that thrash is severe. ## How the delay works When a worker leaves: 1. The leader notices the worker's connectors/tasks are now **orphaned**. 2. Rather than reassigning them, it marks them to be assigned **after** a scheduled delay and leaves them **unassigned** in the meantime. 3. A timer counts down up to `scheduled.rebalance.max.delay.ms`. 4. **If the worker rejoins before the timer expires**, it reclaims its original tasks — no migration, minimal disruption. 5. **If the timer expires first**, the leader reassigns the orphaned tasks to the surviving workers. Note it's a **max** delay: if the worker returns early, the wait ends early. ## The core trade-off - **Higher value** (e.g. 5+ min): great for clusters with frequent rolling restarts/redeploys — avoids churn, keeps assignments stable. **Cost:** when a worker *genuinely* dies, its tasks stay **down** (no data flowing for its source/sink partitions) for the full delay. - **Lower value** (e.g. 0–60 s): faster recovery from real failures — tasks restart on survivors quickly. **Cost:** every quick restart or flap triggers move-out-then-move-back churn. `0` disables the delay: orphaned tasks are reassigned immediately (closer to pre-delay behavior, maximizing availability but maximizing churn). ## How to tune it Match it to your **typical worker downtime**: - **Kubernetes with rolling deploys**: set it comfortably above your pod restart/reschedule time (often a few minutes) so a normal redeploy doesn't migrate tasks. But cap it at the longest data stall you can tolerate. - **Latency-sensitive sinks** (feeding live dashboards/alerts): lower it so a dead worker's partitions resume quickly, accepting more churn. - **Large, stable clusters doing frequent rollouts**: keep the default or higher to minimize rebalance thrash. ## Interaction with other settings - It only matters under **cooperative** rebalancing (`connect.protocol=compatible` or `sessioned`); under `eager` there's no scheduled-delay concept. - It works **with** the group membership `session.timeout.ms` machinery: the worker is first declared *gone* by the group coordinator (session timeout / heartbeat loss), and *then* the scheduled delay decides how long its orphaned tasks wait. They're distinct timers serving different purposes. ## Common mistakes - Believing the delay keeps the whole cluster paused — it doesn't; only the **departed worker's** tasks are idle, everything else keeps running. - Setting it very high to 'reduce rebalances' without realizing a dead worker's data flow stalls for that whole window. - Confusing it with `session.timeout.ms` (how long before a missing worker is even considered gone).

  • If a worker truly crashes and never returns, when do its tasks restart on other workers?
    After scheduled.rebalance.max.delay.ms elapses (default 5 min). During that window those specific tasks are idle and their data doesn't flow; the rest of the cluster is unaffected.
  • How is scheduled.rebalance.max.delay.ms different from session.timeout.ms?
    session.timeout.ms decides when a missing worker is declared gone (via lost heartbeats). scheduled.rebalance.max.delay.ms then decides how long that gone worker's orphaned tasks wait before being reassigned, in case it rejoins.

saying these in an interview costs you the question

  • Saying the delay pauses the whole cluster — only the departed worker's tasks are affected.
  • Setting it very high without acknowledging the real-failure data-stall cost.
  • Confusing it with session.timeout.ms (membership detection vs. reassignment delay).
  • Claiming it applies under eager rebalancing — it's a cooperative-protocol feature.
  • Thinking 0 means 'never reassign' rather than 'reassign immediately'.

context