System Design
C3

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.

Not startedSaved in this browser only.
  1. 1A hash ring with many tokens per node places N = 3 copies and moves about 1/(n+1) of the data on a join.
  2. 2R + W > N gives fresh reads; a sloppy quorum with hints keeps writes available when replicas are down.
  3. 3Version vectors keep concurrent writes as siblings, so no write is lost.
  4. 4Three repair paths: hinted handoff, read repair, and Merkle-tree anti-entropy.
C3
    A

    The prompt

    as the interviewer says it
    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.

    B

    Clarifying questions

    ask, then assume
    questionanswer assumedit decides
    How many keys, how large?10 billion keys, 1 KB valuesNumber of nodes
    Reads and writes a day?2 billion reads, 0.5 billion writesNodes for throughput
    Can a read be a little stale?Yes, unless the client asks for R + W > NQuorum sizes
    Must writes succeed during a failure?Yes: never down for writesLeaderless, sloppy quorum
    Range scans or only point lookups?Point lookups onlyHash partitioning
    Two writers on one key at once?Rare, but no write may be lostVersion vectors, siblings
    One region or many?One region, 3 zonesReplica placement

    I ask about size, read-write mix, consistency and range scans first, because each answer changes the partitioning or the replication.

    C

    Estimates

    arithmetic shown
    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.

    D

    API

    one resource: a key
    callsuccesserrors
    GET /v1/kv/{key}?r=2200 values[] + context404, 503 < R
    PUT /v1/kv/{key}?w=2 + X-Context204413, 503 < W
    DELETE /v1/kv/{key} + X-Context204503 < 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.

    E

    Data model

    per node, then the one-node version
    one key on one nodekeycart:42sibling 1 · vector{B:1,E:1}milk, breadsibling 2 · vector{B:2}milk, eggsa delete is a sibling with no value: a tombstoneStorage node: an LSM treecommit logappend, fsyncmemtablesorted, in memoryfull: flush onceas a sorted fileL0L1L2SSTables; the dark end is each file's Bloom filter
    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.
    step 1: the whole store in one Postgres tablesql
    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;
    1. 1A point lookup is one index probe. This is the store until one machine is too small.
    2. 2One counter per key. With one node there is no concurrent writer elsewhere, so a counter is enough.
    3. 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.

    F

    The architecture

    click a step; its path lights up
    Clientapp serverLoad balanceror client libraryCoordinatorany node, ring in RAMHome replica 1LSM engineHome replica 2LSM engineHome replica 3down in step 3Stand-innext node on ringPeer nodegossip, trees

    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.

    G

    Capabilities used

    what each part gives you
    toolcapabilitywhat it gives this designalso used for
    Store nodeConsistent hash ring, many tokens per nodeEven spread; a join moves about 1/(n+1) of the copies.Cache clusters, load balancers
    Store nodePreference listN copies on N distinct nodes; skip a node already taken.Zone-aware placement
    Store nodeQuorums: N, R, W per requestTrade freshness against latency and availability per call.Leaderless databases
    Store nodeSloppy quorum and hinted handoffWrites succeed while home replicas are down.Retry queues
    Store nodeRead repairA read fixes the stale replicas it saw.Caches with versioned entries
    Store nodeMerkle treesFind differing key ranges with few hashes, not a full scan.Git, file sync, blockchains
    Store nodeVersion vectorsTell a later write from a concurrent one; keep siblings.CRDTs, sync engines
    Store nodeGossip and a failure detectorEvery node knows the ring and who is down, with no central server.Service discovery, clusters
    LSM engineAppend-only writes, sorted files, Bloom filtersFast writes; a miss usually costs no disk read.Time series, logs, queues
    PostgresPrimary key and UPDATE ... WHERE versionStep 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.

    H

    Deep dive: partitioning

    where does a key live?
    ringcart:42 hashes to 0.7874 of the ring.6 nodes × 16 tokens = 96 points on the ring.The coloured arc is the stretch the walk covers.positioncart:420.78740.7957Bhome10.8123Ehome20.8157Chome30.8231Astand-in10.8388Cskip0.8417Askip0.8450Fstand-in20.8501Bskip0.8591Dstand-in3clockwise; each node once
    methodstatusa seventh node moves
    hash(key) mod NNot approved85.4% of keys
    Ring, 1 token per nodeUneven18.6% of copies
    Ring, 16 tokens per nodeApproved10.2% of copies
    Ring, 64 tokens per nodeApproved14.4% of copies
    Fixed slots and a slot mapApprovedwhole slots, by plan
    Key rangesRange scans neededsplit or merge ranges

    Lab: 6 nodes, 6,000 keys, N = 3; the ideal share is 1/7 = 14.3%.

    preference list and joinpseudo code
    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
    1. 1The ring is small and every node holds it, so any node can route any key. No directory lookup.
    2. 2The lab reuses the sharding ring hash: FNV-1a, then a mixing step.
    3. 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.
    4. 4Where a write goes when a home replica is down.
    5. 5Only keys whose home list changed move: about 1/(n+1) of all copies.
    Tested source Go: preference list · Go: join
    Go: preference listgo
    // 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
    }
    
    Go: joingo
    // 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.

    I

    Deep dive: replication and quorums

    N, R and W
    N = 3writes okreads okstale readsstatus
    R = 1, W = 199.9%99.9%9.5%Stale reads are fine
    R = 1, W = 373.5%99.9%0.0%Read-heavy, rare writes
    R = 3, W = 199.9%73.5%0.0%Not approved
    R = 2, W = 2, strict97.2%97.2%0.0%Approved
    R = 2, W = 2, sloppy100.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.
    write, read, handoffpseudo code
    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
    1. 1A coordinator that is also a replica saves one network hop.
    2. 2Sloppy quorum: the stand-in keeps the write apart from its own data, so reads never see it there.
    3. 3The client waits for W acks only. The slowest replica does not add latency.
    4. 4There is no rollback. A failed write can still appear later; the client retries with the same context.
    5. 5R replies are enough to answer. The rest arrive and feed read repair.
    6. 6Delivery is a merge, so a hint delivered twice does no harm.
    Tested source Go: put · Go: get with read repair · Go: handoff
    Go: putgo
    // 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
    }
    
    Go: get with read repairgo
    // 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
    }
    
    Go: handoffgo
    // 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.

    J

    Deep dive: versions and conflicts

    which write wins?
    methodstatuscost
    Last write wins, by timestampA lost update is fineClock skew drops writes without a trace.
    Version vectors, siblingsApprovedThe client merges; one counter per coordinator.
    CRDT values (counter, set)Type has a mergeMetadata grows with writers and deletes.
    One leader per key rangeFailover pause is fineNo 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.
    version vectorspseudo code
    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
    1. 1The new write covers everything the client read.
    2. 2Two clients that write the same context through one coordinator still get different vectors.
    3. 3A covered sibling is older and goes. Two uncovered siblings are concurrent and both stay.
    4. 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
    Go: put (stamps the vector)go
    // 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
    }
    
    Go: register merge, reused from the CRDT labgo
    
    // 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.

    K

    Deep dive: anti-entropy

    recorded from the lab
    A 07f9B 4909A d9fdB d9fdA 47e0B 19d4A 4b7fB 4b7fA 82a0B 82a0A 164eB 164eA d52bB 35a9A 7027B 7027A 0099B 0099A e3b0B e3b0A 3d4fB 3d4fA d333B d333A e3b0B e3b0A 8772B 5be0A d127B d12713 keys1 key0 keys9 keys3 keys0 keys8 keys14 keysequal: stopdiffer: open both childrennever compared
    methodstatuscost for 1 difference
    Send every key and versionNot approved11,621 entries
    One hash for the whole rangeDetects only1 hash, then a full scan
    Merkle tree, depth 10Approved21 hashes, 41 keys
    Merkle tree, depth 15Approved31 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, compare, syncpseudo code
    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
    1. 1Both replicas cut the same ranges, so leaf i means the same keys on both.
    2. 2One comparison clears a whole subtree. Equal replicas cost 1 hash.
    3. 3A deeper tree means fewer keys per leaf to send, but more hashes to keep up to date.
    4. 4A merge, not a copy: siblings on either side survive.
    Tested source Go: build · Go: compare · Go: sync
    Go: buildgo
    // 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
    }
    
    Go: comparego
    // 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
    }
    
    Go: syncgo
    // 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.

    L

    Try it: five runs on one key

    recorded from the lab cluster

    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 1Bv1 {B:1}
    home 2Ev1 {B:1}
    home 3Cv1 {B:1}
    stand-in 1Ano copy
    stand-in 2Fno copy
    stand-in 3Dno copy
    Write "v1" through coordinator B. Acks from B, E and C. 3 acks, W = 2: the write succeeds.

    ✓ 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.

    M

    Failure cases

    what breaks, and why it stays correct
    eventresultwhy it is safesaved by
    A home replica is down during a writeA 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 handoffOne replica misses the write.Read repair or the next Merkle comparison copies the key.Merkle trees
    A read reaches a stale replicaWith 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 writesTwo concurrent versions of one key.Version vectors keep both as siblings. The client merges them.Version vectors
    Clocks disagree under last write winsA later write can lose.Use version vectors, or accept LWW only where a lost update is fine.Version vectors
    A node loses its diskIts ranges have 2 copies left.A new node takes its tokens and streams the ranges from the other replicas.Streaming
    A replica misses a deleteThe 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 trafficIts 3 replicas run hot.Cache it in front of the store, or split the key into several.Cache
    N

    Scale ladder

    start simple; climb only on a signal
    Each step adds one component1One table2+ replicas3+ shards4Leaderless ring5+ regionsmore load →
    Capacity against demand1k10k100k1MDemand, average: 28,935 operations per secondDemand, average28,935Demand, 3× peak: 86,805 operations per secondDemand, 3× peak86,805One node, writes: 30,600 operations per secondOne node, writes30,600One node, reads: 69,750 operations per secondOne node, reads69,75018 nodes, writes: 183,600 operations per second18 nodes, writes183,60018 nodes, reads: 627,750 operations per second18 nodes, reads627,750operations per second, log scale
    stepaddit handlesmove up when you see
    1One 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.
    2Replicas and failover.Reads grow with replicas. A failover takes seconds; writes pause during it.Data or writes outgrow one primary.
    3Hash 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.
    4A 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.
    5A 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.

    O

    Follow-ups

    what the interviewer asks next
    questionanswergo 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
    P

    Drill

    predict, then reveal

    0 of 9 known

    1. N = 3, W = 2, R = 2. Why does a read see the last acknowledged write?

    2. Two of the three home replicas are down. Does a write with W = 2 succeed?

    3. Why can a read miss a write that a sloppy quorum acknowledged?

    4. A hint is lost because the stand-in died. What repairs the replica?

    5. Two replicas share 11,621 keys and differ in one. How many hashes does a depth-10 tree comparison read?

    6. Why does the store keep two siblings instead of the newer value?

    7. Why use 16 to 256 tokens per node instead of 1?

    8. A seventh node joins. How much data moves?

    9. A deleted key comes back a week later. What went wrong?

    Q

    Numbers to say

    measured, derived or cited
    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.