SysDesignPrep.com
Interviewer kit

Design a Distributed Key-Value Store

Run this for someone else. You hold the answers; they do not. Read the prompt, keep the clock, and use the probes below when an answer is thin. Do not show them this page.

The candidate should have a blank page and a whiteboard, not this.

Open with this

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. Take a couple of minutes on requirements, then we will do some numbers, then the design. I will interrupt to keep us moving.

The clock

  • 4 min: functional requirements and scope
  • 4 min: non-functional requirements, with numbers
  • 5 min: back-of-envelope estimates
  • 16 min: high-level design and one or two flows
  • 16 min: deep dives and the close

Move them on out loud when a section overruns. The commonest failure is spending twenty minutes on requirements and never reaching a deep dive, and preventing that is your job as much as theirs.

Requirements · 8 min

Listen for: a scoped set of capabilities, an explicit out-of-scope list, and numeric targets rather than adjectives. Prompt with “what are you not building?” if they never scope, and “what number would make that requirement real?” if they say “fast” or “highly available”.

Functional (6)
  • 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 (6)
  • 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.

Estimates · 5 min

Ask for two or three numbers, not all of them. What matters is whether they state assumptions, round sensibly, and say what the number implies. Push once with “where did that come from?”

The numbers (6)
  • 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.

High-level design · 16 min

Let them draw. Interrupt only to ask what backs a component or what a box actually does. Then pick one flow below and ask them to walk it end to end.

Components (13)
  • 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.
Flows to ask them to walk (5)
  1. 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.
    1. The client library hashes the key and sends the put to a node that holds it.
    2. The coordinator looks up the key’s preference list: three nodes in three racks.
    3. It stamps a new version and sends the write to all three replicas at once.
    4. Each replica appends to its commit log, syncs it, and inserts into the memtable.
    5. When two replicas acknowledge, the coordinator returns success.
  2. 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.
    1. The coordinator sends the read to the replicas in the preference list.
    2. Each replica checks its memtable, then its SSTables, using Bloom filters to skip files.
    3. With two answers in hand, the coordinator compares versions.
    4. A replica that returned an older version is sent the newest one: read repair.
  3. A replica is down during writes: Sloppy quorum and hinted handoff keep the key writable, then deliver the missed writes when the node returns.
    1. Gossip marks replica C as suspected down after its heartbeats stop.
    2. The coordinator writes to A and B, and to the next healthy node on the ring with a hint for C.
    3. When C comes back, gossip spreads the news and the fallback hands the hinted writes to C.
    4. If C never returns, it is removed from the ring and its ranges are rebuilt from the other replicas.
  4. 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.
    1. A new node joins and announces itself through gossip.
    2. For each slice it will own, it streams the data from a current replica.
    3. While streaming, writes for those ranges go to both old and new owners.
    4. Once complete, the ring marks it normal; clients pick up the new ring within seconds.
  5. 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.
    1. When the memtable fills, it is flushed to a new immutable SSTable.
    2. Compaction merges SSTables, keeping the newest version of each key.
    3. Anti-entropy builds a Merkle tree per range on each replica and compares them.
    4. Only the differing ranges are streamed between replicas.

Deep dives · 16 min

Pick two. Ask the headline question, let them answer, then use the follow-ups. The follow-ups are where the level gets decided, so leave time for at least three of them.

Spreading keys across nodes

Ask: How do you decide which nodes hold a key, so that adding a node moves as little data as possible?

Good answers name: Consistent hashing with many virtual nodes per physical node, rack-aware replica placement, Fixed number of slots (say 16,384) assigned to nodes, Range partitioning on the key itself, hash(key) mod N.

Our pick: 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.

  1. 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.
  2. 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.
  3. 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.
  4. 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

Ask: What do N, R and W mean, and which values would you choose?

Good answers name: N = 3, W = 2, R = 2 by default; per-request overrides, W = 1, R = 1, W = 3 (all), R = 1, Consensus per key range (Raft or Paxos leader).

Our pick: 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.

  1. 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.
  2. 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.
  3. 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.
  4. 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

Ask: Two clients update the same key at the same time on different replicas. Which write wins?

Good answers name: Version vectors per key; return siblings to the client to merge, Last write wins by timestamp, CRDTs (counters, sets, maps that merge automatically).

Our pick: 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.

  1. 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.
  2. 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.
  3. 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

Ask: With no master, how does each node know who is alive, and what happens to writes for a dead replica?

Good answers name: Gossip membership with phi-accrual detection; sloppy quorum with hinted handoff for temporary failures; explicit removal for permanent ones, Central membership service (ZooKeeper / etcd), Strict quorum only (no fallbacks).

Our pick: 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.

  1. 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.
  2. 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.
  3. 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

Ask: Each node takes tens of thousands of writes a second. B-tree or LSM tree, and why?

Good answers name: LSM tree: commit log + memtable + immutable SSTables with Bloom filters + compaction, B-tree (InnoDB-style), Pure in-memory hash table with snapshots.

Our pick: 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.

  1. 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.
  2. 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.
  3. 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

Ask: Read repair fixes keys that are read. What about the billions that are not?

Good answers name: Merkle trees per token range, compared periodically; stream only differing leaves, Read repair only, Full data comparison.

Our pick: 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.

  1. 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.
  2. 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.
  3. 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.

Close · 5 min

Ask what breaks first at ten times the load, and what they would build next. Then give them your read: one thing that was strong, one thing that was missing, one thing to practise. Be specific; “good job” helps nobody.

Open in your browser to sign in

Google does not allow sign-in inside this app's built-in browser. Open this page in Safari and sign in there. The link opens this same page.

Tap the ⋯ or share button at the top or bottom of the screen, then Open in browser. Or copy the link and paste it into Safari.