Learn Labs
6. Replication

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
replica 1write
replica 2write
replica 3both
replica 4read
replica 5read
Safe

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.

Writes go to the leftmost w replicas, reads contact the rightmost r. Overlap is what guarantees a read sees the latest write — it is a counting argument, not a coordination protocol.

Read repair

replica 3
mechanism
run it
2 of 3
quorum
1
stale
v7v6clientreplica 1version 7replica 2version 7replica 3version 6
Problem

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.

With no leader there is no failover — a replica that misses a write is simply behind, and something has to notice and fix it. There are two mechanisms and they cover different cases.

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 3
Figure 6.4.14.1 Writing when a node is down

Read 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

MechanismHow it worksCoverage
Read repairA client reading from several nodes in parallel detects stale responses and writes the newer value back to the stale replicaWorks well for values that are READ OFTEN
Hinted handoffIf 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 hintsCovers values that are NEVER READ, which read repair misses
Anti-entropyA background process periodically looks for differences between replicas and copies missing dataDoes 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.

12345write set · w = 3read set · r = 3n = 5, w + r = 6 > 5 — the sets must overlap in at least one node⇒ at least one read replica has seen the latest write ⇒ tolerates 2 unavailable nodes
Figure 6.4.24.3 Quorums

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:

ConfigTolerates
w < nwrites continue if a node is unavailable
r < nreads continue if a node is unavailable
n=3, w=2, r=21 unavailable node
n=5, w=3, r=32 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.

  1. 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.
  2. During rebalancing (Ch 7), nodes may disagree about which nodes hold the n replicas for a value, so read and write quorums may no longer overlap.
  3. 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.
  4. A write that succeeded on some replicas but fewer than w overall is not rolled back where it succeeded, so a write reported as failed may or may not be returned by subsequent reads.
  5. With real-time clock timestamps (Cassandra, ScyllaDB), writes may be silently dropped if another node with a faster clock wrote the same key.
  6. 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:

  1. Read throughput is limited by the leader's capacity
  2. 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
  3. 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:

  1. 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
  2. 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
  3. 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 interruptionStaleness risk
Multi-leaderGreatest — reads and writes need only ONE leader, which can be co-located with the clientReads can be ARBITRARILY out of date
Leaderless (quorum)GoodGood fault tolerance AND a high likelihood of reading up-to-date data — a compromise
Single-leaderWeakestStrongest 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.


On this page