Problem & requirements
The target is a store with small values (say under 10 KB), the ability to hold very large datasets, high availability so it responds even during failures, automatic scaling as servers are added and removed, tunable consistency, and low latency. On a single server the problem is easy: a hash table in memory, with compression and spilling of cold data to disk. That runs out of room quickly, so the real design is distributed.
Distribution forces trade-offs, and the chapter is organized around them: how to partition data, how to replicate it, how strongly to keep replicas in agreement, how to resolve conflicting writes, and how to detect and recover from failures of every size, from a single slow node to a whole data center.
Back-of-the-envelope
Assume 1 billion keys with an average value of 1 KB: about 1 TB of raw data. With a replication factor of 3, the cluster stores 3 TB. Spread over 30 nodes with 1 TB disks, each holds about 100 GB, leaving plenty of room for growth, compaction overhead and rebalancing.
Bloom filters help the read path skip disk files that cannot contain a key. With about 10 bits per key the false-positive rate is roughly 1%. For 1 billion keys that is 10 billion bits, or 1.25 GB of filter memory per full copy of the data (about 3.75 GB cluster-wide with three replicas, or roughly 125 MB per node across 30 nodes), a small price for avoiding most unnecessary disk reads. At 100,000 reads per second, a 1% false-positive rate means about 1,000 wasted SSTable probes per second instead of 100,000.
CAP theorem and the consistency choice
The CAP theorem says that when a network partition occurs, a distributed system must choose between consistency (every read sees the latest write) and availability (every request gets a non-error response). Since partitions cannot be ruled out in practice, the real choice is between CP and AP behavior during a partition. A CP system, such as one guarding bank balances, rejects or blocks requests on the minority side to avoid divergence. An AP system keeps accepting reads and writes everywhere and reconciles later.
Many key-value stores built for shopping carts, sessions and user profiles choose AP with eventual consistency: replicas may briefly disagree, but given no new writes they converge. Others, like those built on Raft or Paxos, choose strong consistency. The requirements of the use case, not the technology, should drive this decision.
Partitioning and replication
Data is partitioned with consistent hashing and virtual nodes, so the cluster can grow automatically and higher-capacity machines can take more virtual nodes. To replicate, a key is stored on the first N servers encountered walking clockwise from its position. With virtual nodes, the first N virtual nodes might belong to fewer than N physical machines, so the walk skips positions until it has N distinct physical servers.
For resilience against larger failures, replicas are placed across racks and data centers connected by fast links. A node handling a request for a key, called the coordinator, forwards it to the replicas responsible for that key.
K hashes to 80°, so C is its first owner and the walk continues clockwise to D and E for three replicas. With virtual nodes the walk would skip any position whose physical server is already in the list.Quorum consensus: N, W and R
With N replicas, a write is considered successful once W replicas acknowledge it, and a read returns once R replicas respond. If W + R > N, every read quorum overlaps every write quorum in at least one node, so a read sees the latest acknowledged write (assuming no concurrent conflicts). The coordinator waits only for the fastest W or R responses, so slow replicas do not dominate latency.
These knobs trade consistency for latency. N=3, W=2, R=2 is a common balanced setting. R=1 and W=N favors fast reads; W=1 and R=N favors fast writes. W + R ≤ N gives lower latency but no overlap guarantee, so stale reads are possible. Consistency models range from strong (reads always see the latest write) through weak to eventual; Dynamo-style stores default to eventual consistency with quorums that make most reads fresh.
- N=3, W=1, R=3: fast writes, slower reads.
- N=3, W=3, R=1: fast reads, writes fail if any replica is down.
- N=3, W=2, R=2: overlap guaranteed, tolerates one slow replica.
v2 and returns the newest version.Inconsistency resolution with vector clocks
Because writes can land on different replicas during partitions, two versions of the same key may be created independently. Simple last-write-wins by wall clock silently discards data and depends on clock accuracy. A vector clock instead attaches to each version a list of [server, counter] pairs. When server Si writes, it increments its own counter. Version X is an ancestor of Y if every counter in X is less than or equal to the matching counter in Y; then Y simply supersedes X. If neither dominates, the versions are siblings in conflict.
Conflicts are returned to the client, which merges them using application knowledge; a shopping cart might take the union of items. The costs are real: clients carry merge logic, and the vector can grow with many writers, so systems cap its length and drop the oldest pairs, accepting that reconciliation may occasionally be imprecise.
Handling failures: gossip, hinted handoff and Merkle trees
A node should not be declared dead because one peer says so. All-to-all heartbeats work but scale poorly, so Dynamo-style systems use a gossip protocol: each node keeps a membership list with heartbeat counters, increments its own periodically, and sends the list to a few random peers. If a node's counter has not advanced for a while according to many peers, it is marked down.
For temporary failures, strict quorums would block writes, so the system uses a sloppy quorum: it picks the first W healthy servers on the ring, skipping unreachable ones. A stand-in node keeps the data with a hint and pushes it back when the owner recovers, a technique called hinted handoff. For permanent failures, replicas run anti-entropy using Merkle trees: keys are grouped into buckets, each bucket is hashed, and parent nodes hash their children. Two replicas compare root hashes and descend only into differing subtrees, so the data transferred is proportional to the difference, not the dataset.
- Data center outage: replicate across data centers so users can be served from another site.
Write and read paths
On a write, the node first appends the request to a commit log on disk for durability, then inserts it into an in-memory memtable. When the memtable reaches a threshold, it is flushed to disk as an immutable, sorted SSTable (sorted strings table). Writes are therefore sequential disk appends plus memory operations, which is why log-structured stores sustain very high write throughput. Background compaction later merges SSTables and discards overwritten or deleted values.
On a read, the node checks the memtable first. On a miss, it must consult SSTables from newest to oldest, and a Bloom filter per SSTable answers whether the key is definitely absent or possibly present, letting most files be skipped without disk I/O. The value found in the newest matching SSTable is returned to the client.
Key numbers
Key terms
- CAP theorem
- During a network partition a distributed system must sacrifice either consistency or availability.
- Quorum
- The minimum number of replicas (W for writes, R for reads) that must respond for an operation to succeed.
- Vector clock
- Per-version list of server counters used to tell whether two versions are ordered or concurrent.
- Gossip protocol
- Decentralized membership and failure detection where nodes periodically exchange state with random peers.
- Sloppy quorum
- A quorum formed from the first healthy nodes on the ring rather than the strict owners.
- Hinted handoff
- A stand-in node temporarily stores writes for a down replica and returns them when it recovers.
- Merkle tree
- A hash tree that lets two replicas find differing data by comparing hashes top-down.
- SSTable
- An immutable on-disk file of key-value pairs sorted by key, produced by flushing a memtable.
- Bloom filter
- A compact probabilistic set that answers definitely-not-present or maybe-present.
Common mistakes
- Stating CAP as pick any two, ignoring that partitions are not optional.
- Placing N replicas on N virtual nodes that share the same physical machine.
- Using wall-clock last-write-wins without acknowledging silent data loss.
- Declaring a node dead based on a single missed heartbeat from one observer.
- Repairing replicas by comparing full datasets instead of Merkle trees.
- Forgetting the commit log, so memtable contents vanish on a crash.
Further study
- Amazon Dynamo paper (2007)
- Cassandra: A Decentralized Structured Storage System (2009)
- Bigtable: A Distributed Storage System for Structured Data (OSDI 2006)
- Brewer's CAP conjecture and the Gilbert-Lynch proof (2002)
- The Log-Structured Merge-Tree (O'Neil et al., 1996)