6.4 Leaderless Replication
In some implementations the client sends writes to several replicas directly; in others a coordinator node does this on the client's behalf.
Quorum overlap
- 6
- w + r
- 5
- n
- 1
- overlap
w + r > n — at least 1 replica is in both sets, so any read contacts a node that saw the latest write. Note this bounds staleness, not concurrency: the read still has to pick the newer version.
Read repair
- 2 of 3
- quorum
- 1
- stale
Replica 3 is stale at version 6. The read still returns the correct value, because the client takes the greatest version even if only one node returned it — but the staleness persists until something fixes it. In leaderless systems this is hard to monitor, unlike leader-based replication where lag is a single number.
Abandon the leader entirely — any replica directly accepts writes from clients.
History: some of the earliest replicated data systems were leaderless, but the idea was mostly forgotten during the era of relational-database dominance. It became fashionable again after Amazon used it for its in-house Dynamo system in 2007. Riak, Cassandra, ScyllaDB are open source Dynamo-inspired datastores — hence "Dynamo-style."
⚠️ The original Dynamo was described in a paper but never released outside Amazon. The similarly named DynamoDB has a COMPLETELY DIFFERENT architecture: single-leader replication based on the Multi-Paxos consensus algorithm.
In some implementations the client sends writes to several replicas directly; in others a coordinator node does this on the client's behalf. Unlike a leader, that coordinator does NOT enforce a particular ordering of writes — and this difference has profound consequences.
4.1 Writing when a node is down
3 replicas, one down for a reboot. NO FAILOVER — all replicas are equal.
user 1234 WRITE ──┬──▶ [replica 1] ok ─┐
├──▶ [replica 2] ok ─┼─ 2 of 3 acknowledged ⇒ SUCCESS
└──▶ [replica 3] ✗ ┘ (client simply IGNORES the missed replica)
(down)
… replica 3 comes back, MISSING the write …
user 2345 READ ───┬──▶ [replica 1] → version 7 ─┐
├──▶ [replica 2] → version 7 ─┼─ client takes the value with the
└──▶ [replica 3] → version 6 ─┘ GREATEST VERSION/TIMESTAMP
(even if returned by only ONE node)
│
└──▶ READ REPAIR: client writes version 7 back to replica 3Read requests are also sent to SEVERAL NODES IN PARALLEL. Every value written must be tagged with a version number or timestamp so the client can tell which responses are up to date.
4.2 Three catch-up mechanisms
| Mechanism | How it works | Coverage |
|---|---|---|
| Read repair | A client reading from several nodes in parallel detects stale responses and writes the newer value back to the stale replica | Works well for values that are READ OFTEN |
| Hinted handoff | If a replica is unavailable, another replica stores writes on its behalf as HINTS. When the intended replica returns, the hint-storing replica sends them over and deletes the hints | Covers values that are NEVER READ, which read repair misses |
| Anti-entropy | A background process periodically looks for differences between replicas and copies missing data | Does NOT copy writes in any particular order, and there may be a SIGNIFICANT DELAY before data is copied |
4.3 Quorums
With n replicas: every write must be confirmed by w nodes; every read must query at least r nodes.
As long as w + r > n, we expect to get an up-to-date value when reading, because at least one of the r nodes we read from must be up to date.
Common choice: n odd (3 or 5), w = r = (n+1)/2 rounded up. But you can vary them: a workload with few writes and many reads may set w = n, r = 1 — faster reads, but ONE failed node causes ALL writes to fail.
Tolerance:
| Config | Tolerates |
|---|---|
| w < n | writes continue if a node is unavailable |
| r < n | reads continue if a node is unavailable |
| n=3, w=2, r=2 | 1 unavailable node |
| n=5, w=3, r=3 | 2 unavailable nodes |
Mechanics: reads and writes are normally sent to ALL n replicas in parallel; w and r determine how many you WAIT FOR. If fewer than w or r are available, the operation returns an error. A node could be unavailable for many reasons — crashed, powered down, disk full, network interruption — and we care only whether it returned a successful response.
Note: there may be more than n nodes in the cluster, but any given value is stored on only n nodes — which is what allows sharding (Ch 7).
Quorums need not be majorities. It matters only that the read and write sets OVERLAP in at least one node. Majorities are common because r = w = majority ensures w + r > n while tolerating up to ⌊n/2⌋ failures, but other assignments allow flexibility in algorithm design.
Deliberately violating the quorum condition (w + r ≤ n): reads and writes still go to n nodes, but fewer successes are required.
- ✗ More likely to read stale values
- ✔ Lower latency (particularly beneficial with synchronous replication)
- ✔ More highly available — during a network interruption there's a higher chance you can continue. The database becomes unavailable only when reachable replicas fall below w or r
4.4 The six ways quorums lie to you
Although quorums APPEAR to guarantee that a read returns the latest written value, in practice it is not so simple.
- A node carrying a new value fails and is restored from a replica carrying an old value, so the number of replicas holding the new value can fall below
w, breaking the quorum condition. - During rebalancing (Ch 7), nodes may disagree about which nodes hold the
nreplicas for a value, so read and write quorums may no longer overlap. - A read concurrent with a write may or may not see the written value. In particular, one read may see the new value and a subsequent read the old one.
- A write that succeeded on some replicas but fewer than
woverall is not rolled back where it succeeded, so a write reported as failed may or may not be returned by subsequent reads. - With real-time clock timestamps (Cassandra, ScyllaDB), writes may be silently dropped if another node with a faster clock wrote the same key.
- Two concurrent writes may be processed in different orders on different replicas — a conflict, exactly as in multi-leader replication.
Dynamo-style databases are generally optimized for use cases that can tolerate eventual consistency. The parameters w and r let you ADJUST THE PROBABILITY of stale reads — but it's wise NOT TO TAKE THEM AS ABSOLUTE GUARANTEES.
4.5 Monitoring staleness
Leader-based: easy. Writes are applied to leader and followers IN THE SAME ORDER, and each node has a position in the replication log. Subtract the follower's position from the leader's ⇒ replication lag, exposed as a metric.
Leaderless: hard. There is NO FIXED ORDER in which writes are applied. The number of hints a replica stores for handoff can be one measure of system health, but it's difficult to interpret usefully.
Eventual consistency is a deliberately vague guarantee, but for OPERABILITY it's important to be able to QUANTIFY "eventual."
4.6 Single-leader vs leaderless performance
Reading from the leader ensures up-to-date responses but has three performance problems:
- Read throughput is limited by the leader's capacity
- On leader failure you must wait for detection and failover. Even a quick failover is noticed by users as increased response times; a long one means the system is unavailable for its duration
- Very sensitive to performance problems on the leader — if the leader is slow (overload, resource contention), increased response times immediately affect users
The leaderless advantage:
Because there is no failover, and requests go to multiple replicas in parallel anyway, one replica becoming slow or unavailable has very little impact on response times — the client simply uses the responses from the faster replicas. Using the fastest responses is called REQUEST HEDGING, and it can significantly reduce tail latency.
At its core, the resilience of a leaderless system comes from the fact that IT DOESN'T DISTINGUISH BETWEEN THE NORMAL CASE AND THE FAILURE CASE.
This is especially helpful for GRAY FAILURES — a node that isn't completely down but is running in a degraded state, unusually slow to handle requests — or a node that's simply overloaded (e.g. recovery via hinted handoff can cause a lot of additional load). A leader-based system has to DECIDE whether the situation is bad enough to warrant a failover (which can itself cause further disruption); in a leaderless system that question doesn't even arise.
But leaderless has its own performance problems:
- One replica must still detect when another is unavailable to store hints, and the handoff must send them — putting additional load on the replicas AT A TIME WHEN THE SYSTEM IS ALREADY UNDER STRAIN
- The more replicas, the bigger the quorums and the more responses to wait for. Even waiting only for the fastest r or w, a bigger r or w raises the chance of hitting a slow replica, increasing overall response time. In practice, quorums are seldom more than 4 of 7 or 5 of 9 nodes
- A large-scale network interruption disconnecting a client from many replicas can make it impossible to form a quorum
Sloppy quorums are the escape hatch: allow any reachable replica to accept writes even if it's not one of the usual n replicas for that key. (Riak/Dynamo: "sloppy quorum"; Cassandra/ScyllaDB: consistency level ANY.) There is NO GUARANTEE that subsequent reads will see the written value, but depending on the application it may still be better than having the write fail.
The three-way summary:
| Resilience to network interruption | Staleness risk | |
|---|---|---|
| Multi-leader | Greatest — reads and writes need only ONE leader, which can be co-located with the client | Reads can be ARBITRARILY out of date |
| Leaderless (quorum) | Good | Good fault tolerance AND a high likelihood of reading up-to-date data — a compromise |
| Single-leader | Weakest | Strongest consistency available |
4.7 Multi-region with leaderless
Cassandra / ScyllaDB: the client picks a node in its local region — the coordinator node — and sends the write there. The coordinator forwards to all replicas in its OWN region AND TO ONE REPLICA IN EVERY OTHER REGION, which then forwards to the other replicas in that region. This optimization avoids making the cross-region request multiple times.
Consistency levels determine how many responses are required: a quorum across all regions, a separate quorum in each region, or a quorum only in the client's local region. A LOCAL quorum avoids waiting for slow cross-region requests but is more likely to return stale results.
Riak keeps all client↔node communication local to one region (so n describes replicas within one region); cross-region replication happens asynchronously in the background, in a style similar to multi-leader replication.