skip to content

Why is removing a node from a broker cluster an evacuation with a completion condition, rather than simply stopping its process?

level: middleimportance: must knowfreq 64%

answer

  1. the cluster cannot tell you from a crash
  2. move ownership before removing the node
  3. redundancy must never dip first
  4. it ends when the node holds nothing
  5. duration is bytes over bandwidth

basics

~20 s

Stopping the process is indistinguishable from that node failing: its copies stop being current and the units it led must be promoted in a hurry. An evacuation moves ownership off first and finishes only when the node holds nothing.

solid answer

~50 s

A cluster cannot tell a deliberate shutdown from a crash. If you stop the process, every **unit of ownership** — every partition (or queue) — it was serving must be promoted elsewhere at that instant, every copy it held stops being current, and the cluster starts re-creating those copies anyway, unplanned, at whatever moment you chose. An evacuation inverts the order: ownership is moved off the node first, while it is still healthy and still serving, and the node is stopped only afterwards. The completion condition is concrete — it holds no unit of ownership, and no unit it used to hold is left short of a current copy. The duration is the records it held divided by the bandwidth available to move them, so it is measured in minutes or hours, not seconds. Where storage is detached the same evacuation still happens, but it is a handover rather than a transfer, so it is far shorter.

go deeper

for a junior

The point to recall is the order: move the work off the node first, then stop it. A node that is simply switched off looks to the cluster exactly like one that crashed.

for a middle

Be able to state the completion condition out loud — the node holds no unit of ownership and nothing it held is short of a current copy — and to estimate the duration as bytes over available bandwidth.

for a senior

Show the sequencing judgment: which node you evacuate first when several must go, why you do not batch them when copies are already thin, and how you keep the change reversible until the last step.

for a principal

Own the standard rather than the run: what evidence a team must produce before any node in the estate is stopped, and who is allowed to waive it when a machine has to go tonight.

## A stopped node is a failed node A broker or streaming cluster reacts to what it observes, and what it observes when you stop a process is exactly what it observes when the machine dies: a member that no longer answers. Every consequence of a failure follows. - Each **unit of ownership** — a partition (or queue) that exactly one node serves at a time — that this node led must be handed to another node immediately, while writers retry into the gap in which nobody serves that unit. - Every **copy** the node held — its complete stored duplicate of a unit — stops being current the moment the process stops, and falls out of **the caught-up set**, the subset of copies current enough to serve a write or be promoted. - Any unit whose remaining current copies are now fewer than the cluster wants will have new ones built somewhere else. That is the same copying an evacuation would have done, except unplanned, unpaced and starting at the second you chose. So "just stop it" does not avoid the data movement. It reorders it: the movement happens *after* the reduction in redundancy rather than before it, which is the wrong way round. ## What an evacuation does instead Evacuating a node means moving every unit of ownership off it and confirming it is empty before the process stops. The node stays up and keeps serving throughout, so nothing is promoted in a hurry and no unit is ever short of a current copy because of you. The order matters more than the mechanism: 1. Declare that this node should hold nothing, and that its units belong elsewhere. 2. Let the receiving nodes fill their copies while the leaving node continues to serve. 3. As each receiving copy becomes current, leadership and ownership move to it — within the same cluster, one unit at a time. 4. When the node holds nothing, stop it. | | Stop the process | Evacuate first | |---|---|---| | Redundancy during the change | drops first, recovers later | never drops | | Promotions | many, all at once, unplanned | one at a time, as each copy becomes current | | Copying performed | the same bytes, afterwards, unpaced | the same bytes, beforehand, paced | | Writers' experience | a gap while owners are promoted | ownership moves under a healthy node | | When you know it is done | you do not | the node holds nothing | ## The completion condition and the duration The two things an evacuation has that a shutdown does not are a **completion condition** and a **duration**. The completion condition is checkable: *this node holds no unit of ownership, and no unit it used to hold is left short of a current copy.* An empty node is not enough on its own — a unit can be off the leaving node while the copy that replaced it is still filling. Both halves have to be true before the process stops. The duration is arithmetic: the records the node holds, divided by the bandwidth the cluster can give to moving them. That is why an evacuation is scheduled rather than issued. A node holding a few hundred gigabytes is a coffee break; a node holding many terabytes is an overnight change, and it is competing for the same disks and links that production is using the whole time. ## Where storage is detached, and where it is not On designs that keep records on shared or remote storage, a leaving node is still evacuated — ownership still has to be handed over, and anything it has not yet pushed to the shared store still has to be settled — but almost no bytes move, so the completion condition is met in seconds rather than hours. The discipline is identical; only the duration collapses. On designs where each node keeps its own copies on local disk, the duration is the whole story. ## Sequencing traps worth naming - **The thin unit.** If some unit of ownership is already down to very few current copies, the node you are about to evacuate may be holding one of them. Evacuating it is still the right move — but the receiving copy has to be built *before* this one goes away, which is precisely the order an evacuation enforces and a shutdown does not. - **Two at once.** Evacuating a batch of nodes together is only safe if no unit loses its remaining current copies and the bandwidth exists to fill the replacements. One after another is slower but bounded. - **The empty-looking node.** A node that appears idle may still be the only current copy of something rarely written. Emptiness of traffic is not emptiness of data. - **The reversal.** Until the node is stopped, an evacuation can be abandoned and the node left in place; after it is stopped, putting it back is a join, with all the copying a join implies.

  • What if the node you must evacuate holds one of the last current copies of some unit of ownership?
    That is the case the ordering exists for. The replacement copy has to be built elsewhere and become current before this node gives the unit up, so redundancy never dips. Stopping the process instead would take that copy away first and leave the unit exposed while the cluster rebuilt it.
  • Can you evacuate several nodes at the same time?
    Only if no unit of ownership is left short of a current copy by the batch, and only if there is bandwidth to fill all the replacements at once. Otherwise the batch competes with itself and with production traffic. Sequential evacuations take longer but the exposure at each step is bounded and known.
  • Is the evacuation reversible?
    Up to the moment the process stops, yes — you can abandon it and leave the node in the cluster, though units already moved off stay moved. Once it is stopped, bringing the machine back is a join, and it starts empty and has to be filled like any newcomer.

saying these in an interview costs you the question

  • Thinks stopping a node deliberately costs less than the node failing
  • Believes removal avoids copying rather than reordering it
  • Calls the removal done when the process has stopped answering
  • Treats an empty-looking node as safe to stop without checking what it holds
  • Assumes redundancy can dip during a planned removal and recover afterwards