skip to content

Spreading the Tier

Two answers to one machine not being enough - copies of the whole keyspace, and the keyspace split across nodes - plus what changes when the machines sit in different places. Each fails differently.

on this pageshow

questions

30

A tier whose nodes each hold a share of the keyspace, not a copy: what limits an operation naming three keys?

level: juniorimportance: must knowfreq 58%

answer

  1. two ways to add machines
  2. each key has exactly one home
  3. one operation runs on one node
  4. co-locate, split, fold, or accept
  5. refusal, composed reply, or never offered

basics

~20 s

One operation runs on one node, so it can only name keys that are all in the same partition. Split across partitions, the call is refused, composed by an intervening layer, or has to be split by the caller.

solid answer

~50 s

When the keyspace is cut into partitions rather than copied, every key belongs to exactly one partition and a partition is served by one node at a time. A node executes against memory it owns, so an operation naming several keys can only run if all of them are in one partition. What the caller then sees varies: where the caller holds the partition map the store declines the call outright, an intervening proxy may instead read each partition and compose one reply that was never indivisible, and some stores never offered a multi-key operation at all, so the library has always issued one call per key. The responses are to co-locate the keys deliberately, split the call and compose in the caller, fold the group into one entry, or accept that it is no longer one unit.

go deeper

for a junior

Recall the one rule: a key lives in exactly one partition, and an operation naming several keys needs them all in the same one. Say which arrangement you are talking about — copies of the whole keyspace, or a share each — before answering.

for a middle

Explain why the limit exists: a node works on memory it owns, and these stores keep coordination out of the operation path on purpose. Then name the responses — co-locate, split and compose, fold into one entry, or accept a torn view.

for a senior

Show that you know what the caller actually experiences differs by routing shape: an outright refusal, a composed reply that was never indivisible, or a library that always issued one call per key. Say which of those a design is silently relying on.

for a principal

Frame it as a portability decision. Cross-partition behaviour is exactly the kind of guarantee that changes when the tier is swapped for a hosted or proxied equivalent, so a design that depends on it is a design that cannot move.

## Two different meanings of "we added nodes" Adding machines to a volatile tier means one of two things, and the two behave in opposite ways. In the first, each added node holds **a copy of the whole keyspace**: every node can answer for every key, and what you gain is read capacity and something to promote when the primary is lost. In the second, the keyspace is **cut into partitions**, and each node holds a share of it: a key belongs to exactly one partition, and a partition is served by one node at a time. Only the second arrangement raises this problem. It is worth saying out loud in an interview, because "we scaled the tier out" describes both and they fail in opposite directions — copies serve answers that are behind and lose recent writes on a promotion, while pieces refuse multi-key work and misroute when a caller's copy of the partition map is out of date. ## Why one operation cannot reach two partitions A node executes an operation against memory it owns. For one operation to read or write keys living on two nodes, something would have to fetch the remote values, combine them, and hold both nodes still while it did so. Stores in this class deliberately do not put that in the path. Their whole value is that an operation is a few microseconds of work against one machine's memory with no coordination in it; making every operation able to span nodes would make the common case pay for a case most workloads can design away by choosing key names better. Making one change across independent parties atomic is a genuine subject with genuine protocols — it is simply not what this tier offers. So the rule is blunt: **an operation that names several keys can run only if every one of those keys is in one partition.** Anything else the store executes as a single unit inherits the same rule — a group of steps submitted together, or a script the server runs on your behalf — and the store can only enforce the rule if it is told every key the unit will touch before it begins. ## What the caller actually observes — and this varies | How a key is resolved to a node | A multi-key operation spanning partitions | What the caller sees | |---|---|---| | A map held by the caller | The store declines it | A cross-partition refusal: an error on a path that worked yesterday | | An intervening proxy | Some decline; some read each partition and compose | Either a refusal, or a reply that was never one indivisible operation | | A directory consulted per key | Usually there is no such operation to offer | The caller resolves each key and issues one call per key | | Independent nodes, spread by the caller's library | Not offered at all | Nothing changes: the library always issued one call per key | The row that catches people is the second. A composed reply is neither a refusal nor an indivisible read: the values came from different nodes at slightly different moments, so the group can be observed part-old and part-new. A design that quietly leans on that behaviour also breaks if the routing shape is ever changed. ## The four honest responses - **Co-locate deliberately.** Put a marker in the key names that the placement rule uses instead of the whole key, so related keys land in one partition on purpose. The operation keeps working; the group can then never be spread. - **Split the call.** Issue one call per key and compose in the caller. It works everywhere, and the values are no longer observed at one instant. - **Fold the group into one entry.** One key is one partition by construction. The entry grows, and on stores that only hand back opaque bytes every change becomes a read of the whole value and a write of the whole value. - **Accept that it is no longer one unit.** Sometimes the right answer is to decide explicitly what a reader does when it sees part of the group changed and part not. ## What the constraint is not It is not about load: the call is not refused because the tier is busy, and adding nodes does not relieve it — adding nodes makes it more likely, because related keys spread further apart. It is not about the size of the values. And it does not arise at all on a tier of copies, where every node holds the whole keyspace: there nothing is refused for this reason, and the caller's problem is an answer that is behind, which is a different subject entirely. Naming the arrangement before answering is most of the mark on this question.

  • Does the same limit apply when the added nodes each hold a copy of the whole keyspace?
    No. Where every node holds the whole keyspace, any node can serve any key, so no operation is refused for where its keys live. That arrangement has a different problem — a copy that is behind the primary — and confusing the two is the usual failure on this question.
  • If an intervening proxy composes a multi-key read from several partitions, what has the caller lost?
    Indivisibility. The values were read from different nodes at slightly different moments, so the caller can see part of the group as it was before another writer and part as it was after. The call also stops working the day the routing shape changes to one that refuses instead.

saying these in an interview costs you the question

  • Thinks adding nodes only adds capacity and leaves existing calls working
  • Believes a node will fetch the missing keys from its peers for you
  • Says any node can serve any key once the tier has several nodes
  • Thinks splitting the call into single-key calls keeps it indivisible
  • Blames the refusal on load rather than on where the keys live
open as a page

A volatile tier acknowledges a write before any replica holds a copy of it: what has the caller been promised?

level: juniorimportance: must knowfreq 66%

basics

~10 s

Acknowledgment means one node has the write in memory, nothing more. Until it propagates, the write exists in exactly one place, and if that node is lost the write is simply gone.

open as a page

In a volatile tier with a primary and two replicas of the whole keyspace, why must a replica be promoted before writes resume?

level: juniorimportance: must knowfreq 62%

basics

~20 s

A replica holds a copy of the keyspace but is not the address that accepts writes, so writes resume only after some deciding party promotes one of them to primary. That promotion is a deliberate step, and it takes time.

open as a page

When a keyspace is split so each key lives on one node, why must every caller resolve keys the same way?

level: juniorimportance: must knowfreq 68%

basics

~20 s

Assignment must be a deterministic function of the key, evaluated identically by every caller. If two callers disagree about where a key belongs, each reads and writes its own copy on a different node and neither sees the other's.

open as a page

A caller updates a key on the primary of an in-memory store, then reads it from a replica and sees the previous value — why?

level: juniorimportance: must knowfreq 62%

basics

~20 s

A replica answers from its own copy, which trails the primary by the propagation lag, so a read issued inside that window returns the value the key held before the write. The write is not lost, only not yet visible there.

open as a page

With the keyspace split across partitions and no way to co-locate, how do you rework an operation over four related keys?

level: middleimportance: must knowfreq 62%

basics

~20 s

Three reworks are honest: one call per key composed in the caller, folding the four values into one entry under one key, or accepting the operation is no longer one unit and saying what a reader does with a mixed view.

open as a page

A service expands into a second region but leaves its volatile tier in the first: what does each lookup now cost?

level: middleimportance: must knowfreq 54%

basics

~20 s

Each lookup turns into a cross-region round trip - tens of milliseconds instead of a fraction of one - which can exceed the cost of the work the tier was avoiding, making the remote path slower than having no tier.

open as a page

A tier stops acknowledging a write until at least one replica holds it: what does that buy, and what does it cost?

level: middleimportance: must knowfreq 58%

basics

~20 s

Waiting for a copy narrows the un-propagated window and adds a network hop to every write. It does not close the window: a write only one copy holds is still lost if that copy is not the one promoted.

open as a page

An automatic failover of a volatile tier's primary is often reported as one number; which separate windows does that total outage actually consist of?

level: middleimportance: must knowfreq 58%

basics

~10 s

A failover outage is the sum of three windows: detection, before anything concludes the primary is gone; promotion, while a replica is reconfigured; and rediscovery, until the last caller stops addressing the dead node.

open as a page

Which component resolves a key to its node in a partitioned in-memory tier, and what does each option cost?

level: middleimportance: must knowfreq 58%

basics

~20 s

Three shapes exist: the caller holds a copy of the partition map, an intervening proxy holds it, or a directory is consulted per lookup. They differ in hops, in who must learn of a change, and in who can be stale.

open as a page

In an in-memory store replicated as one primary and two replicas, why is promotion gated on a majority of the deciding members?

level: middleimportance: must knowfreq 60%

basics

~20 s

A majority gate guarantees that only one side of a network partition can promote, because two separated sides cannot both hold more than half of one membership. Without it, both halves promote, both accept writes, and the two keyspaces diverge.

open as a page

After assignments move, what does a caller with a stale partition map observe, and why is that worse than an error?

level: seniorimportance: must knowfreq 52%

basics

~20 s

It depends on whether the node reached checks ownership. Where it does not, a misrouted read returns an ordinary absence and the caller treats live state as gone; where it does, a reply names the right node to follow and learn from.

open as a page

Two nodes of an in-memory store each accepted writes as primary during a network partition, so when it heals why is one side usually discarded rather than merged?

level: seniorimportance: must knowfreq 52%

basics

~20 s

The tier keeps the current value of each key and no history of how it got there, and where values are opaque bytes there is no merge rule. So the losing node discards its keyspace, copies the winner's, and its acknowledged writes vanish.

open as a page

An in-memory store is overloaded, so replicas each holding the whole keyspace are added — which load does that relieve, and which not?

level: middleimportance: should knowfreq 50%

basics

~20 s

Fan-out across copies buys read request capacity and nothing else. Writes still land on one primary and are then propagated to every copy, and each copy holds the whole keyspace, so neither write throughput nor memory headroom improves.

open as a page

A naming convention pins every key of one tenant to a single partition: what does that buy, and what does it cost?

level: seniorimportance: should knowfreq 48%

basics

~20 s

It buys multi-key work over that tenant's keys after the split. It costs placement freedom: the group is placed whole, the largest tenant sizes a node, and one node's loss now removes everything for a few tenants rather than a slice of everyone.

open as a page

Two regions each run their own volatile tier over one shared system of record and the same key holds different values in each - why is that usually accepted?

level: seniorimportance: should knowfreq 42%

basics

~20 s

Both values are locally derived views of one shared system of record, so neither is authoritative and there is nothing to merge. Making them agree would put a cross-region exchange back on the write path.

open as a page

A per-user quota counter and a deduplication record live only in each region's own volatile tier while callers can be routed to either region - what breaks?

level: seniorimportance: should knowfreq 46%

basics

~20 s

Neither entry is a copy of anything, so two independent tiers hold two different facts about one user: the quota is granted once per region, and a retry landing in the other region finds no record and is treated as new.

open as a page

A tier refuses writes when fewer than a set number of replicas are current: what does that gate buy, and what does it not?

level: seniorimportance: should knowfreq 38%

basics

~10 s

A minimum-healthy-copies write gate converts a silent loss risk into a loud refusal: writes are accepted only while enough copies are current. It narrows the un-propagated window and does not close it.

open as a page

A tier acknowledges writes before a replica holds them: how do you put a defensible number on what a promotion would vaporise?

level: seniorimportance: should knowfreq 47%

basics

~20 s

Multiply the tail propagation lag at peak write rate by that write rate to get a count of writes at risk, then say what those writes were. The count alone is not an answer; the state they held is.

open as a page

A volatile tier's promotion finished in seconds, yet one service kept addressing the dead primary for an hour — what did it get wrong?

level: seniorimportance: should knowfreq 52%

basics

~20 s

Promotion only moves which node accepts writes; it does not reach into callers. A service that resolved the write address once at start-up and never again keeps talking to the dead node, however clean the failover was.

open as a page

A healthy primary in a volatile tier stalls for four seconds and the tier promotes a replica; what did that cost, and how do you choose the detection window?

level: seniorimportance: should knowfreq 46%

basics

~20 s

A stall longer than the detection window is indistinguishable from death, so the tier pays a full outage for nothing and must then demote a node that returns believing it leads. Size the window above the measured stall tail.

open as a page

In a volatile tier that fails over without a human, which parties can be entitled to declare the primary dead, and why is one observer never enough?

level: seniorimportance: should knowfreq 50%

basics

~20 s

Three parties can own the decision: the surviving peer nodes, a set of dedicated observer processes, or a controller outside the data plane. A single observer is never enough — it cannot tell a dead node from its own severed link.

open as a page

An in-memory store holding both re-derivable values and sole-copy records serves all reads from replicas; which reads must return to the primary, and how do you bound the rest?

level: seniorimportance: should knowfreq 46%

basics

~20 s

Reads whose correctness is the presence or absence of a key at this instant cannot be answered by a lagging copy: route that class to the primary. Every other read gets a named staleness budget, measured against propagation lag rather than assumed.

open as a page

In a replicated in-memory store cut in two by a network partition, why must the minority side refuse writes or demote itself unprompted?

level: seniorimportance: should knowfreq 50%

basics

~20 s

Because the majority side cannot reach it to tell it anything. Safety depends on the isolated primary acting against itself on a local timer, refusing writes or demoting, since a promotion is very likely already happening on the other side.

open as a page

Several teams share one volatile tier that will be split into partitions next quarter; what standing position do you set for multi-key designs?

level: principalimportance: should knowfreq 38%

basics

~20 s

Rule: no design may assume one address space. Every call naming more than one key is classified - co-located under a capped group, split, folded, or explicitly accepted as non-indivisible - and reworked before the split.

open as a page

Where should the partition map live for tiers reached by services in several languages, and what should you refuse to depend on?

level: principalimportance: should knowfreq 38%

basics

~20 s

Decide by who bears the cost of a change, not by hop count. Caller-side maps spread correctness across every client version deployed; a proxy concentrates it into one component you run; a directory is authoritative but priced per lookup.

open as a page

One in-memory store spans two data halls and the business wants writes accepted in both during a network partition, so what do you tell them and what do you require first?

level: principalimportance: should knowfreq 40%

basics

~20 s

For one keyspace the request cannot be granted safely: during a split either one hall stops writing, or both write and one side's acknowledged writes are later discarded unmerged. Pick which, deliberately, after inventorying which keys are the only copy.

open as a page

Your product runs an independent volatile tier per region and one region is lost, with its traffic shifted to the survivor - what should the team expect?

level: principalimportance: nice to knowfreq 31%

basics

~20 s

Three separate losses, not one: derived entries for the arriving population are absent and get rebuilt; sole-copy state held only in the lost tier is unrecoverable; and if that region accepted the writes, the write path is a second, larger failure.

open as a page

Your organisation runs volatile tiers on several products whose acknowledgment postures differ: what standing rule do you set, and what do you refuse to depend on?

level: principalimportance: nice to knowfreq 29%

basics

~10 s

Standardise the number, not the mechanism: each tier declares the un-propagated window it accepts and proves it with a real promotion. Refuse to depend on guarantees the fleet's products implement differently.

open as a page

An in-memory store got replicas for read capacity, yet half its reads are pinned to the primary for correctness — what is your call?

level: principalimportance: nice to knowfreq 30%

basics

~20 s

Fan-out only buys capacity for reads you will answer from a lagging copy, so each pin returns part of it. Measure which classes are pinned and why, then change what the tier holds or how it is split rather than adding copies.

open as a page