skip to content

questions

4

Each cluster node holds copies of partitions (or queues) - why does a maintenance restart take one node at a time, waiting between?

level: juniorimportance: must knowfreq 70%

answer

  1. redundancy, not throughput
  2. a stopped copy stops following
  3. catch-up distance grows with downtime
  4. full redundancy before the next stop

basics

~20 s

A restart takes that node's copies out of service, and while it is down they fall behind. Stopping the next node before they have caught up removes a second current copy and can leave a partition with none.

solid answer

~40 s

A cluster keeps every unit of ownership - a partition, or on queue-shaped brokers a queue - on several nodes. One node serves it and the others keep their copies current by continuously fetching what the serving node accepted. Stopping a node moves service elsewhere quickly, but it also stops every copy that node holds; while it is down those copies fall behind by whatever traffic arrives, and on return they have to fetch that `catch-up distance` before they count as current again. The wait between steps exists so the cluster is back to the redundancy it had before that step is finished with. Without it, each step leaves some units with one fewer current copy, and the erosion compounds across the loop. Note this is a durability argument, not a throughput one.

go deeper

for a junior

Recall that a node holds copies of data, that a stopped copy falls behind while writes continue, and that it must catch up before the next node is stopped.

for a middle

Explain the mechanics: what leaves the caught-up set at the stop, what the returning node has to fetch, and why catch-up time tracks write rate rather than downtime.

for a senior

Show you have watched this go wrong - redundancy eroding one step at a time across a long loop, and a unit that was already short before maintenance began.

for a principal

Frame the trade: how much redundancy the estate is willing to run on during routine maintenance, and whether detached storage would remove the cost entirely.

## What one restart actually removes A broker or streaming cluster keeps each **unit of ownership** - a partition, or on queue-shaped brokers a queue - on more than one node, so that losing a node does not lose records. At any moment one node serves that unit, and the other nodes holding it keep their **copies** current by continuously fetching whatever the serving node has accepted. The subset of copies current enough to be promoted, or to satisfy a write, is the **caught-up set**. Stopping one node for maintenance does two different things, and they are easy to confuse: - **Service moves.** Every unit that node was serving is taken over by another node holding a copy. On most designs this takes seconds and writers see little more than a brief retry. - **Copies stop following.** Every copy on that node, for every unit it holds, stops fetching. This is the part that does not undo itself the moment the process comes back. While the process is down the cluster keeps accepting writes. Each copy on the stopped node therefore falls behind by exactly the traffic accepted meanwhile - its **catch-up distance**. When the process returns it must fetch that distance before those copies are back in the caught-up set. Until then, those units are running with fewer current copies than the cluster was designed to keep. ## Why the wait, and not the order, is the safety One safe step looks like this: 1. Stop node A. 2. Every copy A holds leaves the caught-up set; every unit A was serving is served elsewhere. 3. A returns and begins fetching what it missed. 4. Some time later - seconds, or much longer - every copy A holds is current again. 5. Only now is it safe to stop node B. Skipping step 4 does not fail immediately, which is exactly why the mistake survives in scripts for years. It **erodes** redundancy rather than breaking it: after one skipped wait some units have one fewer current copy, after the next some have two fewer, and a long enough loop eventually reaches a unit whose only current copy sits on the node about to be stopped. Restart order matters too - a node can be deliberately left for last - but the wait is what makes any order safe. ## What the wait depends on The honest answer to how long to wait is that it depends on the design and on the traffic, not on a number you can carry between clusters: | what the design does | what the returning node must do | effect on the wait | |---|---|---| | each node keeps local storage, one leader with following copies | fetch everything written while it was down | wait grows with write rate and with downtime | | replication by majority write | reach the point a majority already holds | similar, and a majority must stay available throughout | | records on shared or remote storage (detached storage) | little or nothing to re-fetch | wait is short, sometimes effectively absent | | designs that resynchronise a whole unit rather than resuming | refetch the unit from the beginning | wait far longer than the downtime that caused it | That last row is the one people are most often surprised by: on some designs a copy that has been away resumes from where it stopped, and on others it starts again, so a two-minute restart can owe an hour of fetching. ## It is not the capacity argument Losing one node out of many barely dents throughput; a cluster is sized so that it can. The reason for the wait is **durability** - how many current copies of each record exist while the loop is running - plus the brief unavailability of any unit that is momentarily unserved. A candidate who answers only in terms of serving capacity has answered a different question, the one about replacing stateless workers behind a shared name. ## What this looks like when it goes wrong - A script restarts a whole batch of nodes inside the scheduled maintenance hour; nothing errors, and some units spend the entire loop on a single current copy. - A restart that took one minute is followed by a wait of one minute, on a cluster whose busiest unit needs ten. - A node whose copies were already short before maintenance began is restarted anyway, because nothing in the loop ever looked. - The loop finishes, the dashboard is green by the time anyone looks, and the only trace is a gap in one stream. The junior-level takeaway is small and load-bearing: a returned process is not a restored copy, and the gap between those two facts is the whole of rolling restart discipline.

  • Does the wait get shorter if the restart itself is quicker?
    Only partly. A shorter outage means less traffic missed, so less to fetch. But the fetching competes with live reads and writes on the same disks and links, so a two-minute absence on a busy cluster can still owe far more than two minutes of catch-up - and on designs that resynchronise a whole unit rather than resuming, the restart length barely matters at all.
  • Is the wait still needed if the node being restarted was not serving anything?
    Yes. Serving is only one of the two roles a node plays for a unit. A node that merely holds copies still stops following when it goes down, and those copies still leave the caught-up set. A loop that only considers which units a node was serving will happily strip the current copies out from under units it never served.
  • Why does a design with detached storage change this picture?
    Where a unit's records live on shared or remote storage rather than a node's local disk, a returning node has little or nothing to re-fetch, because the records were never only on it. Ownership can move back almost immediately. The step is then cheap, which is one of the main operational arguments for that design - though metadata and caches still need a moment.

saying these in an interview costs you the question

  • Says a restarted node is fine as soon as the process is up
  • Treats the wait as being about serving capacity only
  • Assumes copies resume instantly because the stored files were never deleted
  • Thinks the restart order alone makes a maintenance loop safe
  • Believes a broker never allows an operator command that could lose records
open as a page

In a rolling restart of a broker cluster, why is the process answering again a weaker gate than a cluster-side catch-up check?

level: middleimportance: must knowfreq 63%

basics

~20 s

A process answering is a node-side fact about itself; what the next step needs is a cluster-side fact about data - that every copy the restarted node holds is back in the caught-up set. The two become true minutes apart.

open as a page

A nightly script restarts each broker node with a fixed 60-second sleep between steps; afterwards acknowledged records are missing - why?

level: seniorimportance: should knowfreq 54%

basics

~20 s

The sleep was shorter than the catch-up each returning node owed, so current copies were never restored between steps. The count fell step by step until the loop stopped the node holding the last current copy of a unit.

open as a page

As the owner of a change-safety standard, what evidence would you require before any broker cluster node is stopped for maintenance?

level: principalimportance: should knowfreq 42%

basics

~20 s

Machine-checkable evidence, gathered before the stop and enforced by the tool: no unit of ownership is short of its current copies, the previous step has been paid back, and the loop aborts rather than proceeds when that cannot be shown.

open as a page