On this page
System Design — Key-Value Store
Last reviewed 11 Sept 2026
Part of the system design series. See the framework and building blocks first if you haven’t.
1. Requirements
Functional
put(key, value)andget(key)— single-key operations, no joins, no multi-key transactions.- Support replication for durability and read availability.
- Support horizontal scaling — adding a node increases capacity without downtime.
Non-functional
- Highly available for both reads and writes, even during network partitions (the classic Dynamo motivation: Amazon’s shopping cart couldn’t reject a write just because a few nodes were unreachable).
- Low, predictable latency (single-digit ms p99 for point lookups).
- Tunable consistency — the same store should support “give me the latest write, I’ll wait” and “give me something now, I don’t care if it’s a few seconds stale.”
- Scale to billions of keys / tens of TB without a full-cluster rebuild per node added.
2. High-level architecture
flowchart LR
C[Client] --> CO[Coordinator node<br/>any node can take the request]
CO -->|hash key onto ring| R{Consistent hash ring}
R --> N1[(Replica 1)]
R --> N2[(Replica 2)]
R --> N3[(Replica 3)]
N1 & N2 & N3 -->|ack| CO
CO -->|success once W acks received| C
Any node can act as coordinator for any request — there’s no single leader to bottleneck or fail over. The coordinator hashes the key onto a consistent-hash ring (this is where that piece plugs in directly) to find the key’s preference list — the N nodes responsible for it — and fans the request out to them.
3. Partitioning and replication — the core trade-off
This is Dynamo’s central design, and it’s exactly what Cassandra, Riak, and (internally) DynamoDB inherited.
| Parameter | What it controls |
|---|---|
| N | Replication factor — how many nodes hold a copy of each key |
| W | Write quorum — how many replicas must ack before a write succeeds |
| R | Read quorum — how many replicas must respond before a read returns |
The rule that makes reads see the latest write: R + W > N. With N=3, W=2, R=2, any read quorum overlaps any write quorum by at least one node, so a read is guaranteed to see the most recent successful write — this is “quorum consistency,” strictly weaker than linearizability but strong enough for most applications and much more available than waiting on all N.
4. Deep dive — conflict resolution and hinted handoff
Why conflicts happen at all. With W < N and no single leader, two clients can concurrently write to different replicas of the same key during a partition, or a replica can miss a write while down and later rejoin with a stale value. A leaderless system with high availability has to accept this rather than block on consensus for every write.
Vector clocks (Dynamo’s original answer). Each value carries a vector clock — a set of (node, counter) pairs recording which nodes have written it and how many times. On read, if one version’s vector clock strictly dominates another’s, the older is discarded automatically. If two versions have divergent clocks (concurrent writes neither node knew about), both are returned to the client — “siblings” — and the application resolves them (Amazon’s cart example: union the two carts). This pushes conflict resolution up to whoever understands the business semantics, rather than guessing at the storage layer.
Last-write-wins (the simpler, more common answer today). Attach a timestamp to each write and keep the highest one on conflict. Simpler to reason about and implement, but silently drops data on true concurrent writes — acceptable for caches and many workloads, not acceptable for anything like a shopping cart or a counter. Cassandra defaults to LWW at the cell level.
Hinted handoff. If a replica in the preference list is down, the coordinator writes to the next healthy node instead, tagged with a “hint” saying which node it’s really meant for. When the original node comes back, the hint is replayed to it. This keeps writes succeeding at W even during a transient node outage, without blocking.
5. What real systems do today
- Amazon DynamoDB — the managed evolution of the original Dynamo paper; still consistent-hash partitioned internally, but replaced manual vector-clock conflict resolution with a single-master-per-partition model plus optional strongly-consistent reads, trading some of the original design’s write-availability-under-partition for a much simpler API (no client-side sibling resolution) — a deliberate simplification once AWS controlled the whole deployment instead of shipping the model to arbitrary operators.
- Apache Cassandra — the most direct architectural descendant of the Dynamo paper (its docs literally cite it): consistent hashing with virtual nodes (“tokens”,
num_tokensper node), tunable per-query consistency levels (ONE,QUORUM,LOCAL_QUORUM,ALL), gossip-based membership/failure detection between peer nodes (no coordinator election needed), and LWW-with-timestamp conflict resolution at the cell level rather than vector clocks. - Multi-datacenter reads — both Cassandra and DynamoDB Global Tables favor
LOCAL_QUORUM-style semantics in practice: satisfy the quorum within the local datacenter for latency, replicate cross-region asynchronously, and accept that a client can briefly read stale data from a follower region — the same availability-over-consistency call the original Dynamo paper made, just pushed to the datacenter level instead of the node level.
6. Scaling & failure
- Bottleneck: hot partition (one key getting disproportionate traffic) → fix: salt the key with a small random suffix for very hot keys and fan reads out across the salted copies, or cache the hot key in front of the store entirely — the ring balances key space, not per-key request volume (see the consistent hashing notes).
- Bottleneck: read amplification from R > 1 on every request → fix: most systems only send the “digest” request (a hash, not full data) to the extra R-1 replicas and only fetch full payload from one, comparing digests to detect divergence — cuts network cost while still doing read repair.
- Read repair — when a quorum read detects replicas disagree, the coordinator pushes the winning value back to the stale replicas in the background, so divergence self-heals over time without a dedicated repair job for every case (anti-entropy processes like Merkle-tree comparison handle the rest on a schedule).
- What happens when a replica node dies: writes at W keep succeeding against the remaining healthy replicas in the preference list (as long as enough are up to satisfy W) plus hinted handoff to a stand-in node; reads at R similarly route around it. The system stays available; what degrades is the effective replication factor until the node recovers or is replaced, which is why operators alert on “nodes down” well before it threatens quorum, not after.
- What happens when an entire preference list is unreachable (network partition splits the responsible nodes from the coordinator): with W or R unsatisfiable, that specific request fails or, in
ONE/ANY-consistency configurations, is best-effort accepted anyway and reconciled later — this is the direct CAP trade-off: Dynamo-family stores choose availability, accepting temporary inconsistency, over rejecting the write.
Interview follow-ups
- “Walk me through R + W > N and why it matters.” — Quorum overlap guarantees a read sees the latest committed write; use N=3, W=2, R=2 as the concrete example.
- “Two clients write the same key concurrently during a partition — what happens?” — Depends on conflict strategy: vector clocks return both versions as siblings for the app to merge; LWW keeps the later timestamp and silently drops the other.
- “What’s hinted handoff for, and what’s its failure mode?” — Keeps writes succeeding through a transient node outage by parking the write elsewhere temporarily; combined with LWW it can cause silent data loss on true concurrent writes to different stand-ins.
- “Why would you pick vector clocks over last-write-wins, or vice versa?” — Vector clocks preserve all concurrent writes for app-level merge (needed when losing data is unacceptable, e.g. a cart); LWW is simpler and fine when last-write-wins semantics are actually correct for the data (e.g. a cache, a “last seen” timestamp).
- “How do you avoid a hot key overwhelming one partition?” — Key salting or an in-front cache; the ring balances key space, not per-key traffic.
- “Single-leader relational replica vs this leaderless design — when would you pick which?” — Leaderless/quorum for high write availability and simple key-value access patterns; single-leader when you need strong consistency or multi-key transactions the KV model doesn’t offer.
- “DynamoDB moved away from vector clocks and client-resolved siblings — why would a managed service do that?” — Controlling the whole deployment let AWS trade some of the original availability-under-partition guarantee for a much simpler, foot-gun-free client API.
Sources: Dynamo: Amazon’s Highly Available Key-value Store (paper) · Dynamo — Apache Cassandra Documentation · Leaderless Replication: Quorums, Hinted Handoff and Read Repair · Internals of DynamoDB — Medium · DynamoDB vs. Cassandra: a Complete Comparison in 2025 — Bytebase