skip to content

Before moving key ownership between nodes of a live partitioned tier, what headroom must the remaining nodes hold?

level: seniorimportance: should knowfreq 48%

answer

  1. the move needs room to run
  2. entries exist twice in flight
  3. freed is not returned
  4. plan the peak, not the end state
  5. shrinking is the constrained direction

basics

~20 s

Enough for the entries in flight to exist twice, plus the move's own working room. A move copies before it deletes, memory freed on the source may not return to the operating system, and the move competes with live traffic on nodes that are still serving.

solid answer

~50 s

A move of key ownership is copy-then-hand-over-then-delete, so for the length of the transfer the entries in flight are held on both the source and the destination. The destination therefore needs room for its current share plus the incoming share plus the transfer's own overhead, and the source will not shrink as fast as you expect: its stored-data size falls as entries go, but its resident footprint often stays high because the allocator keeps freed memory for reuse rather than returning it. On top of that, both nodes spend processor time and bandwidth on the move while still serving. The consequence is a sequencing rule: start a move while there is room, not because you have run out. A node at its ceiling mid-move begins removing entries under pressure or refusing writes, depending on the posture it is on — and either way the move caused the incident.

go deeper

for a junior

Know that entries can be moved between nodes while the tier keeps serving, and that the transfer takes time and memory on both ends rather than happening instantly.

for a middle

Explain why the peak is above the end state: entries are held on both nodes while a range is in flight, and freed memory on the source is retained by the allocator rather than returned.

for a senior

Demonstrate the sequencing judgment — move before the ceiling is near, size the destination for the peak, throttle the transfer against live traffic, and confirm completion by ownership rather than by footprint.

for a principal

Set the rule the team plans against: rebalancing budget is capacity you buy in advance, and a tier whose only spare room is its safety margin has no move available to it when it needs one.

Moving which node owns which keys, without taking the tier down, is the operation that separates people who have run a partitioned in-memory tier from people who have read about one. The mechanism that decides *which* node owns a key belongs elsewhere; what matters here is that the move consumes the very resource it is meant to relieve, and that this is what makes the timing decision. ## What is in flight, and where it lives A move proceeds as a transfer, a hand-over, and a release: 1. the destination receives the entries for a range of key assignments while the source keeps serving them; 2. ownership of that range is handed over at a point in time; 3. the source releases the entries it no longer owns. Between step 1 and step 3 the **data exists in two places at once**. Ownership does not: at any instant exactly one node is answerable for a key. Those two facts sit together and are the whole shape of the operation — you are paying for a duplicate copy of the moving range in exchange for never having a moment when nobody owns it. ## Sizing the headroom | what the headroom covers | why it is needed | |---|---| | the incoming share on the destination | it arrives before the source releases it | | the transfer's own working memory | buffers on both sides while entries are in transit | | the traffic the node is already serving | the move does not pause live writes, which keep arriving | | what the allocator will not return | the source's resident footprint lags its stored-data size | The honest figure to plan against is the **peak**, not the new steady state. A plan that says "after the move each node holds a quarter of the keyspace" has described the end, not the middle, and the middle is where the move fails. The direction matters too. Adding a node is the easy direction: the destination is empty, so the headroom question answers itself. **Removing** a node is the constrained one — the survivors must absorb its entries, the per-entry overhead of holding them, the move transient, and its share of the traffic, all at once. Shrinking is where a move gets refused halfway. ## Why the source does not shrink immediately When entries are released, the memory they occupied is returned to the allocator, and the allocator typically keeps it for the next allocation rather than handing it back to the operating system. The node's **stored-data size** — what the server counts for its own entries — falls straight away; its **resident footprint** — what the operating system sees the process holding — may stay near its old value for a long time, and falls only as the freed space gets reused or the process is replaced. So a resident footprint that has not moved is not evidence that the move failed, and it is a poor completion signal. Use ownership reported consistently by every node, and the settled stored-data size, to say the move is done. ## The window where ownership changes hands For as long as a range is in flight, requests for those keys have to reach exactly one owner, and how callers are steered there is a point of genuine variation across this class of store: some designs answer the caller with a redirect it follows, some put a proxy in front that holds or forwards the request so the caller sees only latency, and some briefly fail requests for the moving range and expect the caller to try again. What all of them have in common is that a caller which cannot tolerate one slow or failed attempt for a subset of keys will notice the move. The same window is why hand-rolling a move is a mistake. A copy-then-delete performed by a script has no hand-over point: a write that lands on the source after its key was copied is thrown away when the source's copy is deleted, and while both nodes hold the key, two readers can legitimately see different values. The move has to transfer ownership of the range as a single decision, not key by key at the mercy of whatever traffic arrives. ## What differs between stores - The unit moved: some designs move a whole range of assignments at a time, others move entry by entry, which changes how coarse the interruption is. - Who runs it: on a self-operated tier you trigger and watch the move; on a managed instance it may be a control you press, with neither a view of the transfer nor a way to pause it — in which case your headroom planning has to happen entirely before you press it. - What throttles it: some moves can be rate-limited so they compete less with live traffic, and a slower move is often the right choice on a tier close to its limits. ## The rule that falls out Move while you have room. The capacity trigger for a partitioned tier is not "the ceiling", it is "the ceiling minus the growth expected during the move minus what the move itself consumes" — which is usually a much lower number than teams expect, and the reason a well-run tier is rebalanced during ordinary weeks rather than during an incident.

  • How does removing a node differ from adding one, for headroom?
    Adding is easy: the destination is empty, so it has room by definition. Removing is constrained — the surviving nodes must absorb the leaving node's entries, their per-entry overhead, the transient of the move, and its share of the traffic, while still serving. Shrinking is where moves get refused or start pressure removal.
  • Why not move ownership by copying entries with a script and deleting the originals?
    Because nothing coordinates the hand-over. A write that lands on the source after you copied its key is lost when you delete it, and while both nodes hold a key, two readers may see different values. The transfer has to change ownership of the range as one decision, with requests routed to exactly one owner throughout.
  • What signal tells you the move finished rather than stalled?
    Every node reporting the same ownership for every range, and the source's stored-data size settled at its new share. Resident footprint is a bad completion signal: a node that has handed entries away can hold the same footprint for a long time because the allocator kept the freed memory.

Moving between two flats while you go on living in them: for a few days each box is in both rooms at once — its contents in the new place, its space not yet reclaimed in the old — and only one address is answerable for your post. You cannot start on the day the old flat is packed to the ceiling, because the packing itself needs somewhere to stand.

saying these in an interview costs you the question

  • A move can be started once memory is already at the ceiling
  • Releasing entries on the source immediately returns memory to the operating system
  • Moving ownership is invisible to callers and costs the nodes nothing
  • The destination only needs room for the new steady-state share
  • Copy the entries, delete the originals — the ordering does not matter