Design a KV store like Dynamo: PUT key → value, GET key → value, at billions of keys and millions of ops/sec with node failures handled gracefully.
Requirements
Functional requirements
put(key, value)— store or overwrite a value.get(key)— retrieve value or not-found.delete(key)— remove a key.- Support values from bytes to megabytes.
Non-functional requirements
- Availability: 99.99% — tolerate node and rack failures.
- Latency: p99 GET under 10 ms; p99 PUT under 20 ms.
- Scale: Billions of keys; horizontal scale by adding nodes.
- Durability: Survive node crash without data loss (replication factor 3).
- Consistency: Tunable — strong for critical keys, eventual for cache-like keys.
Estimation
- 1B keys × 1 KB average = 1 TB raw data × replication factor 3 = 3 TB cluster.
- 1M ops/sec peak → ~1M GET/s + 100K PUT/s — each node handles ~50K ops with 20 nodes.
Hash ring0 … 2^32-1
Node Akeys 0–25%
Node Bkeys 25–50%
Node Ckeys 50–75%
Node Dkeys 75–100%
New node E joins → only ~1/N keys move, not all keys.
Data partitioning
Consistent hashing in practice
- Hash key onto ring; assign to next node clockwise.
- Virtual nodes (100–200 per physical node) for even load distribution.
- Adding a node moves only ~1/N keys — not full reshuffle.
- Hot keys: replicate to multiple nodes or add random suffix sharding.
def get_responsible_nodes(key: str, rf: int = 3) -> list[str]:
hash_val = hash_key(key)
ring_position = ring.bisect(hash_val)
nodes = []
for i in range(rf):
nodes.append(ring[(ring_position + i) % len(ring)])
return nodes
Replication and consistency
| Aspect | Quorum read/write (R + W > N) | Leader-follower (strong consistency) |
|---|---|---|
| Write | W nodes must ack (e.g., W=2 of N=3) | Leader acks; async replicate to followers |
| Read | R nodes must respond (e.g., R=2); merge by version vector | Read from leader or linearizable follower |
| Consistency | Tunable — R+W>N gives strong | Strong by default; followers may lag |
| Failure | Tolerates N-W node failures on write | Leader election needed on leader failure |
Write
Quorum read/write (R + W > N)W nodes must ack (e.g., W=2 of N=3)Leader-follower (strong consistency)Leader acks; async replicate to followersRead
Quorum read/write (R + W > N)R nodes must respond (e.g., R=2); merge by version vectorLeader-follower (strong consistency)Read from leader or linearizable followerConsistency
Quorum read/write (R + W > N)Tunable — R+W>N gives strongLeader-follower (strong consistency)Strong by default; followers may lagFailure
Quorum read/write (R + W > N)Tolerates N-W node failures on writeLeader-follower (strong consistency)Leader election needed on leader failure
Architecture
StreamHub production architecture (AWS)
Node internals
- Coordinator: Receives client request; hashes key; routes to responsible nodes.
- Storage engine: LSM-tree (LevelDB/RocksDB) for write-heavy or B-tree for read-heavy.
- Gossip protocol: Nodes discover membership and failure without central registry.
- Hinted handoff: If target node is down, temporary node stores writes until recovery.
- Anti-entropy: Merkle tree comparison between replicas to repair divergence.
# Tunable consistency (pseudo-client config)
consistency:
write_quorum: 2 # W
read_quorum: 2 # R
replication_factor: 3 # N
# R + W > N → strong consistency (2 + 2 > 3)
Handling failures
| Failure | Detection | Recovery |
|---|---|---|
| Node crash | Gossip heartbeat timeout | Promote replica; hinted handoff replays writes |
| Network partition | Quorum can't be reached | Choose CP (reject writes) or AP (accept, repair later) |
| Data divergence | Periodic Merkle tree scan | Anti-entropy repair syncs missing keys |
| Adding node | Admin command or auto-scale | Consistent hash remaps ~1/N keys to new node |
Node crash
DetectionGossip heartbeat timeoutRecoveryPromote replica; hinted handoff replays writesNetwork partition
DetectionQuorum can't be reachedRecoveryChoose CP (reject writes) or AP (accept, repair later)Data divergence
DetectionPeriodic Merkle tree scanRecoveryAnti-entropy repair syncs missing keysAdding node
DetectionAdmin command or auto-scaleRecoveryConsistent hash remaps ~1/N keys to new node
Quick recall
Everything you need if you only revisit this box.
- Partition keys via consistent hashing with virtual nodes for even load.
- Replication factor 3; quorum W + R > N gives strong consistency.
- Gossip protocol for membership; hinted handoff for temporary node failures.
- Anti-entropy (Merkle trees) repairs replica divergence in background.
- Hot keys need replication beyond consistent hashing — random suffix or local cache.
- LSM-tree storage engine for write-heavy; B-tree for read-heavy workloads.
Test yourself
Answer these before moving on — recall is what makes it stick.