Short answers to what readers ask most about this topic.
01How do you design a distributed key-value store like Dynamo?
Partition keys across nodes with consistent hashing, then copy each key to N replicas on the next nodes around the ring. Read and write through a quorum of R and W replicas, version values with vector clocks, and repair failures with hinted handoff and Merkle-tree anti-entropy. Nodes find each other through gossip, and each node stores data in a local engine such as an LSM tree.
02What do N, W and R mean in a quorum system?
N is the number of replicas that hold a key, W is how many must acknowledge a write, and R is how many must answer a read. When R plus W is greater than N, every read set overlaps every write set. The Dynamo paper reports (3, 2, 2) as a common configuration.
03What is the difference between vector clocks and last-write-wins?
Last-write-wins keeps the version with the larger timestamp, which is simple but silently drops one of two concurrent updates and trusts synchronised clocks. A vector clock stores a counter per coordinating node, so the store can tell whether one version descends from another or whether they conflict. Conflicting versions are returned together and the application merges them.
04What is hinted handoff, and how is it different from anti-entropy?
Hinted handoff handles short outages: when a replica is down, another node accepts its write with a hint and delivers it when the replica recovers. Anti-entropy handles longer divergence by comparing Merkle trees between replicas in the background and syncing only the keys that differ. Hints can be lost if the stand-in node dies, which is why both exist.
05How does an LSM tree store data on a single node?
A write is appended to a write-ahead log and inserted into an in-memory memtable. When the memtable fills, it is flushed as an immutable sorted file called an SSTable. Reads check the memtable and then SSTables from newest to oldest, deletes are written as tombstones, and compaction merges files so old versions disappear.
Design a Distributed Key-Value Store, Dynamo-Style
How to design a distributed key-value store the Dynamo way: consistent hashing, N/W/R quorums, vector clocks, hinted handoff, Merkle trees, gossip and an LSM node.
A Dynamo-style key-value store hashes each key onto a ring, copies it to N replicas, and answers once W replicas acknowledge a write or R replicas answer a read. Vector clocks expose conflicting writes, hinted handoff and Merkle-tree anti-entropy repair failures, gossip tracks membership, and each node stores data in an LSM tree.
I run Postgres, Redis and Docker on a single VPS for ERP and POS work, and one Postgres primary with a replica has always been enough. So I have never operated a Dynamo-style cluster, and I will not pretend otherwise. I wrote this post, and the two small TypeScript programs in it, to understand why the design looks the way it does.
The authority here is the 2007 Amazon Dynamo paper, the 1996 LSM-tree paper and the reference pages on vector clocks and Merkle trees, all linked at the end. Every number is either quoted from those sources or derived with the arithmetic shown, and the two code outputs were produced by running the code.
What is a Dynamo-style key-value store, and what does it give up?
It is a database with one shape of operation, get(key) and put(key, value), spread across many equal nodes with no leader. The Dynamo paper built it for services that must respond within 300 ms for 99.9 percent of requests, and it deliberately chose availability over strong consistency. Replicas can disagree for a while, and the system promises only eventual consistency.
The paper's own summary table is the best map of the design, because each problem gets exactly one technique. I use it as the outline of this post, one section per row.
Problem
Technique
What it buys
Partitioning
Consistent hashing
Incremental scalability
High availability for writes
Vector clocks, reconciled on reads
Version size is decoupled from update rates
Handling temporary failures
Sloppy quorum and hinted handoff
High availability and durability when some replicas are unreachable
Recovering from permanent failures
Anti-entropy using Merkle trees
Synchronises divergent replicas in the background
Membership and failure detection
Gossip-based membership protocol
Preserves symmetry and avoids a centralised registry
Storage on each node is a separate concern. The paper treats the local persistence engine as pluggable and lists Berkeley DB Transactional Data Store, BDB Java Edition, MySQL and an in-memory buffer with a persistent backing store. This post uses an LSM tree for that layer, because it is what most modern descendants use.
How does the store decide which node owns a key?
Hash the key onto a ring and walk clockwise. The first node is the coordinator, and the next N minus 1 nodes hold the other replicas. That ordered list is the preference list. Adding or removing a node only moves the keys next to it, which is the whole reason for the ring.
One detail matters for durability. The preference list must name N distinct physical machines. With virtual nodes the walk can land on two tokens of the same server, so the paper skips positions that belong to a machine already in the list. Three replicas on one disk are one replica with three names.
How do N, W and R quorums trade consistency for latency?
N is how many nodes hold a copy of a key. W is how many must acknowledge a write, and R is how many must answer a read. When R plus W is greater than N, every read set shares at least one node with every write set, so the read touches a replica that saw the latest acknowledged write. The paper notes that latency is set by the slowest of the R or W replicas, so both are usually below N, and the common configuration it reports is (3, 2, 2). The program below checks the overlap by brute force.
// Can a read set of R replicas miss every replica that a write set of W replicas touched?
const N = 3;
function subsets(size: number, from = 0): number[][] {
if (size === 0) return [[]];
const out: number[][] = [];
for (let i = from; i < N; i++) {
for (const rest of subsets(size - 1, i + 1)) out.push([i, ...rest]);
}
return out;
}
function everyReadSeesEveryWrite(W: number, R: number): boolean {
return subsets(W).every((w) => subsets(R).every((r) => r.some((x) => w.includes(x))));
}
for (const [W, R] of [[2, 2], [2, 1], [1, 1]]) {
console.log(`N=${N} W=${W} R=${R} R+W=${R + W} overlap guaranteed: ${everyReadSeesEveryWrite(W, R)}`);
}
// Output (run with tsx):
// N=3 W=2 R=2 R+W=4 overlap guaranteed: true
// N=3 W=2 R=1 R+W=3 overlap guaranteed: false
// N=3 W=1 R=1 R+W=2 overlap guaranteed: false
The arithmetic for N equal to 3: with W and R both 2, the sum is 4, which is greater than 3, and the program confirms that every pair of sets overlaps. With W equal to 2 and R equal to 1 the sum is 3, so a read can land on the one replica that a write skipped. The table turns the same arithmetic into a choice.
Config (N, W, R)
R plus W
Does a read see the last acknowledged write?
Failure behaviour with replicas down
(3, 2, 2)
4
Yes, overlap is guaranteed
Reads and writes each survive one replica down
(3, 1, 1)
2
No, stale reads are possible
Fastest; writes survive two replicas down
(3, 3, 1)
4
Yes, overlap is guaranteed
Reads survive two down, but a write fails if any one replica is down
(3, 1, 3)
4
Yes, overlap is guaranteed
Writes survive two down, but a read fails if any one replica is down
The overlap guarantee holds only for a strict quorum on the same N nodes. Dynamo uses a sloppy quorum: reads and writes go to the first N healthy nodes, not the first N on the ring. During a failure, R plus W greater than N no longer proves a read sees the last write.
Vector clocks or last-write-wins: how are conflicting writes resolved?
Last-write-wins attaches a timestamp and keeps the larger one. It is simple, it trusts clocks that drift between machines, and it silently discards one of two concurrent updates. The Dynamo paper names it as the policy a data store falls back on when the store itself resolves conflicts. A vector clock instead keeps a counter per coordinating node, stored as (node, counter) pairs, so two versions can be compared causally.
If every counter in one clock is at least the matching counter in the other, that version descends from the other and the older one can be dropped. If neither descends from the other, the versions are concurrent. The store keeps both, returns both on read, and leaves the merge to the application. The paper's example is a shopping cart merged by union, because an add to cart must never be forgotten.
type Clock = Record<string, number>; // coordinator node -> counter
type Order = "equal" | "descends" | "ancestor" | "concurrent";
function compare(a: Clock, b: Clock): Order {
const nodes = new Set([...Object.keys(a), ...Object.keys(b)]);
let aAhead = false;
let bAhead = false;
for (const n of nodes) {
if ((a[n] ?? 0) > (b[n] ?? 0)) aAhead = true;
if ((b[n] ?? 0) > (a[n] ?? 0)) bAhead = true;
}
if (aAhead && bAhead) return "concurrent"; // neither saw the other: keep BOTH versions
if (aAhead) return "descends"; // a is newer, b can be dropped
if (bAhead) return "ancestor";
return "equal";
}
const v1: Clock = { a: 1 }; // ticket written via node a
const v2: Clock = { a: 2 }; // updated again via a
const v3: Clock = { a: 2, b: 1 }; // one client updates v2 via node b
const v4: Clock = { a: 2, c: 1 }; // another client updates v2 via node c, during a partition
console.log("v2 vs v1:", compare(v2, v1));
console.log("v3 vs v2:", compare(v3, v2));
console.log("v3 vs v4:", compare(v3, v4));
// Output:
// v2 vs v1: descends
// v3 vs v2: descends
// v3 vs v4: concurrent <- a real conflict; return both, let the application merge
The cost is growth. The paper limits clock size by storing a timestamp beside each pair and, once the pair count reaches a threshold (it says 10 as an example), removing the oldest pair. It admits this can lead to inefficiencies in reconciliation, because some ancestry can no longer be proven.
How does the store survive a down replica: hinted handoff and anti-entropy?
With a sloppy quorum, the coordinator writes to the first N healthy nodes. If replica A is down, its copy goes to the next node D with a hint naming A in the metadata. D keeps hinted replicas in a separate local store and delivers them to A once A recovers, then deletes its own copy. Writes stay available, and the paper notes this works best when membership churn is low.
Hints are lost if D dies before A returns, so replicas also compare data in the background. Each node keeps a Merkle tree per key range: leaves are hashes of individual keys and every parent is a hash of its children. Two replicas exchange the root hash. If it matches, the range is identical, and if not they descend only into the differing subtrees. A range of 1,024 keys has a tree 10 levels deep, since 2 to the power of 10 is 1,024, so locating one differing key takes about 10 levels times 2 child hashes, roughly 20 comparisons instead of 1,024.
The paper lists the cost of Merkle trees: when a node joins or leaves, key ranges change and the affected trees must be recalculated. Plan repair around that, and do not treat anti-entropy as free just because the compare step is cheap.
How do nodes know who is alive without a central registry?
Gossip. In the paper, each node contacts a peer chosen at random every second, and the two nodes reconcile their persisted membership change histories. Every node therefore ends up with an eventually consistent view of who is in the ring. Failure detection stays local too: a node that cannot reach a peer treats it as failed and routes around it, with no coordinator to ask.
A worked estimate shows why this scales. If every informed node tells one new node per round, the informed count at best doubles each round, so 1,000 nodes need about 10 rounds, since 2 to the power of 10 is 1,024. At one exchange per second that is on the order of 10 seconds. Real spread is slower than the doubling ideal, so read it as a lower bound on convergence time, not a measurement.
What does one node do with a write: memtable, WAL and SSTables?
The LSM-tree paper's idea is to absorb writes in memory and merge them to disk in sequential batches. In node-level terms: append the write to a write-ahead log (WAL), insert it into the in-memory memtable, and when the memtable is full flush it as an immutable sorted file, an SSTable. A read checks the memtable, then tables from newest to oldest. A delete is a tombstone, and compaction merges tables. This sketch holds all of that in about 60 lines.
import { appendFileSync, existsSync, readFileSync, writeFileSync, rmSync } from "node:fs";
const TOMBSTONE = "\u0000deleted";
const FLUSH_AT = 3; // entries in the memtable before it becomes an SSTable
class TinyLsm {
private mem = new Map<string, string>();
private sstables: Array<Array<[string, string]>> = []; // newest first, each sorted by key
constructor(private walPath: string) {
if (existsSync(walPath)) {
for (const line of readFileSync(walPath, "utf8").split("\n").filter(Boolean)) {
const [k, v] = JSON.parse(line) as [string, string];
this.mem.set(k, v); // replay: the WAL rebuilds the memtable after a crash
}
}
}
put(key: string, value: string): void {
appendFileSync(this.walPath, JSON.stringify([key, value]) + "\n"); // log first, then memory
this.mem.set(key, value);
if (this.mem.size >= FLUSH_AT) this.flush();
}
delete(key: string): void { this.put(key, TOMBSTONE); } // a delete is just a newer write
get(key: string): string | undefined {
const hit = this.mem.get(key) ?? this.fromTables(key);
return hit === TOMBSTONE ? undefined : hit;
}
private fromTables(key: string): string | undefined {
for (const t of this.sstables) { // newest table wins
const row = t.find(([k]) => k === key);
if (row) return row[1];
}
return undefined;
}
private flush(): void {
this.sstables.unshift([...this.mem.entries()].sort(([a], [b]) => (a < b ? -1 : 1)));
this.mem.clear();
writeFileSync(this.walPath, ""); // data now lives in a table, so the log can be truncated
}
compact(): void { // merge every table: newest value per key, tombstones dropped
const merged = new Map<string, string>();
for (const t of [...this.sstables].reverse()) for (const [k, v] of t) merged.set(k, v);
this.sstables = [[...merged.entries()]
.filter(([, v]) => v !== TOMBSTONE)
.sort(([a], [b]) => (a < b ? -1 : 1))];
}
stats() { return { memtable: this.mem.size, sstables: this.sstables.map((t) => t.length) }; }
}
const wal = "/tmp/tiny-lsm.wal";
rmSync(wal, { force: true });
const db = new TinyLsm(wal);
db.put("wash:101", "basic");
db.put("wash:102", "premium");
db.put("wash:103", "basic"); // third entry: memtable flushes to SSTable 1
db.put("wash:101", "wax"); // the newer version lives in the memtable
db.delete("wash:102"); // a tombstone, not an in-place delete
console.log("after writes ", db.stats(), db.get("wash:101"), db.get("wash:102"));
console.log("after restart ", new TinyLsm(wal).stats(), new TinyLsm(wal).get("wash:101"));
db.put("wash:104", "basic"); // second flush: two SSTables
console.log("before compact", db.stats());
db.compact();
console.log("after compact ", db.stats(), db.get("wash:101"), db.get("wash:102"));
// Output:
// after writes { memtable: 2, sstables: [ 3 ] } wax undefined
// after restart { memtable: 2, sstables: [] } wax
// before compact { memtable: 0, sstables: [ 3, 3 ] }
// after compact { memtable: 0, sstables: [ 3 ] } wax undefined
Read the output as evidence. After the writes, the memtable holds 2 entries (the newer wash:101 and the tombstone for wash:102) and one SSTable holds 3. A fresh instance built from the WAL alone still answers wax for wash:101, proving the log replays. The WAL contained only 2 entries because it was truncated at the flush. Compaction then merged two 3-row tables into one 3-row table: the old wash:101 lost to the new one and the tombstoned wash:102 vanished.
The sketch is honest about what it skips. It keeps SSTables in memory rather than files, never calls fsync, and has no checksums, index or Bloom filter. Those are exactly the parts a real engine spends its code on, which is the argument for using an existing one.
Should you build one, and what is the decision checklist?
For a single-VPS ERP workload, almost certainly not. A Dynamo-style store pays for write availability with conflict handling and repair work that a primary-plus-replica Postgres never asks of you. Before reaching for one, walk this list in order.
Name the access pattern. If you only look rows up by key, continue. If you need joins or multi-row transactions, stay relational.
Decide what a conflict means. If concurrent versions can be merged, like a cart union, vector clocks fit. If exactly one writer must win, choose leader-based replication.
Pick N, W and R from a failure budget. (3, 2, 2) tolerates one replica down for both reads and writes, and the sum 4 exceeds 3.
Plan repair before launch: hinted handoff for short outages, Merkle anti-entropy for long ones, and a tombstone retention rule that outlasts your longest outage.
Prefer a mature engine for the node-level storage. The paper itself ran Dynamo on existing storage engines rather than writing its own.
The rule I take from this: every box in a Dynamo-style design exists because the one before it chose availability, and each choice leaves a bill. Replicas diverge, so you need versions. Nodes fail, so you need hints and repair. Writing it down first tells you whether you are willing to pay.