Design a Distributed Key-Value Store
A Dynamo-style store that keeps 10 billion keys on dozens of machines, survives node and rack failures, and lets each caller trade consistency for latency with N, R and W.
Last updated 2026-09-30. Difficulty: hard. Patterns: consistent-hashing, quorum-replication, conflict-resolution, gossip, lsm-tree. Reported at Amazon, Google, Meta, Microsoft, Databricks, Snowflake.
The interviewer asks, the candidate answers and draws, and you press Next. Pause to answer yourself at the key decisions, and ask the AI Mentor anything along the way.
Functional requirements
- put(key, value) and get(key). Opaque values up to ~1 MB, keys up to a few hundred bytes. delete(key) is a put of a tombstone.
- Tunable consistency per request. Callers choose how many replicas must answer a read (R) and acknowledge a write (W), trading latency against freshness.
- Always writable. Writes succeed even while some replicas are down or partitioned, as in Amazon’s shopping-cart use case.
- Conflict handling. Concurrent writes to the same key on different replicas are detected and resolved, or returned to the caller to merge.
- Elastic membership. Add and remove nodes without downtime; data moves in the background.
- Out of scope. Range scans and secondary indexes, multi-key transactions, and cross-region replication beyond a mention.
Non-functional requirements
- Scale (10 B keys · ~10 TB raw · 1 M ops/s). About 70 % reads. With 3 replicas, 30 TB stored and roughly 900 k replica writes a second.
- Latency (p99 < 10 ms for get and put). Clients are other services on the same network. Every hop and every disk seek counts.
- Availability (writes succeed with any W replicas reachable). This is the requirement that drives the hardest trade-off: staying writable during failures means accepting that replicas can diverge, so the system must detect, repair and resolve divergence on its own.
- Durability (no acknowledged write lost). An acknowledged write is on W replicas’ commit logs, in different racks, before the client hears success.
- Decentralised. No single master: any node can coordinate any request, and there is no node whose loss stops the cluster.
- Operability. Nodes fail weekly at this size. Detection, repair and rebalancing must be automatic.
Back-of-envelope estimates
- Stored data with replication: ~30 TB. 10 B keys × ~1 KB average value and overhead = 10 TB; × 3 replicas = 30 TB.
- Nodes for capacity: ~20. With 2 TB of usable SSD per node and room for compaction (keep disks half full): 30 TB ÷ 1.5 TB ≈ 20 nodes. Capacity is not the binding constraint.
- Nodes for throughput: ~60. 1 M client ops/s: 700 k reads touching R = 2 replicas and 300 k writes touching N = 3 ≈ 2.3 M replica ops/s. At ~40 k ops/s per node, ~58 nodes. Throughput, not disk, sizes the cluster.
- Commit log bandwidth per node: ~15 MB/s. 300 k writes/s × 3 replicas × 1 KB ≈ 900 MB/s across 60 nodes ≈ 15 MB/s per node of sequential appends. Easy for an SSD; this is why logs are append-only.
- Virtual nodes per physical node: ~256. With 256 vnodes each on 60 nodes the ring has ~15 k ranges. Each node owns 256 small ranges scattered around the ring, so load evens out and a failed node’s ranges are picked up by many peers.
- Failures to expect: ~1 node a week. Servers fail a few percent a year each; across 60 nodes plus disks, network cards and kernel panics, about 1 node is down at some point in any given week. Failure handling is the normal path.
Components
- Client (partition-aware library): A service using the store. Its client library caches the ring and sends each request straight to a node that holds the key, avoiding an extra hop through a load balancer.
- Coordinator (any node): The node that receives a request. It finds the key’s replicas on the ring, fans the request out, waits for R or W answers, resolves versions, and triggers read repair. Every node can play this role.
- Partition ring (consistent hashing + vnodes): Maps hash(key) to a position on a ring of virtual nodes; the next N distinct physical nodes clockwise, in different racks, form the key’s preference list. Every node holds a copy, kept current by gossip.
- Gossip + failure detector: Each node exchanges membership and heartbeat state with a few random peers every second, so changes spread through the cluster in seconds. A phi-accrual detector turns missed heartbeats into a suspicion level per peer.
- Replica A: First node in the key’s preference list. Like every node, it runs the storage engine shown on the right: commit log, memtable, SSTables and compaction.
- Replica B: Second replica, in a different rack from A, so one rack failure never takes out two copies.
- Replica C: Third replica, in a third rack. With N = 3, W = 2 and R = 2, any read overlaps any acknowledged write on at least one replica.
- Fallback node (holds hints): The next healthy node after the preference list. When a replica is down it accepts that replica’s writes with a hint saying who they are for, and hands them back when the replica returns.
- Anti-entropy (Merkle trees): Periodically compares replicas range by range using Merkle trees, so only the parts that differ are exchanged. Repairs divergence that hints and read repair missed.
- Commit log (append-only, fsync): Every write is appended here and synced before the replica acknowledges it. Replayed on restart to rebuild the memtable. Sequential, so cheap.
- Memtable (sorted, in memory): The latest writes in a sorted in-memory structure. Reads check it first. When it reaches ~128 MB it is frozen and flushed to disk as an SSTable.
- SSTables (immutable · Bloom filters): Immutable sorted files on SSD, each with a sparse index and a Bloom filter, so a read skips files that cannot contain the key and touches about one file per lookup.
- Compaction: Merges SSTables in the background, keeping the newest version of each key and dropping tombstones once they are older than the repair window. Keeps read amplification and disk usage bounded.
User flows
- Write a key with N = 3, W = 2. The happy path: hash to a preference list, write to all replicas in parallel, acknowledge once two have the write on disk. The third catches up on its own.
- The client library hashes the key and sends the put to a node that holds it. The library keeps a copy of the ring refreshed every few seconds, so it can pick the first replica directly. If its copy is stale, any node can still coordinate the request.
- The coordinator looks up the key’s preference list: three nodes in three racks. hash(key) lands on a position; walking clockwise, the first three vnodes that belong to distinct physical nodes in distinct racks are the replicas.
- It stamps a new version and sends the write to all three replicas at once. The new version vector increments the coordinator’s own counter on top of the context the client supplied, recording that this write descends from what the client read.
- Each replica appends to its commit log, syncs it, and inserts into the memtable. The append is sequential and the fsync is batched every couple of milliseconds across writes (group commit), which keeps durability cheap. Nothing is written to SSTables on the write path.
- When two replicas acknowledge, the coordinator returns success. W = 2 means the write survives the loss of any one replica. The slowest replica does not hold up the response, which keeps p99 low; it acknowledges a moment later, or is repaired if it never does.
- Read a key with R = 2, and repair on the way. Ask enough replicas to overlap the last write, return the newest version, and quietly fix any replica that was behind.
- The coordinator sends the read to the replicas in the preference list. It asks for the full value from the closest replica and only a digest (hash and version) from the others, which saves bandwidth when they agree.
- Each replica checks its memtable, then its SSTables, using Bloom filters to skip files. A Bloom filter answers "definitely not here" for most files, so a read usually touches one SSTable. The newest version found wins locally.
- With two answers in hand, the coordinator compares versions. If one version vector descends from the other, the newer one wins. If neither descends from the other, the writes were concurrent: both values (siblings) are returned with a merged context, and the client resolves them, such as by merging two shopping carts.
- A replica that returned an older version is sent the newest one: read repair. Done asynchronously after responding, so it never adds latency. Hot keys are therefore repaired almost immediately after any divergence.
- A replica is down during writes. Sloppy quorum and hinted handoff keep the key writable, then deliver the missed writes when the node returns.
- Gossip marks replica C as suspected down after its heartbeats stop. The phi-accrual detector raises suspicion as heartbeats go missing relative to their normal timing, instead of using a fixed timeout, so a slow network does not flap nodes up and down.
- The coordinator writes to A and B, and to the next healthy node on the ring with a hint for C. This is a sloppy quorum: W = 2 can be met by A and B alone, and the fallback still takes a copy so N copies exist. The hint records that the data belongs to C.
- When C comes back, gossip spreads the news and the fallback hands the hinted writes to C. Hints are replayed in the background at a throttled rate, so a node returning after an hour is not flattened by a burst. Once delivered, the fallback deletes them.
- If C never returns, it is removed from the ring and its ranges are rebuilt from the other replicas. Removal is an operator action (or an automated one after a long timeout), never automatic on a short failure. With vnodes, many nodes each stream a small part of C’s data, so the rebuild is fast.
- Add ten nodes to a busy cluster. Elastic growth with consistent hashing: new nodes take small ranges from many peers, data streams in the background, and clients never see an error.
- A new node joins and announces itself through gossip. It picks 256 random tokens on the ring (its vnodes). Each token takes over a small slice of a range from whichever node owned it before.
- For each slice it will own, it streams the data from a current replica. Whole SSTables are streamed rather than individual keys, which is fast and avoids loading the write path. Because the new node takes slices from dozens of peers, no single peer is overloaded.
- While streaming, writes for those ranges go to both old and new owners. The new node is a "pending" replica: it receives writes but is not counted toward W or used for reads until streaming finishes, so a half-loaded node can never answer a read.
- Once complete, the ring marks it normal; clients pick up the new ring within seconds. Old owners drop the data they no longer own in a cleanup pass. About 1/7 of the data moved for a 60-to-70 node expansion, which is the minimum possible.
- Compaction and anti-entropy in the background. Two maintenance loops that keep the store healthy: compaction keeps reads fast and disks bounded; Merkle-tree repair makes replicas converge even for keys nobody reads.
- When the memtable fills, it is flushed to a new immutable SSTable. A flush writes one sorted file sequentially, then the matching commit log segment can be deleted.
- Compaction merges SSTables, keeping the newest version of each key. Leveled compaction keeps reads to about one file per level at the cost of more rewriting; size-tiered compaction writes less but reads more. Tombstones are dropped only after the repair window, so a deleted key cannot be resurrected by a replica that missed the delete.
- Anti-entropy builds a Merkle tree per range on each replica and compares them. Matching roots mean the whole range is identical, compared with one hash. Where they differ, the comparison walks down the tree to the few leaves (small key ranges) that differ.
- Only the differing ranges are streamed between replicas. A full repair cycle over all data runs within the tombstone grace period (typically every few days), which is the rule that keeps deletes from coming back.
Deep dives
Spreading keys across nodes
How do you decide which nodes hold a key, so that adding a node moves as little data as possible?
The simple answer, hash(key) mod number_of_nodes, moves almost every key when a node is added: going from 60 to 61 nodes relocates about 98 % of the data. In a cluster that grows and loses nodes every week, that is unusable.
Consistent hashing places nodes and keys on the same ring, and a key belongs to the next node clockwise. Adding a node only takes keys from its neighbours. The remaining problem is unevenness: with one position per node, ranges vary wildly in size and a failed node dumps all its load on one neighbour.
- Consistent hashing with many virtual nodes per physical node, rack-aware replica placement chosen
- Fixed number of slots (say 16,384) assigned to nodes situational: systems with a central control plane, such as Redis Cluster or Distributed Cache
- Range partitioning on the key itself situational: when range queries are required, as in Bigtable or Spanner
- hash(key) mod N rejected
The answer: Each physical node takes ~256 random tokens on a 64-bit hash ring. A key’s replicas are found by walking clockwise from hash(key) and taking the first N tokens that belong to distinct physical nodes in distinct racks (and, if needed, distinct zones). The ring is small (60 nodes × 256 tokens) and every node and client library holds a copy, updated by gossip, so routing needs no lookup service. Adding a node moves about 1/N of the data, streamed from many peers at once.
Why random tokens rather than evenly spaced ones?
Random tokens need no coordination: a joining node just picks its own. With 256 per node, the law of large numbers keeps ownership within a few percent of even. Some systems use a smarter token allocator to get tighter balance with fewer vnodes, which reduces repair and gossip overhead.
One key is extremely hot. Does consistent hashing help?
No: a key lives on exactly N replicas, so a single hot key is bounded by those nodes. Mitigate above the store: cache it in the client, read from any replica (R = 1) if staleness is acceptable, or split the key into sub-keys in the application.
How do you make sure two replicas never end up in the same rack?
The placement walk skips tokens whose node is in a rack already used for this key. The ring metadata includes each node’s rack and zone, so this is a local computation on every node.
What does the ring cost to keep in sync?
About 15 k tokens with a few bytes each: well under a megabyte. Gossip sends only versions and deltas, so steady-state traffic is tiny. Ring changes converge across 60 nodes in a few gossip rounds, a few seconds.
N, R and W
What do N, R and W mean, and which values would you choose?
N is how many replicas hold each key. A write waits for W of them to acknowledge; a read waits for R of them to answer. If R + W > N, every read set overlaps every acknowledged write set on at least one replica, so a read sees the latest acknowledged write, as long as the set of replicas has not changed.
Lower R or W means faster and more available operations but more chance of stale reads. Different callers want different trade-offs, so the store lets them choose per request.
- N = 3, W = 2, R = 2 by default; per-request overrides chosen
- W = 1, R = 1 situational: caches, metrics and session data where staleness is harmless
- W = 3 (all), R = 1 situational: read-heavy data that is rarely written and must be fresh
- Consensus per key range (Raft or Paxos leader) situational: when strong consistency matters more than always-writable, as in etcd or Spanner
The answer: Default to N = 3 across three racks, W = 2, R = 2. Writes acknowledge as soon as two replicas have them on their commit logs; reads return once two replicas answer, with digests from the others for read repair. Callers can lower R or W for latency-sensitive, staleness-tolerant data, or raise them. Be explicit that this is not linearizable: overlapping quorums give "you see the latest acknowledged write" in steady state, while sloppy quorums during failures and concurrent writers can still produce stale reads or siblings, which the conflict handling covers. This is the model of Amazon’s Dynamo paper and of Cassandra and Riak.
With R + W > N, can a read still be stale?
Yes, in three ways: during a sloppy quorum the write may have gone to a fallback node outside the read set; a write that failed to reach W may still have landed on one replica and be seen by some reads but not others; and concurrent writes have no single latest. R + W > N is a steady-state guarantee, not linearizability.
How would you offer strong consistency for a few keys?
Route those operations through a consensus path: lightweight transactions using Paxos on the key’s replicas (as Cassandra does for compare-and-set), at the cost of several round trips. Keep it opt-in, because making every write pay for consensus would lose the latency and availability the design was built for.
Why not wait for all three replicas on writes to be safe?
Because then the slowest of three disks or networks sets your p99, and any single failure stops writes. W = 2 already survives one failure; the third replica catches up through hints, read repair and anti-entropy.
How do you replicate across regions?
Treat each region as a set of replicas (for example three per region) and let callers choose LOCAL_QUORUM (fast, region-local) or a cross-region quorum. Cross-region writes add 50–150 ms, so most data uses local quorums with asynchronous replication to other regions and conflict handling for the rare concurrent writes.
Concurrent writes and conflicts
Two clients update the same key at the same time on different replicas. Which write wins?
Because the store stays writable during partitions, two replicas can each accept a different write for the same key, neither knowing about the other. When they meet again, the system must decide whether one write supersedes the other or whether they are truly concurrent.
Timestamps seem to answer it, but clocks on different machines disagree by milliseconds or worse, so "last write wins" by wall clock silently drops writes that happened later in real time. Detecting concurrency requires tracking causality, not time.
- Version vectors per key; return siblings to the client to merge chosen
- Last write wins by timestamp situational: immutable or idempotent data, or where losing a concurrent update is acceptable (caches, profiles last edited by one device)
- CRDTs (counters, sets, maps that merge automatically) situational: counters, sets and flags where the merge rule is known, as in Riak data types
The answer: Each value carries a version vector: a small map from coordinating node to counter. A write carries the context of the value the client last read; the coordinator increments its own entry. On read, a version that descends from another replaces it; versions where neither descends are concurrent and both are returned as siblings, along with a merged context. The client resolves them (for a cart: union of items, with removals tracked) and writes back, which collapses the siblings. Vectors are pruned by dropping the oldest entries beyond a size limit. Data types that merge naturally (counters, sets) can use CRDTs instead, and data where losing a concurrent write is acceptable can opt into last-write-wins.
Why did Amazon’s shopping cart sometimes resurrect deleted items?
Because merging two concurrent carts by union keeps an item removed in one version and present in the other. The fix is to record removals explicitly (an observed-remove set), so a merge knows the removal happened after the add it removes. It is the classic example of why merges need the right data type.
Isn’t last-write-wins good enough with NTP-synchronised clocks?
NTP keeps clocks within milliseconds on a good day, but skew of tens of milliseconds and occasional jumps happen. Two writes from different clients within that window are ordered arbitrarily, and one is silently lost. For a profile edited from one device, fine; for a shared counter or cart, not.
How big do version vectors get?
One entry per node that has coordinated a write to that key. With clients routing to the first replica, that is usually one to three entries. Pruning beyond ten or so entries, by dropping the oldest, risks a false conflict but never a lost write.
Membership and failure handling
With no master, how does each node know who is alive, and what happens to writes for a dead replica?
A central coordinator would know exactly who is up, but it would also be a single point of failure and a bottleneck, which the design rules out. Every node must instead build its own view of the cluster, and those views must converge quickly.
Failures are mostly temporary: a restart, a GC pause, a flaky switch. Treating every missed heartbeat as permanent would trigger expensive data movement constantly. The design needs to separate "temporarily unreachable" from "gone".
- Gossip membership with phi-accrual detection; sloppy quorum with hinted handoff for temporary failures; explicit removal for permanent ones chosen
- Central membership service (ZooKeeper / etcd) situational: systems that already need consensus, such as those with leader-based replication
- Strict quorum only (no fallbacks) rejected
The answer: Every second, each node gossips its membership table (node, status, heartbeat counter, tokens) with a few random peers; changes reach the whole cluster in O(log n) rounds. A phi-accrual detector on each node estimates how unusual the current silence is given past heartbeat timing, and marks a peer down above a threshold, which adapts to slow networks better than a fixed timeout. When a replica is down, coordinators write to the next healthy node with a hint; the hint is handed back when the node returns. Hints are kept for a few hours; beyond that, anti-entropy repairs the node. Permanently removing a node, and rebuilding its ranges, is an explicit decision.
A node has a long GC pause and is marked down, then comes back. What happens?
Coordinators had sent its writes to fallbacks with hints. When gossip sees its heartbeat counter advance again, it is marked up and the hints are replayed to it. Reads may have skipped it meanwhile. No data moved permanently, which is why short outages must not trigger rebalancing.
What if the fallback node holding hints also dies?
The hints are lost, but the write is still on the other W replicas. The returning node is behind on those keys until read repair or anti-entropy catches it up. Hints are an optimisation for fast catch-up; anti-entropy is the guarantee.
How do you avoid a network partition splitting the cluster into two that both think they are the whole?
In this design both sides keep serving: each side writes to whatever replicas it can reach (sloppy quorum), and conflicts are reconciled when the partition heals. That is the deliberate availability choice. Operations that cannot tolerate it use the consensus path, which only the majority side can complete.
The storage engine on each node
Each node takes tens of thousands of writes a second. B-tree or LSM tree, and why?
A B-tree updates pages in place: each write may read and rewrite a page somewhere on disk, which is random I/O. At 15 k writes a second per node that is a lot of random writes, and SSDs wear faster under them.
A log-structured merge tree never updates in place. Writes go to a sequential log and an in-memory table; on disk, data lives in immutable sorted files that are merged in the background. Writes become sequential and cheap, at the cost of reads that may check several files and of background compaction work.
- LSM tree: commit log + memtable + immutable SSTables with Bloom filters + compaction chosen
- B-tree (InnoDB-style) situational: read-heavy workloads with range queries
- Pure in-memory hash table with snapshots rejected
The answer: Each node runs an LSM engine (RocksDB-style). A write appends to the commit log with group commit (fsync every couple of milliseconds) and inserts into a sorted memtable; at ~128 MB the memtable is flushed to an SSTable. Each SSTable has a sparse index and a Bloom filter (about 10 bits per key, 1 % false positives), so a read checks the memtable and usually one file. Leveled compaction keeps read amplification low for this read-heavy (70 %) workload; tombstones are dropped only after the repair grace period. Compaction throughput is rate-limited so it does not hurt p99.
Why can’t compaction simply drop tombstones immediately?
Because a replica that missed the delete still has the old value. If every other replica forgets the tombstone, anti-entropy would see the old value on that replica and copy it back: the deleted key is resurrected. Keeping tombstones longer than a full repair cycle guarantees every replica has seen the delete first.
Compaction falls behind under heavy writes. What do users see?
More SSTables per read, so read latency rises, and disk usage grows. The engine throttles writes (stalls) when level-0 file counts get too high, which protects reads but raises write latency. Monitor pending compaction bytes and add nodes before it becomes chronic.
How does a node recover after a crash?
It replays the commit log segments not yet flushed into a fresh memtable, then rejoins. SSTables are immutable, so they need no recovery. Writes it missed while down come back through hints, read repair and anti-entropy.
Making replicas converge
Read repair fixes keys that are read. What about the billions that are not?
Divergence happens: lost hints, failed writes that reached only one replica, disks restored from old snapshots. Read repair only fixes keys that are read, and cold keys may never be. Over months, replicas would quietly drift apart.
Comparing replicas key by key means shipping billions of keys or hashes over the network. The trick is to compare hashes of ranges first and only descend into the ranges that differ.
- Merkle trees per token range, compared periodically; stream only differing leaves chosen
- Read repair only rejected
- Full data comparison rejected
The answer: Each replica builds, per token range, a Merkle tree whose leaves hash small sub-ranges of keys and whose parents hash their children. Two replicas compare roots; if equal, the range is identical. Otherwise they walk down to the differing leaves and stream just those keys, keeping the newer versions (or siblings). A full cycle across all ranges runs every few days, staggered and rate-limited, and always within the tombstone grace period. Incremental repair marks already-repaired SSTables so later cycles only look at new data.
How deep should the tree be?
Deep enough that a differing leaf covers a small range: with 2^15 leaves over a range holding a hundred million keys, each leaf covers a few thousand keys. Deeper trees find differences more precisely but cost more memory and build time.
How would you know repair is keeping up?
Track, per range, when it was last fully repaired, and alert if any range approaches the tombstone grace period. Also track how much data each repair streams: a sudden rise means something upstream (hints, writes) is failing.
Isn’t building trees expensive?
It reads all the data in the range, so it is throttled and staggered across nodes and ranges. Incremental repair helps most: data already confirmed identical is skipped, so each cycle mainly reads recent SSTables.
The theory behind it
- Replication and consistency: Leader-follower, multi-leader and leaderless replication; synchronous vs asynchronous; quorums; CAP and PACELC; consistency models from eventual to linearizable; read-your-writes and failover.
- Sharding and partitioning: Horizontal partitioning strategies (range, hash, directory), consistent hashing, choosing a shard key, hot shards, cross-shard queries, and resharding without downtime.
- Consensus, leases and coordination: Leader election, Raft and Paxos at interview depth, ZooKeeper and etcd, distributed locks and why they are fencing tokens, plus how to avoid needing coordination at all.