System design · Partitioning, replication, quorums

How to design a key-value store

A key-value store has the smallest API in system design: get, put, delete. Everything interesting happens underneath - spreading keys across hundreds of machines, keeping copies in sync, staying writable when a node dies, and deciding which value wins when two writes collide. This guide follows the design in Amazon's 2007 Dynamo paper, which also shaped Cassandra and Riak.

Updated · 7 min read

Requirements

Pin down the scope and, above all, the consistency the caller needs. That one answer drives most of the design.

  • Functional: put(key, value); get(key); delete(key). Keys are short strings; values are opaque blobs up to about 1 MB.
  • Functional, out of scope: range scans, secondary indexes and multi-key transactions. Say so out loud.
  • Non-functional: always writable - a put should succeed even when some nodes are down or the network is split.
  • Non-functional: single-digit millisecond p99; scale by adding nodes; no single point of failure.
  • Non-functional: tunable consistency - the caller chooses between faster, possibly stale reads and slower, up-to-date ones.

Capacity estimates

These are assumptions, stated out loud. The goal is to know how many machines you need and how much traffic each one takes.

QuantityAssumptionResult
Data10B keys × 1 KB average value≈ 10 TB of unique data
Replicasreplication factor N = 3≈ 30 TB stored across the cluster
Nodes≈ 2 TB usable disk per node, kept half full for compaction and growth≈ 30 nodes
Traffic500,000 gets/s and 100,000 puts/s at peak600,000 client requests/s
Replica workeach put goes to 3 replicas; each quorum get reads 2≈ 300,000 replica writes/s and 1M replica reads/s ≈ 43,000 operations/s per node

Three conclusions. No single machine holds the data, so partitioning is mandatory. Replication triples the write load, so the storage engine must handle writes cheaply. And with 30 nodes, something is always failing - failure handling is the normal path, not an edge case.

API design

The client talks to any node, which acts as the coordinator for that request. Consistency is a per-request option, not a cluster-wide setting.

PUT /kv/{key}?w=2
  If-Match: {context}
  body: <bytes>
→ 200 OK  { "context": "vc:[A:3,B:1]" }

GET /kv/{key}?r=2
→ 200 OK  { "values": ["<bytes>"], "context": "vc:[A:3,B:1]" }

DELETE /kv/{key}?w=2
  If-Match: {context}
→ 204 No Content

The context is the version the client last saw. A get can return more than one value when replicas hold conflicting versions; the client merges them and writes the result back.

Data model

Each node is an independent store with an LSM-tree engine. Writes go to an append-only log for durability and to a sorted in-memory table; full tables are flushed to immutable sorted files on disk.

commit log        append-only, fsync'd, replayed on restart
memtable          in memory, sorted by key, ≈ 64 MB
SSTable files     immutable, sorted by key, on disk
  data block      key | version | timestamp | tombstone flag | value
  index block     first key of each data block → offset
  bloom filter    "is this key possibly in this file?"

record version    vector clock [(node, counter), ...] or a timestamp

Deletes don't remove data - they write a tombstone. The tombstone must live long enough to reach every replica, or a stale replica will bring the deleted key back during repair.

High-level design

There is no master. Every node runs the same code and can coordinate any request.

  • Partitioning: consistent hashing places keys on a ring. Each physical node owns many virtual nodes, so load spreads evenly and a new node takes small slices from everyone.
  • Replication: a key is stored on the first N distinct physical nodes clockwise from its hash - its preference list.
  • Coordination: the node that receives a request forwards it to the N replicas and waits for W acks on a put or R replies on a get.
  • Membership: nodes gossip with a few random peers every second, so ring changes and node states spread through the cluster without a central registry.
  • Background repair: hinted handoff covers short outages; anti-entropy with Merkle trees fixes everything else.

Where it breaks

Hash keys evenly and traffic still isn't even. One key - a celebrity profile, a flash-sale item - can take a large share of all reads, and every request for it lands on the same N replicas. Adding nodes doesn't help, because the key lives in one place. DynamoDB, the managed service, limits a single partition to about 3,000 reads and 1,000 writes per second; past that, requests are throttled no matter how much capacity the table has.

For hot reads, cache the key in front of the store or in the client, and read from any replica rather than a quorum when slight staleness is fine. For hot writes, split the key: write to key#1 ... key#10 at random and sum or merge them on read. The cost is more complex reads.

Quorums and tunable consistency

  • With N replicas, a put waits for W acks and a get waits for R replies. If R + W > N, every read set overlaps every write set, so a read sees the latest acknowledged write.
  • N = 3, R = 2, W = 2 is the common default. R = 1, W = 1 is fastest but may return stale data. W = 3 makes writes fail when any replica is down.
  • Cassandra exposes this as per-request levels such as ONE, QUORUM and ALL - the same idea, chosen by the caller.
  • A quorum is not linearizability. Concurrent writes, clock skew and sloppy quorums can still produce anomalies; say this before the interviewer does.

Conflicts: vector clocks or last-write-wins

When two clients write the same key on different replicas, the store must decide what survives. Last-write-wins keeps the value with the highest timestamp. It's simple and Cassandra uses it, but clock skew can silently drop a newer write. Vector clocks record a counter per writing node, so the store can tell whether one version descends from another or whether they are truly concurrent. The original Dynamo returned concurrent versions to the client to merge - a shopping cart takes the union of both. Vector clocks are safer but push work onto the client.

Failures: hinted handoff, Merkle trees and gossip

  • Failure detection: a node is marked down when peers stop hearing from it. Cassandra uses a phi accrual detector, which outputs a suspicion level rather than a yes or no.
  • Hinted handoff: if a replica is down, the coordinator writes to the next healthy node with a hint naming the intended owner. When the owner recovers, the hint is replayed. This is a sloppy quorum - writes stay available, but R + W > N no longer guarantees overlap.
  • Anti-entropy: each replica builds a Merkle tree of hashes over its key ranges. Two replicas compare root hashes and walk down only the branches that differ, so they find and copy the few divergent keys without shipping the whole dataset.

What interviewers look for

  • Consistent hashing with virtual nodes, and why plain modulo hashing reshuffles almost every key when a node is added.
  • The R + W > N rule with concrete values, and a clear trade-off between latency, availability and freshness.
  • A named conflict strategy - last-write-wins or vector clocks - and what each one loses.
  • A failure story covering temporary outages and permanent divergence.
  • A storage engine choice with a reason: LSM trees for write-heavy loads, and a plan for hot keys.

Frequently asked questions

Why use consistent hashing instead of hash modulo the number of nodes?

+

With hash(key) mod N, changing N moves almost every key to a different node. Consistent hashing moves only the keys between the new node and its neighbour on the ring - roughly 1/N of the data. Virtual nodes keep the ranges even.

What does R + W > N mean?

+

N is the number of replicas, W is how many must acknowledge a write, and R is how many must answer a read. If R + W > N, the set of replicas you read from always overlaps the set that took the latest write, so at least one reply holds the newest acknowledged value.

Should a key-value store use an LSM tree or a B-tree?

+

An LSM tree turns every write into a sequential append, which suits write-heavy, replicated workloads - Cassandra, RocksDB and LevelDB use it. A B-tree updates pages in place and gives more predictable reads, which suits read-heavy workloads.

How does a Dynamo-style store stay available when a node is down?

+

Writes go to the next healthy node on the ring with a hint, then get handed back when the owner returns. Reads succeed as long as R replicas answer. Merkle-tree anti-entropy and read repair fix any copies that drift apart.

How do you handle a hot key?

+

Hashing spreads keys, not traffic, so one popular key can overload its replicas. Cache hot reads in front of the store, read from any one replica when staleness is acceptable, and split hot write keys into several sub-keys that are merged on read.

Now break one yourself.

The first challenge takes about two minutes. No signup.