skip to content

Membership & Data Movement

What actually has to move when a cluster gains, loses or restarts a node. Probed because most data-loss stories start with a planned change, not with a failure.

part ofBroker & streaming operationsoverview, primer and where to startread it →
on this pageshow

questions

25

A cluster node joins a busy broker cluster and reports healthy — does the load on the existing nodes drop yet?

level: juniorimportance: must knowfreq 72%

answer

  1. membership is not capacity
  2. an announcement, not a transfer
  3. it holds nothing until assigned
  4. relief arrives with the copying
  5. detached storage is the exception

basics

~20 s

No. A newly joined node holds no records, so it serves nothing and relieves nothing until units of ownership are assigned to it and their records have been copied across. Detached-storage designs are the exception.

solid answer

~50 s

Joining a cluster is an announcement, not a transfer. The new node is now a member, but it holds no **unit of ownership** — no partition (or queue) — so it answers neither reads nor writes, and every existing node is working exactly as hard as before. Relief needs two more steps: a reassignment that says which units should move onto it, and then the copying of those units' records, which is bytes over available bandwidth and can run for hours on a large node. Only once a copy is current enough to join the caught-up set can it serve. Two exceptions: on designs where records live on shared or remote storage, ownership can move without copying and the node serves almost at once; and where new units are created continuously, placing the new ones there relieves the cluster gradually with no copying.

go deeper

for a junior

Remember the order: a node joins, then units of ownership are assigned to it, then its records are copied, and only then does it serve. Membership going up by one changes nothing on its own.

for a middle

Explain what has to be filled before a copy may serve — it must become current enough to join the caught-up set — and that the duration is simply bytes over the bandwidth available for copying them.

for a senior

Show that you plan for the lead time: you expand before the need, you know the copying competes with production reads on the same disks, and you can say what the cluster looks like mid-fill.

for a principal

The angle to own is whether the estate should keep paying this lead time at all: where records live on shared storage, adding capacity becomes an ownership change of seconds instead of a data movement of hours.

## Membership and capacity are two different events When a cluster node starts and joins a broker or streaming cluster, it announces itself: the cluster records that a machine is present, reachable and willing to hold data. That is all joining does. It does not hand the newcomer a single **unit of ownership** — a partition (or queue) that exactly one node serves at a time. Until some unit is assigned to it, the node stores no records, accepts no writes and answers no reads. This is the gap that surprises people. The machine is provisioned, the process is healthy, the member count went up by one — and every existing node is working exactly as hard as it was a minute ago. The member list grew; the capacity did not. ## Why the relief arrives late, and costs something on the way Getting work onto the newcomer takes two further steps. 1. **Declare** a new mapping from units of ownership to the nodes that should hold them — a **reassignment plan**. Some platforms expect an operator to produce it, some compute one themselves, and some hosted ones do it invisibly. 2. **Execute** it. Executing means the newcomer fetches the records of each unit assigned to it from whichever node holds them now, until its **copy** — one complete stored duplicate of that unit — is current enough to join **the caught-up set** and be allowed to serve. Step 1 takes a second. Step 2 takes the bytes involved divided by whatever bandwidth is available for copying them, and it reads those bytes off disks that are already serving production traffic. A node that will eventually hold terabytes is useless for hours. Worse, while the copying runs the cluster is doing *more* work than before, not less: the existing nodes serve their normal traffic plus the outbound copying. | Stage | Typical elapsed | What the new node serves | |---|---|---| | Process up, listed as a member | seconds | nothing | | Units of ownership assigned to it | seconds | nothing yet | | Its copies are being filled | minutes to hours | nothing yet | | Its copies reach the caught-up set | — | reads, and leadership may move to it | Only the last row is relief. ## The first exception: detached storage Some designs keep a unit's records on shared or remote storage rather than on a node's local disk. There, ownership of a unit can be handed to the newcomer without copying anything: the records already sit where the new owner can read them, so it begins serving almost at once. Anything the previous owner has not yet pushed to the shared store still has to be settled, and caches start cold, so it is not literally free — but it is a change of seconds rather than a data movement of hours. That is the real fork in this subject: **where records live on local disks, capacity follows data; where storage is detached, capacity follows ownership.** ## The second exception: letting new work land there Not all relief requires moving anything. Where units of ownership are created continuously — new queues, new streams, a new period's units — a platform can place the *new* ones on the newcomer and leave existing ones untouched. Relief then arrives in proportion to how much of the traffic is new, with no copying at all. It is slow, cheap and only ever partial, so it complements a reassignment rather than replacing one. Queue-shaped brokers where a queue lives on its owning node behave the same way: the newcomer takes the queues placed on it, never the ones already elsewhere. ## What to do with this as an operator - **Expand ahead of the need, not at the moment of it.** The lead time is the copying, and it is longest exactly when the cluster is busiest. - **Measure relief in units served and catch-up distance**, not in node count. Catch-up distance is how far behind its leader a copy still is; while it is large, that copy consumes capacity rather than providing it. - **Expect the bill before the benefit.** The machine is charged from the moment it boots and contributes from the moment its copies are current. - **Do not read an empty newcomer as a broken one.** A node holding nothing an hour after joining is normal if nothing was assigned to it; it is a fault only if a reassignment was declared and no records are arriving. - **Say which of the two shapes you are on** before promising a timeline. The same sentence — "we added a node" — means hours on one design and seconds on another. The compact model: joining is an announcement, assignment is a decision, copying is the work, and capacity appears at the end of the third.

  • If the platform keeps creating new partitions (or queues) rather than moving existing ones, does the newcomer help sooner?
    Yes, but only gradually and only partially. New units of ownership can be placed on it at creation time, which needs no copying at all, so relief arrives in proportion to how much of the traffic is new. The existing hot units stay exactly where they are until a reassignment moves them.
  • How do you tell whether the new node is actually taking load yet?
    Count the units of ownership it currently serves, and check how far behind their leaders its copies still are. Membership tells you nothing; a node can sit in the member list for hours holding nothing. A copy that is still filling is consuming cluster capacity, not adding any.

A new hire joins an overloaded support rota on Monday. They are on the roster immediately, but they answer no tickets until work is handed over — and the handover is done by the very people who are already drowning, so the team gets slower before it gets faster.

saying these in an interview costs you the question

  • Thinks a node starts taking traffic the moment it is listed as a member
  • Assumes capacity rises when the machine is provisioned
  • Believes ownership moves onto a new node without any byte copying on local-disk designs
  • Treats adding a node as the immediate fix for a cluster that is overloaded now
  • Reads an empty newcomer as a fault rather than as an unassigned node
open as a page

When the leader of a partition (or queue) dies, what do writers see before another copy is promoted within the same cluster?

level: juniorimportance: must knowfreq 62%

basics

~20 s

Writes to that partition fail rather than block: the client's cached owner map still names the dead node, so sends are rejected until a replacement copy is promoted. Records survive only if the writer's retries outlast that gap.

open as a page

Each cluster node holds an equal number of copies, yet one node serves most reads and writes — why?

level: juniorimportance: must knowfreq 60%

basics

~20 s

Leadership is a second distribution, separate from where copies sit. Where exactly one node serves each partition (or queue), the work follows leadership rather than storage, so perfectly even copies can still leave one node leading — and serving — most units.

open as a page

A reassignment plan moving a partition (or queue) to another cluster node is accepted in a second — why is the move not finished?

level: juniorimportance: must knowfreq 70%

basics

~20 s

A reassignment plan only declares which node should hold which partition (or queue). Accepting it starts ordinary copy traffic that streams the unit's retained records to the new node, and ownership settles only once that copy has caught up.

open as a page

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%

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.

open as a page

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%

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.

open as a page

What does a cluster gain by stamping writes to a partition (or queue) with a number raised at each leadership change?

level: middleimportance: must knowfreq 62%

basics

~20 s

A number raised at every change of leadership makes a superseded leader's writes rejectable: whoever receives a write compares its stamp with the newest it knows and refuses anything older, without first having to establish whether the sender is alive.

open as a page

A reassignment's copy traffic shares disks and links with live traffic — what is the move's bandwidth ceiling for, and what breaks at each extreme?

level: middleimportance: must knowfreq 62%

basics

~20 s

The ceiling caps how much bandwidth a migration's copy traffic may consume, because that traffic competes with production. Set it too high and live reads and writes slow down; set it too low and the cluster sits half-migrated indefinitely, holding copies on both node sets.

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

The leader of a partition (or queue) is lost and no copy in the caught-up set is reachable; what does promoting a copy outside that set, within the same cluster, cost you?

level: seniorimportance: must knowfreq 66%

basics

~20 s

It costs acknowledged records. A copy outside the caught-up set is missing writes the cluster already confirmed to producers, and promoting it discards them silently — no error reaches anyone. The alternative is keeping every record and staying unwritable for an unknown time.

open as a page

A cluster node leading a partition (or queue) goes silent: how does the cluster tell a crash from a network cut-off?

level: juniorimportance: should knowfreq 58%

basics

~20 s

The cluster cannot tell them apart. A crashed node and a node cut off by the network produce the same observation from outside — silence — so the old leader has to be treated as possibly alive and still accepting writes.

open as a page

How long should the cluster tolerate silence from one of its nodes before promoting new leaders, and what does shortening that tolerance cost?

level: middleimportance: should knowfreq 48%

basics

~20 s

Tolerance must exceed the node's ordinary pauses — garbage collection, a slow disk flush, a congested link — or the cluster promotes replacements for nodes that are alive. Shortening it shortens outages but buys needless promotions, each with its own gap and client churn.

open as a page

Why does leadership pile onto a few cluster nodes after failures and restarts, and why does it stay there?

level: middleimportance: should knowfreq 52%

basics

~20 s

Every stop hands a unit's leadership to a surviving copy that is caught up, and a node that comes back rejoins as an ordinary copy. Leadership accumulates on whichever nodes happened to be alive at the wrong moments and stays until something hands it back.

open as a page

Your cluster is saturated at peak and you add two nodes to relieve it — why does the incident get worse before it gets better?

level: seniorimportance: should knowfreq 57%

basics

~20 s

Because relief has to be copied onto the new nodes, and that copying reads records off the very nodes that are already at their ceiling. The cluster does extra work for hours and gains nothing until the new copies are current.

open as a page

A writer's cached owner map still names the superseded leader for a partition (or queue): why does its next write not quietly succeed?

level: seniorimportance: should knowfreq 47%

basics

~20 s

Because authority is checked where the write must become durable, not where the client sends it. The stale route is refused under an out-of-date generation number, the error travels back, and the client's cached ownership view is corrected rather than the write silently landing on a node nobody reads from.

open as a page

Correcting leadership skew copies no bytes — so what does a leadership-only rebalance actually cost, and when should it run?

level: seniorimportance: should knowfreq 48%

basics

~20 s

A leadership-only rebalance changes which node serves each unit and moves no records, so it costs no bandwidth. It costs a brief interruption per unit as service hands over and clients refresh, which means the real decision is how many units to hand over at once, and when.

open as a page

A move stays inside its bandwidth ceiling, yet reads on unrelated partitions on the same cluster nodes slow down — where is the contention?

level: seniorimportance: should knowfreq 52%

basics

~20 s

A bandwidth ceiling caps migration bytes on the wire, not the work that produces them. The source node must read cold historical records off disk, which competes with live reads and can evict the warm data they were being served from; the destination pays a sustained write.

open as a page

A reassignment has run for hours with many units still in the intermediate mapping — how does cancelling it differ from walking away?

level: seniorimportance: should knowfreq 45%

basics

~20 s

The intermediate mapping is real cluster state, not an operator's intention. Walking away leaves it in force, holding copies on both node sets indefinitely and blocking later plans; cancelling is a deliberate operation that reverts the mapping, discards partial copies and moves finished units back.

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

Should an estate move to detached storage so that adding or removing a cluster node stops costing copy traffic at all?

level: principalimportance: should knowfreq 38%

basics

~20 s

Sometimes. Detaching storage turns every membership change from a data movement of hours into an ownership handover of seconds, which is worth a great deal to an elastic estate — but it makes one shared store a dependency of every node, and it rarely removes the recent records that still live locally.

open as a page

Should leadership restoration to home nodes run continuously on its own, or only when an operator decides to run it?

level: principalimportance: should knowfreq 35%

basics

~20 s

Neither extreme survives contact. Continuous restoration prevents skew accumulating but hands traffic back unattended, including to a node that just returned unwell. Operator-run restoration is predictable but relies on someone noticing an invisible drift. The workable policy is automatic above a threshold, with health preconditions and an alert.

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

Should a cluster node that has lost contact with the coordination role stop serving writes on its own, or rely on later refusal?

level: principalimportance: nice to knowfreq 34%

basics

~20 s

Decide by what a misleading answer costs. Self-demotion after a silence period shortens the time a client can believe it is writing successfully, at the price of stepping aside during blips; relying only on refusal keeps serving through blips but leaves work provisional for longer.

open as a page

Who decides whether a cluster may promote a copy that is behind the lost leader, and what should that policy depend on?

level: principalimportance: nice to knowfreq 34%

basics

~20 s

The accountable owner of each stream decides, in advance and per stream, not the operator during the incident. The policy should turn on what the stream carries, whether it can be republished from upstream, and what an outage costs relative to silent loss.

open as a page

Your estate spends days on reassignments whenever hardware changes — what does holding records on shared or remote storage change about a move's cost, and what does it not?

level: principalimportance: nice to knowfreq 38%

basics

~20 s

Where records live on shared or remote storage rather than a node's local disk, changing which node owns a unit becomes a metadata handover with no bytes in flight: no bandwidth ceiling, no half-migrated state. Placement decisions, cold starts and a new shared dependency remain.

open as a page