Key-value store
Design a key-value store that holds more data than one machine and takes writes even when nodes fail. The design is leaderless: any node coordinates, and N replicas share each key.
- 1A hash ring with many tokens per node places N = 3 copies and moves about 1/(n+1) of the data on a join.
- 2R + W > N gives fresh reads; a sloppy quorum with hints keeps writes available when replicas are down.
- 3Version vectors keep concurrent writes as siblings, so no write is lost.
- 4Three repair paths: hinted handoff, read repair, and Merkle-tree anti-entropy.
Design a distributed key-value store: put and get by key, more data than one machine holds, and never down.
- functional
- put(key, value), get(key), delete(key).
- sizes
- Keys up to 256 B; values up to 1 MB, 1 KB typical.
- available
- Writes succeed while 1 zone of 3 is down.
- latency
- p99 under 10 ms inside one region.
- durable
- An acknowledged write survives the loss of 1 node.
- consistency
- Tunable per request; eventual by default.
I will design a leaderless store: N = 3 copies, tunable quorums, and repair in the background.
| question | answer assumed | it decides |
|---|---|---|
| How many keys, how large? | 10 billion keys, 1 KB values | Number of nodes |
| Reads and writes a day? | 2 billion reads, 0.5 billion writes | Nodes for throughput |
| Can a read be a little stale? | Yes, unless the client asks for R + W > N | Quorum sizes |
| Must writes succeed during a failure? | Yes: never down for writes | Leaderless, sloppy quorum |
| Range scans or only point lookups? | Point lookups only | Hash partitioning |
| Two writers on one key at once? | Rare, but no write may be lost | Version vectors, siblings |
| One region or many? | One region, 3 zones | Replica placement |
I ask about size, read-write mix, consistency and range scans first, because each answer changes the partitioning or the replication.
- raw data
- 10¹⁰ keys × 1,100 B (1 KB value, 100 B key and metadata) = 11 TB.
- copies
- × N = 3 = 33 TB.
- on disk
- × 1.08 LSM space amplification (D1 lab) ≈ 35.6 TB.
- nodes
- 2 TB per node, a choice that keeps a rebuild short: 35.6 / 2 ≈ 18, 6 per zone.
- reads
- 2 × 10⁹ / 86,400 s ≈ 23,148 a second; 3× peak ≈ 69,444.
- writes
- 0.5 × 10⁹ / 86,400 s ≈ 5,787 a second; 3× peak ≈ 17,361.
- per node
- Peak writes × 3 copies / 18 ≈ 2,894; peak reads × R = 2 / 18 ≈ 7,716.
- bandwidth
- Peak: 52 MB/s of replica writes and 139 MB/s of replica reads, for the cluster.
Storage decides the node count here, not throughput. One lab node takes about 30,600 writes a second.
About 36 TB on disk after 3 copies, so 18 nodes at 2 TB each, and about 3,000 replica writes a second per node at peak.
| call | success | errors |
|---|---|---|
| GET /v1/kv/{key}?r=2 | 200 values[] + context | 404, 503 < R |
| PUT /v1/kv/{key}?w=2 + X-Context | 204 | 413, 503 < W |
| DELETE /v1/kv/{key} + X-Context | 204 | 503 < W |
- values[] has one entry, or several siblings after concurrent writes.
- The context is the opaque, joined version vector of what the client read.
- A retried PUT writes the same value again. The client then merges two equal siblings.
- r and w are per request: r=1 for a fast read, w=3 for a write that must be on every replica.
A read returns a context. The client sends it back on the next write, so the store knows which versions that write replaces.
- key
- Hashed onto the ring; never ordered.
- sibling
- A version vector plus a value or a tombstone.
- engine
- LSM: 9.13× write amplification for random inserts in the D1 lab, against 25.82× for a B-tree.
CREATE TABLE IF NOT EXISTS kv (
key text PRIMARY KEY1,
value bytea NOT NULL CHECK (octet_length(value) <= 1048576),
version bigint2 NOT NULL DEFAULT 1,
updated_at timestamptz NOT NULL DEFAULT now()
);
UPDATE kv SET value = $2, version = version + 1, updated_at = now()
WHERE key = $1 AND version = $33
RETURNING version;- 1A point lookup is one index probe. This is the store until one machine is too small.
- 2One counter per key. With one node there is no concurrent writer elsewhere, so a counter is enough.
- 3Compare-and-set: the write fails if another client wrote since this client read. The lab ran 16 such writers: 1 won.
Each key holds one or more siblings, each a version vector and a value. A node stores them in an LSM tree, because writes only append.
Step 1: Write
- The client sends PUT to any node, or straight to a home replica.
- The coordinator stamps a version vector and sends the write to the 3 home replicas.
- It answers after W = 2 acks. The third copy arrives later.
If it fails
A home replica is down: the next healthy node on the ring stores the write with a hint. Fewer than W acks: the client gets 503 and retries with the same context.
Any node can coordinate. It writes to the key's 3 home replicas and answers after 2. Handoff, read repair and Merkle trees bring the third copy back in line.
| tool | capability | what it gives this design | also used for |
|---|---|---|---|
| Store node | Consistent hash ring, many tokens per node | Even spread; a join moves about 1/(n+1) of the copies. | Cache clusters, load balancers |
| Store node | Preference list | N copies on N distinct nodes; skip a node already taken. | Zone-aware placement |
| Store node | Quorums: N, R, W per request | Trade freshness against latency and availability per call. | Leaderless databases |
| Store node | Sloppy quorum and hinted handoff | Writes succeed while home replicas are down. | Retry queues |
| Store node | Read repair | A read fixes the stale replicas it saw. | Caches with versioned entries |
| Store node | Merkle trees | Find differing key ranges with few hashes, not a full scan. | Git, file sync, blockchains |
| Store node | Version vectors | Tell a later write from a concurrent one; keep siblings. | CRDTs, sync engines |
| Store node | Gossip and a failure detector | Every node knows the ring and who is down, with no central server. | Service discovery, clusters |
| LSM engine | Append-only writes, sorted files, Bloom filters | Fast writes; a miss usually costs no disk read. | Time series, logs, queues |
| Postgres | Primary key and UPDATE ... WHERE version | Step 1: a one-node store with compare-and-set. | Optimistic locking |
The ring places data, quorums tune consistency, and three repair paths make the replicas converge.
| method | status | a seventh node moves |
|---|---|---|
| hash(key) mod N | Not approved | 85.4% of keys |
| Ring, 1 token per node | Uneven | 18.6% of copies |
| Ring, 16 tokens per node | Approved | 10.2% of copies |
| Ring, 64 tokens per node | Approved | 14.4% of copies |
| Fixed slots and a slot map | Approved | whole slots, by plan |
| Key ranges | Range scans needed | split or merge ranges |
Lab: 6 nodes, 6,000 keys, N = 3; the ideal share is 1/7 = 14.3%.
preference(key): // every node computes the same list1
h = hash(key) // a point on a 64-bit ring2
walk the tokens clockwise from h:
take the token's node unless already taken3
first N nodes = home replicas // N = 3
next nodes = stand-ins, in order4
join(node):
place 16 to 256 tokens for the node // each token splits one range
FOR EACH key whose home list changed:
stream the value from an old replica5
the old replica drops its copy- 1The ring is small and every node holds it, so any node can route any key. No directory lookup.
- 2The lab reuses the sharding ring hash: FNV-1a, then a mixing step.
- 3A node owns many tokens, so the walk skips tokens of nodes it has. A real store also skips a second node in the same zone.
- 4Where a write goes when a home replica is down.
- 5Only keys whose home list changed move: about 1/(n+1) of all copies.
Tested source Go: preference list · Go: join
// Preference walks the ring clockwise from the key's hash and lists each node once, in the order
// it meets them. The first N are the key's home replicas; the rest are stand-ins, in order.
func (c *Cluster) Preference(key string) []string {
h := sharding.Hash(key)
start := sort.Search(len(c.tokens), func(i int) bool { return c.tokens[i].Pos >= h })
var out []string
for i := range c.tokens {
t := c.tokens[(start+i)%len(c.tokens)] // wrap past the last token
if !slices.Contains(out, t.Node) {
out = append(out, t.Node)
}
if len(out) == len(c.nodes) {
break
}
}
return out
}
// Join adds a node to the ring. Its tokens take over parts of other nodes' ranges. For every
// key whose home replicas changed, the new home copies the value from the old replicas (merged),
// and a node that is no longer a home replica drops its copy. It returns how many key copies
// moved and how many copies the cluster holds.
func (c *Cluster) Join(name string) (moved, copies int, err error) {
if _, ok := c.nodes[name]; ok {
return 0, 0, fmt.Errorf("join %s: already a member", name)
}
keys := c.Keys()
before := map[string][]string{}
for _, k := range keys {
before[k] = c.Home(k)
}
c.addTokens(name)
for _, k := range keys {
merged := Value{}
for _, h := range before[k] {
merged = merged.Merge(c.nodes[h].Data[k])
}
after := c.Home(k)
for _, h := range after {
if !slices.Contains(before[k], h) {
c.nodes[h].Data[k] = merged // streamed from the old replicas
moved++
}
}
for _, h := range before[k] {
if !slices.Contains(after, h) {
delete(c.nodes[h].Data, k)
}
}
copies += len(after)
}
return moved, copies, nil
}
I hash each key onto a ring where every node owns many tokens. The first 3 distinct nodes clockwise are the key's replicas, and a new node takes small slices from everyone.
| N = 3 | writes ok | reads ok | stale reads | status |
|---|---|---|---|---|
| R = 1, W = 1 | 99.9% | 99.9% | 9.5% | Stale reads are fine |
| R = 1, W = 3 | 73.5% | 99.9% | 0.0% | Read-heavy, rare writes |
| R = 3, W = 1 | 99.9% | 73.5% | 0.0% | Not approved |
| R = 2, W = 2, strict | 97.2% | 97.2% | 0.0% | Approved |
| R = 2, W = 2, sloppy | 100.0% | 97.2% | 0.1% | Approved: never down |
Lab: each node down 10% of the time, at random; 20,000 reads and writes per row, one client that reads, then writes.
- R + W > N makes every read set meet every write set.
- W = 3 loses a quarter of the writes when each node is down 10% of the time.
- The sloppy quorum takes every write. It pays with a few stale reads until handoff.
put(key, value, context):
coordinator = first healthy home replica1
version = context + 1 for the coordinator
FOR EACH home replica r:
IF r is up: store on r
ELSE: store on the next healthy stand-in, with a hint for r2
IF acks >= W: RETURN ok3
RETURN error // copies already stored stay4
get(key):
ask every healthy home replica
RETURN merge of the first R replies5
then write the merge of ALL replies to each replica behind it
handoff(): // every node, every few seconds
FOR EACH hint whose home replica is up:
deliver it, then delete it6- 1A coordinator that is also a replica saves one network hop.
- 2Sloppy quorum: the stand-in keeps the write apart from its own data, so reads never see it there.
- 3The client waits for W acks only. The slowest replica does not add latency.
- 4There is no rollback. A failed write can still appear later; the client retries with the same context.
- 5R replies are enough to answer. The rest arrive and feed read repair.
- 6Delivery is a merge, so a hint delivered twice does no harm.
Tested source Go: put · Go: get with read repair · Go: handoff
// Put writes val under key. ctx is the value the client read before (empty for a blind write):
// the new version supersedes everything in it. The coordinator is the first healthy home
// replica. It sends the write to each home replica; with a sloppy quorum, a down replica's copy
// goes to the next healthy node outside the N, with a hint. The write succeeds at W acks.
func (c *Cluster) Put(key, val string, ctx Value) (PutResult, error) {
coord := c.firstUp(c.Preference(key))
if coord == nil {
return PutResult{}, fmt.Errorf("put %s: %w: every node is down", key, ErrNoQuorum)
}
return c.PutAt(coord.Name, key, val, ctx)
}
// PutAt is Put with a chosen coordinator, as when a load balancer sends the request to any node.
func (c *Cluster) PutAt(coordinator, key, val string, ctx Value) (PutResult, error) {
pref := c.Preference(key)
coord, err := c.Node(coordinator)
if err != nil {
return PutResult{}, fmt.Errorf("put %s: %w", key, err)
}
if !coord.Up {
return PutResult{}, fmt.Errorf("put %s: coordinator %s is down: %w", key, coordinator, ErrNoQuorum)
}
vv := crdt.VV{}
for _, e := range ctx.Entries {
vv = vv.Join(e.VV)
}
for _, e := range coord.Data[key].Entries { // never reuse a counter the coordinator has issued
vv[coord.Name] = max(vv[coord.Name], e.VV[coord.Name])
}
vv[coord.Name]++
w := Value{Entries: []crdt.MVEntry[string]{{Value: val, VV: vv}}}
res := PutResult{Coordinator: coord.Name, Version: vv}
used := map[string]bool{}
for _, home := range pref[:c.cfg.N] {
if n := c.nodes[home]; n.Up {
n.Data[key] = n.Data[key].Merge(w)
res.Acks = append(res.Acks, Ack{Node: home})
continue
}
if !c.cfg.Sloppy {
continue
}
for _, s := range pref[c.cfg.N:] { // the next healthy node not used yet
if n := c.nodes[s]; n.Up && !used[s] {
used[s] = true
n.Hints = append(n.Hints, Hint{For: home, Key: key, Val: w})
res.Acks = append(res.Acks, Ack{Node: s, For: home})
break
}
}
}
if len(res.Acks) < c.cfg.W {
return res, fmt.Errorf("put %s: %w: %d of W=%d acks", key, ErrNoQuorum, len(res.Acks), c.cfg.W)
}
return res, nil
}
// Get reads key from the healthy home replicas. Replies arrive in preference order; the client
// gets the merge of the first R. The coordinator then merges every reply and writes the merge
// back to each replica that was behind (read repair).
func (c *Cluster) Get(key string) (GetResult, error) {
var res GetResult
for _, home := range c.Home(key) {
if n := c.nodes[home]; n.Up {
res.Replies = append(res.Replies, Reply{Node: home, Val: n.Data[key]})
}
}
if len(res.Replies) < c.cfg.R {
return res, fmt.Errorf("get %s: %w: %d of R=%d replies", key, ErrNoQuorum, len(res.Replies), c.cfg.R)
}
all := Value{}
for i, r := range res.Replies {
if i < c.cfg.R {
res.Val = res.Val.Merge(r.Val)
}
all = all.Merge(r.Val)
}
for _, r := range res.Replies {
if !r.Val.Equal(all) {
c.nodes[r.Node].Data[key] = all
res.Repaired = append(res.Repaired, r.Node)
}
}
return res, nil
}
// Handoff runs on every healthy node: each hint whose home replica is up again is delivered,
// merged into that replica, and deleted. Hints for replicas still down stay.
func (c *Cluster) Handoff() []Ack {
var done []Ack
for _, name := range c.Names() {
n := c.nodes[name]
if !n.Up {
continue
}
var keep []Hint
for _, h := range n.Hints {
home := c.nodes[h.For]
if !home.Up {
keep = append(keep, h)
continue
}
home.Data[h.Key] = home.Data[h.Key].Merge(h.Val)
done = append(done, Ack{Node: name, For: h.For})
}
n.Hints = keep
}
return done
}
I default to N = 3, R = 2, W = 2 with a sloppy quorum. Writes stay available, and reads are fresh except during a failure before handoff runs.
| method | status | cost |
|---|---|---|
| Last write wins, by timestamp | A lost update is fine | Clock skew drops writes without a trace. |
| Version vectors, siblings | Approved | The client merges; one counter per coordinator. |
| CRDT values (counter, set) | Type has a merge | Metadata grows with writers and deletes. |
| One leader per key range | Failover pause is fine | No conflicts, but writes stop during an election. |
- The lab coordinator never reuses a counter it has issued, so two writes through one node never share a vector.
- A shopping cart merges by union. A deleted item may then return; the cart needs an OR-Set to stop that.
write at coordinator c with context ctx:
vv = join of every vector in ctx1
vv[c] = max(vv[c], highest c counter seen here2) + 1
store on a replica:
keep each sibling that no other sibling's vector covers3
read:
RETURN every sibling, plus the joined vector as the context
client:
merge the siblings (union of carts), write with that context
// the new vector covers both4- 1The new write covers everything the client read.
- 2Two clients that write the same context through one coordinator still get different vectors.
- 3A covered sibling is older and goes. Two uncovered siblings are concurrent and both stay.
- 4After the merge write, one value is left on every replica.
Tested source Go: put (stamps the vector) · Go: register merge, reused from the CRDT lab
// Put writes val under key. ctx is the value the client read before (empty for a blind write):
// the new version supersedes everything in it. The coordinator is the first healthy home
// replica. It sends the write to each home replica; with a sloppy quorum, a down replica's copy
// goes to the next healthy node outside the N, with a hint. The write succeeds at W acks.
func (c *Cluster) Put(key, val string, ctx Value) (PutResult, error) {
coord := c.firstUp(c.Preference(key))
if coord == nil {
return PutResult{}, fmt.Errorf("put %s: %w: every node is down", key, ErrNoQuorum)
}
return c.PutAt(coord.Name, key, val, ctx)
}
// PutAt is Put with a chosen coordinator, as when a load balancer sends the request to any node.
func (c *Cluster) PutAt(coordinator, key, val string, ctx Value) (PutResult, error) {
pref := c.Preference(key)
coord, err := c.Node(coordinator)
if err != nil {
return PutResult{}, fmt.Errorf("put %s: %w", key, err)
}
if !coord.Up {
return PutResult{}, fmt.Errorf("put %s: coordinator %s is down: %w", key, coordinator, ErrNoQuorum)
}
vv := crdt.VV{}
for _, e := range ctx.Entries {
vv = vv.Join(e.VV)
}
for _, e := range coord.Data[key].Entries { // never reuse a counter the coordinator has issued
vv[coord.Name] = max(vv[coord.Name], e.VV[coord.Name])
}
vv[coord.Name]++
w := Value{Entries: []crdt.MVEntry[string]{{Value: val, VV: vv}}}
res := PutResult{Coordinator: coord.Name, Version: vv}
used := map[string]bool{}
for _, home := range pref[:c.cfg.N] {
if n := c.nodes[home]; n.Up {
n.Data[key] = n.Data[key].Merge(w)
res.Acks = append(res.Acks, Ack{Node: home})
continue
}
if !c.cfg.Sloppy {
continue
}
for _, s := range pref[c.cfg.N:] { // the next healthy node not used yet
if n := c.nodes[s]; n.Up && !used[s] {
used[s] = true
n.Hints = append(n.Hints, Hint{For: home, Key: key, Val: w})
res.Acks = append(res.Acks, Ack{Node: s, For: home})
break
}
}
}
if len(res.Acks) < c.cfg.W {
return res, fmt.Errorf("put %s: %w: %d of W=%d acks", key, ErrNoQuorum, len(res.Acks), c.cfg.W)
}
return res, nil
}
// Set writes v at replica r with timestamp ts from r's clock.
func (l LWW[T]) Set(r string, v T, ts int64) (LWW[T], LWW[T]) {
w := LWW[T]{Value: v, TS: ts, Replica: r}
return l.Merge(w), w
}
// Merge keeps the later write. Equal timestamps fall back to the replica name, so every
// replica picks the same winner. The loser is gone, even if it happened later in real time.
func (l LWW[T]) Merge(o LWW[T]) LWW[T] {
if o.TS > l.TS || (o.TS == l.TS && o.Replica > l.Replica) {
return o
}
return l
}
// MV is a multi-value register. Each write carries a version vector. A write that has seen
// another write replaces it; two concurrent writes are both kept, as siblings, until a later
// write that has seen both replaces them.
type MV[T comparable] struct {
Entries []MVEntry[T] `json:"entries"`
}
// MVEntry is one sibling: a value and the version vector of its write.
type MVEntry[T comparable] struct {
Value T `json:"value"`
VV VV `json:"vv"`
}
// Set writes v at replica r. Its vector is the join of every sibling's vector plus one for r,
// so it supersedes every sibling this replica has seen.
func (m MV[T]) Set(r string, v T) (MV[T], MV[T]) {
vv := VV{}
for _, e := range m.Entries {
vv = vv.Join(e.VV)
}
vv[r]++
w := MV[T]{Entries: []MVEntry[T]{{Value: v, VV: vv}}}
return m.Merge(w), w
}
// Merge keeps every entry that no other entry's vector dominates.
func (m MV[T]) Merge(o MV[T]) MV[T] {
all := append(slices.Clone(m.Entries), o.Entries...)
var out []MVEntry[T]
for i, e := range all {
keep := true
for j, f := range all {
if i != j && (e.VV.Before(f.VV) || (e.VV.Equal(f.VV) && j < i)) {
keep = false // dominated, or a duplicate of an earlier entry
break
}
}
if keep {
out = append(out, e)
}
}
sortEntries(out)
return MV[T]{Entries: out}
}
Each write carries a version vector. When neither vector covers the other, the store keeps both values as siblings and the client merges them.
| method | status | cost for 1 difference |
|---|---|---|
| Send every key and version | Not approved | 11,621 entries |
| One hash for the whole range | Detects only | 1 hash, then a full scan |
| Merkle tree, depth 10 | Approved | 21 hashes, 41 keys |
| Merkle tree, depth 15 | Approved | 31 hashes, 1 key |
The figure: depth 3, 48 shared keys; node B missed one write while it was down. Tables: 11,621 shared keys.
build(node, shared keys, depth):
cut the hash space into 2^depth leaves1
leaf = hash(every key and version in its range, sorted)
inner = hash(left child, right child)
compare(a, b, level, i):
IF a[level][i] == b[level][i]: RETURN // the whole range matches2
IF level == depth: sync the keys of leaf i3; RETURN
compare(a, b, level + 1, 2i)
compare(a, b, level + 1, 2i + 1)
sync(key): both sides store merge(a[key], b[key])4- 1Both replicas cut the same ranges, so leaf i means the same keys on both.
- 2One comparison clears a whole subtree. Equal replicas cost 1 hash.
- 3A deeper tree means fewer keys per leaf to send, but more hashes to keep up to date.
- 4A merge, not a copy: siblings on either side survive.
Tested source Go: build · Go: compare · Go: sync
// BuildMerkle hashes the given keys of one node into a tree of the given depth.
func BuildMerkle(n *Node, keys []string, depth int) *Merkle {
m := &Merkle{Depth: depth, keys: make([][]string, 1<<depth)}
for _, k := range keys {
l := Leaf(k, depth)
m.keys[l] = append(m.keys[l], k)
}
leaves := make([][32]byte, 1<<depth)
for i, ks := range m.keys {
slices.Sort(ks)
h := sha256.New()
for _, k := range ks {
if v, ok := n.Data[k]; ok { // a key this node lacks adds nothing
fmt.Fprintf(h, "%s:%s\n", k, canonical(v))
}
}
copy(leaves[i][:], h.Sum(nil))
}
m.Levels = make([][][32]byte, depth+1)
m.Levels[depth] = leaves
for lv := depth - 1; lv >= 0; lv-- {
below := m.Levels[lv+1]
m.Levels[lv] = make([][32]byte, len(below)/2)
for i := range m.Levels[lv] {
m.Levels[lv][i] = sha256.Sum256(append(below[2*i][:], below[2*i+1][:]...))
}
}
return m
}
// Compare walks both trees from the root. It descends only into children whose hashes differ,
// so matching subtrees cost one comparison each.
func Compare(a, b *Merkle) Diff {
var d Diff
var walk func(lv, i int)
walk = func(lv, i int) {
d.Compared++
d.Path = append(d.Path, [2]int{lv, i})
if a.Levels[lv][i] == b.Levels[lv][i] {
return // the whole range matches
}
d.Differ = append(d.Differ, [2]int{lv, i})
if lv == a.Depth {
d.Leaves = append(d.Leaves, i)
return
}
walk(lv+1, 2*i)
walk(lv+1, 2*i+1)
}
walk(0, 0)
return d
}
// AntiEntropy compares two replicas with Merkle trees and repairs only the differing leaves:
// for each key in them, both sides store the merge of the two values.
func (c *Cluster) AntiEntropy(a, b string, depth int) (SyncResult, error) {
na, err := c.Node(a)
if err != nil {
return SyncResult{}, err
}
nb, err := c.Node(b)
if err != nil {
return SyncResult{}, err
}
shared := c.SharedKeys(a, b)
ta, tb := BuildMerkle(na, shared, depth), BuildMerkle(nb, shared, depth)
res := SyncResult{Diff: Compare(ta, tb), Shared: len(shared), LeafCount: 1 << depth}
for _, l := range res.Diff.Leaves {
for _, k := range ta.KeysIn(l) {
res.Keys++
va, vb := na.Data[k], nb.Data[k]
if va.Equal(vb) {
continue
}
m := va.Merge(vb)
na.Data[k], nb.Data[k] = m, m
res.Repaired = append(res.Repaired, k)
}
}
return res, nil
}
Replicas compare Merkle trees root first. One differing key among 11,621 costs 21 hash comparisons at depth 10, not a scan of every key.
All six nodes are up. A write waits for 2 of 3 acks, and a read for 2 of 3 replies.
key cart:42 · N = 3 · R = 2 · W = 2 · strict quorum, no hints
✓ home replicas agree0 hints held
A node failure costs hints, not writes. The only stale read in these runs comes from the sloppy quorum before handoff, and read repair fixes a replica on the next read.
| event | result | why it is safe | saved by |
|---|---|---|---|
| A home replica is down during a write | A stand-in stores the copy with a hint. | W acks still arrive. The hint goes home when the replica returns. | Hinted handoff |
| The stand-in dies before handoff | One replica misses the write. | Read repair or the next Merkle comparison copies the key. | Merkle trees |
| A read reaches a stale replica | With R = 1, the old value can return. | R + W > N on a strict quorum. Read repair then fixes the replica. | Read repair |
| A network partition; both sides take writes | Two concurrent versions of one key. | Version vectors keep both as siblings. The client merges them. | Version vectors |
| Clocks disagree under last write wins | A later write can lose. | Use version vectors, or accept LWW only where a lost update is fine. | Version vectors |
| A node loses its disk | Its ranges have 2 copies left. | A new node takes its tokens and streams the ranges from the other replicas. | Streaming |
| A replica misses a delete | The old value can come back. | Keep tombstones longer than a full repair cycle. Cassandra defaults to 10 days. | Tombstones |
| One key gets most of the traffic | Its 3 replicas run hot. | Cache it in front of the store, or split the key into several. | Cache |
| step | add | it handles | move up when you see |
|---|---|---|---|
| 1 | One Postgres table, key and value, with a version column. | About 30,600 writes and 69,750 reads a second in the lab; data up to one machine's disk. | Reads load the primary, or one node down is too long an outage. |
| 2 | Replicas and failover. | Reads grow with replicas. A failover takes seconds; writes pause during it. | Data or writes outgrow one primary. |
| 3 | Hash shards, each a primary with replicas. | Writes and data grow with shards. Each shard still pauses for its own failover. | Writes must succeed through node and zone failures, with no pause. |
| 4 | A leaderless ring, N = 3, quorums, hints, Merkle repair. | Add nodes for data or throughput: 18 nodes give about 183,600 writes a second at the one-node rate. | Users in other regions, or a whole region must fail without an outage. |
| 5 | A ring per region, quorums inside a region, replication between regions. | Local latency in each region. Cross-region conflicts become siblings. | Top of the ladder. |
Ring capacity is derived: the lab's one-node rate × 18 nodes ÷ 3 copies for writes, ÷ R = 2 for reads. An LSM engine on server disks has its own rate; use these as orders of magnitude.
I start with one Postgres table. I go leaderless only when the data outgrows a sharded primary setup, or when writes must continue through a failover.
| question | answer | go deeper |
|---|---|---|
| How does a node know who is in the ring? | Gossip spreads membership and heartbeats in a number of rounds that grows with log N. A failure detector turns heartbeat gaps into "down". | X6 |
| Can the store be strongly consistent? | Quorums alone are not linearizable: a write that fails can still appear later. For that, run a consensus group per key range, as Raft-based stores do. | X5, X3 |
| How would you support range scans? | Partition by key range instead of hash, and split hot ranges. Scans then read one or a few nodes in order. | X2 |
| Why an LSM tree on each node? | Writes only append and flush sorted files, so random writes stay cheap. Bloom filters keep most misses off the disk. | D1 |
| Last write wins or vector clocks? | LWW is simple but drops one of two concurrent writes. Version vectors keep both and push the merge to the client or to a CRDT. | X7, X4 |
| How do you replicate across regions? | Use a quorum inside each region and send writes to the other regions without waiting. A region can fail without a write outage. | X1 |
0 of 9 known
N = 3, W = 2, R = 2. Why does a read see the last acknowledged write?
Two of the three home replicas are down. Does a write with W = 2 succeed?
Why can a read miss a write that a sloppy quorum acknowledged?
A hint is lost because the stand-in died. What repairs the replica?
Two replicas share 11,621 keys and differ in one. How many hashes does a depth-10 tree comparison read?
Why does the store keep two siblings instead of the newer value?
Why use 16 to 256 tokens per node instead of 1?
A seventh node joins. How much data moves?
A deleted key comes back a week later. What went wrong?
- quorum
- R = W = 2 of 3: 0 stale reads. R = W = 1: 9.5% stale, each node down 10% of the time.
- writes ok
- W = 3: 73.5%. W = 2: 97.2%. Sloppy W = 2: 100.0%.
- join
- A 7th node moves 14.4% of copies with 64 tokens; mod N moves 85.4%.
- Merkle
- 1 difference in 11,621 keys: 21 hashes at depth 10.
- gossip
- Fanout 3 reaches 64 nodes in 5.2 rounds, 4,096 in 9.6.
- one node
- About 30,600 writes and 69,750 reads a second, 1 KB values.
- LSM
- 9.13× write amplification, 1.08× space, leveled.
Simulations are seeded and repeat exactly. One-node rates: Postgres 16 on an 8-core laptop, 32 clients, shared with other jobs. Tombstone default from the Cassandra documentation.