In 2007 a team at Amazon published Dynamo: Amazon's Highly Available Key-value Store at SOSP. The paper did not invent any single technique it used. Consistent hashing, quorums, vector clocks, Merkle trees and gossip were all known. What it showed was that they could be composed into one production system with a clear goal: a write must succeed even while disks fail, networks split and data centres go dark, and the application must accept that it will sometimes read more than one version of a value.

That trade became the template for a generation of leaderless stores. This article reads it as a system. It starts from the requirements, maps each to the technique chosen, traces one write and one read through all of them, works through a network partition with real version clocks, and ends with what the descendants kept, what they dropped, and how to decide whether a Dynamo-style design fits your problem today.

Advertisement

The requirements that forced the design

Dynamo served services such as the shopping cart and session state, where the business cost of rejecting a write was higher than the cost of occasionally merging two versions. Operations are a simple get and put on a key, with values usually under a megabyte and no operations spanning keys. Writes must always be accepted. The system runs inside a trusted environment, so there is no Byzantine fault handling. And performance is specified at the tail: services agreed latency targets at the 99.9th percentile, with the paper's example being a response within 300 ms for 99.9 percent of requests at a peak of 500 requests per second.

Those last two requirements are what make the design unusual. An always-writable store cannot block on a single leader or on a majority that happens to be unreachable, so it gives up linearizability. A 99.9th-percentile target means every component must avoid waiting on the slowest replica, so requests complete when enough replicas answer, not when all do.

Six problems, six techniques

ProblemTechniqueWhy this one
Spread keys over nodes, add nodes one at a timeConsistent hashing with virtual nodesAdding a node moves only its neighbours' keys; virtual nodes even out load and let big machines own more
Accept writes during failuresVector clocks; reconcile on readConcurrent versions are detected, not silently overwritten
Handle temporary failuresSloppy quorum and hinted handoffWrites go to the next healthy nodes, which return the data later
Recover from permanent failuresAnti-entropy with Merkle treesReplicas compare hash trees and transfer only differing ranges
Know who is in the cluster and who is downGossip membership, local failure detectionNo central coordinator to lose
Meet the tail-latency targetTunable N, R, W; coordinator returns on first W or R repliesThe slowest replica is never on the critical path

Each row is covered in depth on its own page: vector clocks, Merkle trees, gossip membership and quorum reads and writes. This article is about how they fit together, which is where Dynamo's behaviour comes from.

Advertisement

The ring and the preference list

ABCDEFGHkey k = MD5(user:7)walk clockwise from k:B, C, D are the top N=3D is down, so E is nextCoordinator Bfirst healthy node in preference listReplicas B, Chome replicas that answeredD unreachablefailure detector marks it, not membershipE stores a hinted copyhint says: belongs to DW=2 met by B and C (or E)client gets successD returnsE hands the hint back, deletes it
A key hashes to a point on the ring. The first N distinct nodes clockwise form its preference list. When D is unreachable, the write goes to E with a hint naming D, and E hands the data back when D recovers.

Every key is hashed with MD5 to a 128-bit position on a ring. Each physical node owns many positions, called virtual nodes or tokens, so its load is the sum of many small arcs rather than one large one. A key's preference list is found by walking clockwise from its position and collecting the first N distinct physical nodes, skipping extra tokens of a node already chosen. The list is kept longer than N so a request can move down it past unreachable nodes.

Any node can receive a request. If it is not in the key's top N, it forwards to one that is. The node that handles the request is its coordinator.

One write and one read through every layer

The coordinator logic is short, and it is worth reading because every other mechanism hangs off it.

# Coordinator logic as the paper describes it, simplified. N=3, R=2, W=2.
def put(key, value, context):              # context = vector clock from an earlier get
    prefs = preference_list(key)           # N healthy nodes, skipping down ones,
                                           # extended past position N if needed
    coord = prefs[0]
    clock = context.increment(coord.id)    # coordinator bumps its own entry
    acks = 0
    for node in prefs[:N]:
        hint = home_replica_for(node, key) # set only if node is a stand-in
        send_async(node, "store", key, value, clock, hint)
    for reply in replies(timeout_ms=SLA_BUDGET):
        acks += 1
        if acks >= W:
            return OK                      # remaining replicas finish in background
    return FAIL                            # client may retry via another coordinator

def get(key):
    prefs = preference_list(key)
    versions = []
    for reply in gather(prefs[:N], "read", key, need=R, timeout_ms=SLA_BUDGET):
        versions.extend(reply.versions)
    frontier = drop_dominated(versions)    # keep only clocks nobody descends from
    read_repair_async(prefs[:N], frontier) # push newest versions to stale replicas
    return frontier, merge_clocks(frontier)  # app reconciles if len(frontier) > 1

A successful write promises only that W of the N chosen nodes stored this version. With the common configuration of N=3, R=2 and W=2, R+W exceeds N, so in a stable cluster any read quorum overlaps any write quorum. But the chosen nodes are the first N healthy nodes, not the fixed home replicas. During a failure the write set and a later read set can be different groups of nodes with no overlap at all. That is the meaning of sloppy quorum: availability is preserved by giving up the guarantee that R+W greater than N makes reads see the latest write.

Reads return every version that is not dominated by another. When the versions are concurrent, the client receives all of them plus a merged context, and the next put with that context supersedes them. Version detection relies on vector clocks:

def descends(a, b):            # True if clock a has seen everything clock b has
    return all(a.get(n, 0) >= c for n, c in b.items())

def drop_dominated(versions):
    keep = []
    for v in versions:
        if not any(o is not v and descends(o.clock, v.clock) and o.clock != v.clock
                   for o in versions):
            keep.append(v)
    return dedupe_by_clock(keep)

# Example from the worked partition below:
#   [A:2]        vs [A:1, B:1]  -> neither descends the other -> two siblings
#   [A:3, B:1]   vs [A:2]       -> first descends second       -> one version

Clocks grow when many coordinators touch one key, so the paper truncates the oldest entry past a threshold (10 is its example). Truncation can make a descendant look concurrent and produce a spurious sibling.

Worked example: a partition and a merge

Take key cart:42 with preference list A, B, C and N=3, W=2, R=2. Client 1 puts version v1 through A. A writes it locally, B and C acknowledge, and the clock is [A:1].

Now the network splits. A is on one side; B and C are on the other. Client 2, on the B side, reads [A:1] and puts v2 with that context. B coordinates: it increments its own entry to [A:1, B:1], stores it, and C acknowledges; W=2 is met. Client 3, on the A side, also reads [A:1] and puts v3 through A. A cannot reach B or C, so it walks down its preference list to D and E, which store hinted copies meant for B and C. The clock is [A:2]. Both writes succeeded. Neither side waited for the partition to heal.

When the partition heals, D and E deliver their hints to B and C, and B and C now each hold two versions: [A:1, B:1] and [A:2]. Neither descends the other. The next read returns both as siblings with the merged context [A:2, B:1]. The application merges them; the shopping cart article shows why a naive merge resurrects deleted items and how to avoid it. The client then puts the merged value with context [A:2, B:1], coordinated by A, producing [A:3, B:1], which dominates both siblings, and read repair spreads it.

How often does this happen? In the paper's 24-hour measurement of the cart service, about 99.94 percent of requests saw exactly one version. The authors attribute most divergence to concurrent writers, often automated clients, rather than failures. Rare is not never: your merge code runs in production whether or not your tests exercise it.

Permanent failures, membership and partitioning

Hinted handoff handles nodes that come back. For nodes that do not, or for hints lost because the stand-in itself failed, each node keeps a Merkle tree per key range it hosts. Two replicas compare roots; if they match, the range is identical; if not, they descend only into differing subtrees and transfer only differing keys. This anti-entropy process bounds how long replicas can stay divergent.

Membership is changed explicitly by an administrator, and the change is spread by gossip, with each node contacting a random peer periodically to reconcile its view. The paper separates this from failure detection: a node that cannot reach a peer treats it as down locally and routes around it, without removing it from the ring. So a flapping link causes hinted writes, not a rebalancing storm.

The paper's partitioning section is often skipped and is one of its most practical parts. It describes three strategies. First, random tokens per node with ranges following token boundaries, which made adding a node expensive because ranges had to be rescanned and Merkle trees rebuilt. Second, random tokens but equal-sized partitions. Third, Q equal partitions with each of S nodes holding Q/S of them, which balanced load best and made bootstrapping and backup simpler, because a partition can be copied whole. Fixed range boundaries are easier to move, repair and archive than random ones.

What the descendants kept and dropped

SystemKept from DynamoChanged or dropped
Apache CassandraToken ring with virtual nodes, gossip, hinted handoff, tunable consistency per request, read repair, Merkle-tree repairData model and LSM storage follow Bigtable; concurrent writes resolved by last-write-wins timestamps per cell, not vector clocks; Cassandra 4.0 lowered the default token count per node from 256 to 16
RiakClosest to the paper: ring, sloppy quorums, hinted handoff, causal version tracking with siblings, active anti-entropyMoved to dotted version vectors to limit sibling explosion; added CRDT data types so many merges are automatic
Project VoldemortConsistent hashing, vector clocks, client-side routing, pluggable storage enginesSpecialised for read-only bulk-loaded stores alongside the read-write path
Amazon DynamoDBThe name and the key-value API stylePer its 2022 USENIX ATC paper, each partition is a replication group using Multi-Paxos; the leader serves writes and strongly consistent reads. It is not a leaderless Dynamo

The parts that survived are the ones that are mechanical and invisible to developers: ring placement, gossip, hinted handoff, repair. The part most often dropped is client-side reconciliation, because asking every application to write a correct merge function proved to be the expensive part. Systems either accepted last-write-wins and its silent loss of concurrent updates, or replaced hand-written merges with CRDTs. DynamoDB went further and returned to leaders per partition, buying strong consistency and simpler semantics with a managed failover path. If you are choosing a managed store, when to pick DynamoDB covers that decision.

Failure modes the paper teaches you to expect

  • Reading your own write fails. Under a sloppy quorum, a write acknowledged by stand-ins may not be seen by a read that reaches the home replicas. Do not promise read-your-writes unless you route both requests through the same coordinator or use strict quorums.
  • Sibling explosion. A hot key written by many uncoordinated clients accumulates concurrent versions faster than they are merged. Cap sibling counts, alert on them, and fix the write pattern.
  • Deletes come back. A delete is a write of a tombstone. If tombstones are garbage-collected before every replica has seen them, an old replica's value reappears during repair. Keep tombstones longer than your maximum repair interval.
  • Repair never finishes. Anti-entropy that is throttled too hard, or skipped on busy clusters, lets divergence accumulate until a node failure exposes it. Schedule repair and measure that each range completes within its window.
  • Last-write-wins with skewed clocks. Descendants that replaced vector clocks with timestamps silently lose updates when wall clocks disagree. Use client-supplied timestamps only if you control the clients' clocks.

When a Dynamo-style design fits today

Your requirementDynamo-style leaderless storeLeader-per-partition store
Writes must succeed during partitionsYes, by designNo; minority side rejects writes
Read-your-writes and conditional updatesHard; needs extra machineryNatural
Values with a meaningful merge (sets, counters)Good fit, especially with CRDTsPossible but merge is unnecessary
Multi-region active-active writesNatural fitNeeds a separate replication layer
Developers unfamiliar with concurrent versionsRiskySafer

Most applications do not need minority-side writes and pay a large semantic cost for them. Pick a leaderless design when the business case is the one Dynamo had: accepting the write is worth more than reading it consistently, and the data has a merge you can write down.

What to do next

  1. Read the paper itself, especially its partitioning and lessons sections.
  2. For each key family in your system, write down whether a rejected write or a merged write is worse. That answer, not fashion, should pick leaderless or leader-based storage.
  3. If you run a Dynamo descendant, find out how it resolves concurrent writes (vector clocks, dotted version vectors or timestamps) and test that path with two writers during an injected partition.
  4. Set tombstone lifetime longer than your longest repair cycle and alert when repair falls behind.
  5. Measure sibling counts or timestamp conflicts per key in production; a steady rise usually means a write pattern bug, not bad luck.
Key takeaway: Dynamo's contribution was composition: consistent hashing with virtual nodes for placement, sloppy quorums and hinted handoff for availability, vector clocks for detecting concurrent writes, Merkle trees for repair and gossip for membership, all tuned for 99.9th-percentile latency. Its descendants kept the mechanical parts and mostly dropped client-side reconciliation in favour of last-write-wins or CRDTs, while DynamoDB returned to leaders. Use the design when an accepted write is worth more than a consistent read and the data has a merge you can write down.