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
CAP forces a choice when partitions strike: CP blocks writes to stay consistent, AP keeps serving and reconciles later.
Keys land on a hash ring and copy to the next N servers; quorum W plus R above N guarantees readers see the latest write.
Gossip detects failures, sloppy quorum plus hinted handoff cover short outages, Merkle trees resync full replicas.
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.
Hash key0 onto the same ring the servers share.
Walk clockwise: the first server met owns the key.
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.
R equal 1, W equal N: fast reads, slow writes.
W equal 1, R equal N: fast writes, slow reads.
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.