Day 7 B-Building blocks 2026-09-30 ← All lessons

Design a key-value store: partitioning, replication, consistency

A single hash table runs out of room fast. A distributed key value store spreads pairs across servers and tunes consistency against latency.

10Max size of one key value pair
3Copies stored per key
2Quorum for strong consistency

Key points

Outcomes

01Explain the CAP tradeoff with the three replica example and choose CP or AP for a use case.
02Size N, W, and R quorums to trade read speed, write speed, and strong consistency.
03Describe failure handling: gossip detection, hinted handoff, and Merkle tree anti entropy.
01

CAP: pick two, design for it


A key value store exposes put and get on pairs under 10 KB. One server with a hash table runs out of room fast, even with compression and disk overflow. Scale means many servers, and many servers mean partitions. CAP says you cannot hold consistency, availability, and partition tolerance all at once.

Option A

CP: consistency first

  • Block writes during a partition to avoid divergence
  • Bank balances stay correct: errors beat stale reads
Option B

AP: availability first

  • Keep accepting reads and writes, sync when healed
  • Feeds and carts tolerate briefly stale data
Interview tipInterview line: state your CAP choice early, then set N, W, and R to match it.
02

Partition and replicate


Consistent hashing from Chapter 5 places servers on a ring. Each key walks clockwise to its owner, then copies to the first N distinct servers. With N equal to 3, key0 lives on s1, s2, and s3. Replicas sit in distinct data centers so one outage cannot take all copies.

How to read: Watch the dots: the key lands on the ring, then three dots fan out clockwise to owner and replicas.

key0Hash ringSHA spaces1 owners2 replicas3 replica
  1. Hash key0 onto the same ring the servers share.

  2. Walk clockwise: the first server met owns the key.

  3. Copy onward to the next N distinct servers.

Replication journey: key0 hashes onto the ring, then copies clockwise to three servers.
Interview tipMore virtual nodes per server evens out slices, and fatter servers earn more of them. Adding or removing a server moves only neighboring keys.
03

Quorums: tune N, W, and R


Every write waits for W acknowledgments and every read waits for R replies out of N replicas. Small quorums answer fast but risk stale reads. W plus R above N forces an overlapping node with the latest data, so reads see the newest write. Dynamo and Cassandra choose eventual consistency instead and reconcile concurrent writes with vector clocks.

  1. R equal 1, W equal N: fast reads, slow writes.

  2. W equal 1, R equal N: fast writes, slow reads.

  3. W plus R above N, usually 2 plus 2 above 3: strong consistency.

04

Outages are normal: detect, detour, resync


Failures are routine at scale, so detection needs two witnesses, not one claim. Gossip spreads heartbeat counters through random peers, and a counter frozen past its deadline marks a server down. Short outages use sloppy quorum and hinted handoff. Dead replicas resync with Merkle trees, comparing hashes top down.

How to read: Follow arrows left to right: s0 gossips its list, peers confirm the stall, s2 is marked down.

flowchart LR S0["s0: heartbeat list"] --> S1["random peer"] S0 --> S3["random peer"] S1 --> S2["s2 marked down"] S3 --> S2
Gossip failure detection: heartbeat lists spread peer to peer until a stalled counter convicts a server.
GotchaVector clocks push conflict resolution to the client, and long server lists cost more. Cap the list length or reconciliation drifts.
Q&A

Check yourself


Q1A partition splits n3 from n1 and n2. Your store keeps accepting writes everywhere. Which CAP choice did you make?
  • AP: availability over consistency, sync when healed
  • CP: consistency over availability, block writes
  • CA: both at once, partitions cannot happen
✓ AP: availability over consistency, sync when healed — Accepting writes on both sides risks divergence, which is the AP tradeoff.
Q2N equals 3. Which setup guarantees strong consistency?
  • W equal 1, R equal 1
  • W equal 1, R equal 2
  • W equal 2, R equal 2
✓ W equal 2, R equal 2 — W plus R equals 4, above N, so quorums always overlap on fresh data.
Q3s2 goes down briefly during writes. What keeps the store available?
  • Block all writes until s2 returns
  • Sloppy quorum: healthy servers cover, hinted handoff replays later
  • Delete s2 and rehash every key
✓ Sloppy quorum: healthy servers cover, hinted handoff replays later — Temporary stand-ins plus replay beat blocking the whole quorum.
Sources: System Design Interview Vol 1: Ch. 6, Design a key-value store (pp. 87-109)