skip to content

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

level: juniorimportance: must knowfreq 60%

answer

  1. two distributions, not one
  2. storage even, work uneven
  3. only the serving node takes the writes
  4. count units led, per node

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.

solid answer

~50 s

Two things are spread across a cluster and only one of them shows up on a storage chart. Copies decide which nodes *hold* each unit of ownership — a partition, or a queue on queue-shaped brokers. Leadership decides which single node currently *serves* it. On designs where one node serves a unit at a time, that node takes the unit's writes and feeds the other copies, and on many designs it answers the reads too. Leadership changes every time a node dies, is restarted or is declared unreachable, and it is not handed back unless something hands it back. So after a few such events a node can hold its fair third of the copies while leading three quarters of the units: even copies, hot node. Cluster-wide sums average this away — only a per-node count of units led shows it.

code

json · 14 lines
json
{
  "clusterTotals": {
    "nodes": 4,
    "unitsOfOwnership": 48,
    "copiesHeld": 144,
    "averageCopiesPerNode": 36
  },
  "perNode": [
    { "node": "n1", "copiesHeld": 36, "unitsLed": 31 },
    { "node": "n2", "copiesHeld": 36, "unitsLed": 9 },
    { "node": "n3", "copiesHeld": 36, "unitsLed": 5 },
    { "node": "n4", "copiesHeld": 36, "unitsLed": 3 }
  ]
}

go deeper

for a junior

Recall that holding a copy and serving it are different jobs. Where one node serves each partition (or queue), only that node takes its traffic, so equal storage does not mean equal work.

for a middle

Explain both distributions and what the serving node does per unit it leads: takes the writes, feeds the other copies, and on many designs answers the reads. Then explain why a cluster-wide sum cannot reveal an uneven distribution.

for a senior

Show how you confirm it during an incident: a per-node count of units led beside units held, not an average, plus the observation that the saturated node is now the one whose loss costs most.

for a principal

Frame it as a planning risk. Capacity sized on averages hides a machine running at several times its share, and the blast radius of losing it exceeds what the copy count suggests. Decide which per-node view the estate is required to publish.

## Two distributions live on one cluster A cluster spreads two different things across its nodes, and they are easy to confuse because both are described as "balanced". - **Data placement** — which nodes hold a **copy** of each **unit of ownership**: a partition, or on queue-shaped brokers a queue. This is a byte-level fact; changing it means moving records between machines. - **Leadership** — which single node currently **serves** that unit. On designs where exactly one node answers for a unit at a time, this is the assignment that decides where the work lands. Changing it moves no records at all. A cluster can be perfectly balanced on the first and badly lopsided on the second. That is leadership skew: the storage chart is level, one machine is saturated. ## Why the serving node is the hot one On a leader-based design the node serving a unit does work the other copies do not: - it accepts every **write** for that unit and orders it; - it **feeds** the other copies, so its outbound traffic is roughly the write rate multiplied by the number of copies it must serve; - on many designs it also answers every **read** for that unit, though designs differ — some let a reader be served by a nearby copy instead; - it holds the per-unit bookkeeping the other copies do not need to maintain. So a node's real load is roughly proportional to the traffic of the units it **leads**, not to the records it **stores**. Leading thirty-four of forty units means serving most of the cluster's traffic from one machine, however level the disks look. ## Why the totals look healthy Almost every default view is a cluster-wide sum or average, and a sum cannot show a distribution. | View | What it shows | What it hides | |---|---|---| | Total bytes stored | The estate is not filling up | Which node is doing the serving | | Copies per unit | Nothing is under-replicated | That one copy of each is the one working | | Cluster write rate | Aggregate demand | That one node absorbed most of it | | Per-node units **led** | The actual work distribution | Nothing relevant here — this is the view | The fix for visibility is unglamorous: publish, per node, the count of units it **holds a copy of** beside the count of units it **currently leads**. A cluster of five nodes where one leads sixty per cent of the units is skewed no matter how green the totals are. ## How the model varies between designs This is a property of designs that elect one serving node per unit. It is not universal: 1. **Leader-based split streams** — one node serves each unit, the others hold copies and catch up. Skew is at its sharpest here. 2. **Majority-write designs** — a write is accepted once enough copies take it; there is still typically a coordinating node per unit, so concentration is possible but the write path is less lopsided. 3. **Detached storage** — records live on shared or remote storage, so ownership can move without copying. Serving still concentrates, but the node is doing less per unit. 4. **Queue-shaped brokers where consumers compete** — work is pulled by whoever is free, so there is no single serving node per queue to concentrate, though the queue itself may still be homed on one node. Stating which of these you mean is half the answer in an interview; asserting the first as though it were the class is the classic mistake. ## What the skew actually costs 1. **Tail latency on one machine.** Response times for every unit that node leads share one CPU, one network card and one set of disks. 2. **Headroom that is gone where it is needed.** Capacity planned on cluster averages says there is room; the saturated node says there is not. 3. **A blast radius larger than the copy count implies.** Losing the node that leads most units interrupts most of the traffic at once, even though only its share of the data needs replacing. 4. **A misdiagnosis waiting to happen.** The saturated node invites the conclusion that some client or key is hot, and the team goes looking for a workload problem that does not exist. The correction is an assignment change, not a data migration — which is why it is cheap, and why noticing it is the hard part.

  • Does the node leading most of the units also store more data than its peers?
    Usually not. Leadership is an assignment, not a byte count: a node can lead most units while holding exactly its share of copies, which is why a storage chart stays level while the machine saturates. It is also why the correction is an assignment change rather than a data migration.
  • If reads can be served by a copy that is not the serving node, does skew stop mattering?
    It softens rather than disappears. Designs differ: some let a reader be served by a nearby copy, others send every read to the serving node. Writes still land on one node per unit, and that node still feeds the other copies, so a read-heavy workload spreads out while a write-heavy one stays concentrated.

A filing company with five branches keeps exactly the same number of boxes in each — but the phone system routes every incoming call to one branch. The storage audit is perfect; one office is drowning. The fix is rerouting calls, not moving boxes.

saying these in an interview costs you the question

  • Equal copy counts must mean equal load on each node
  • Leadership follows storage, so balanced copies balance traffic
  • A saturated node can only mean a hot key or a heavy tenant
  • Cluster totals would have shown an imbalance like this
  • Leadership settles back where it started once traffic calms down