KnowSys

Replication & Consistency Models

Follow one name change, from "Old Name" to "New Name", across the copies of a database: why a reload can still show the old name, what a system can promise a reader, and what each promise costs in Postgres, Kafka, Dynamo-style stores and CRDTs.

⏱ 45 min read◆ BeginnerAssumes: Docker for the first experiment; chapters 21 (Postgres) and 23 (Kafka) help but aren't required
Start reading

You change the display name on your profile from "Old Name" to "New Name" and press Save. The page tells you it worked. You reload, and the old name is still there. A moment later you reload again and the new name appears, as if nothing had happened.

Nothing was broken. The database behind a profile page usually keeps several copies of the row on different servers, so that it survives a failure and can answer more reads than one server could. When you pressed Save, one copy changed straight away and the others heard about it afterwards, a bit like head office changing a price and posting the new sheet to each branch. Your reload happened to ask a branch that hadn't received the sheet yet.

So the useful question is what the system promised you: after a write, which copies may a reader be sent to, and what is that reader guaranteed to see? Such a promise is called a consistency model, and the promises run from "always the latest" down to "eventually". This chapter reproduces your stale reload on two real Postgres servers, puts names on the promises a system can make, and then follows the same name change through four real designs to see what each promise costs.

01One row, two copies

1.1How a copy follows the original

A single server holding your profile row has two weaknesses. If it dies, the row is gone until someone restores a backup, and if a million people read profiles at once, one machine has to answer all of them. The usual remedy is a second server holding a copy of the same data. We'll call the original the primary, because it's the one that accepts changes, and the copy a replica. (Postgres calls it a standby, and we'll use both words.) Reads can be sent to the replica, which takes load off the primary, and if the primary dies the replica can take its place.

A replica needs a way to follow the primary. Copying the whole table after every change would be far too slow, so databases reuse a log they already keep. Before the primary changes a row, it appends a short description of the change to a file called the write-ahead log, or WAL. Record 7, say, reads "set name to New Name in row 1 of profile". (Chapter 08's journal did the same job for a filesystem.)

The primary streams this log to the replica, and the replica does two things with each record. First it saves the record in its own copy of the log and flushes it, meaning it asks the disk to store the bytes permanently so they survive a power cut (the fsync call from chapter 08). Later it replays the record, meaning it repeats the change on its own copy of the table. Queries on the replica read the table and not the log, so until a record has been replayed, a query on the replica can't see it.

That gap between "the record has arrived" and "the record has been replayed" is where your stale reload comes from. Here is your Save followed through both servers. When the primary has made a change final and told the client, we say the change is committed.

Your Save, and two reloads
Primarytakes the writesYoubrowserReplicaa copy that follows the WALrow · nameOld Namerow · nameOld Nameyour pageOld NameWAL record 7name = New NameWAL record 7flushed
Step 1. Before the Save, the primary and the replica both hold the row with Old Name, and so does the page you're looking at.
1 / 7

On a quiet system the lag is short, and it grows when the replica is busy or the network is slow. It never quite reaches zero, because the primary commits first and tells the replica afterwards. Since the primary doesn't wait for the replica, this arrangement is called asynchronous replication. To see the effect on real servers we'll stretch the lag until we can watch it.

1.2Trying it: update on the primary, read from the replica

You need Docker. The script below starts two Postgres containers on one network (docker network create makes the network, and docker run -d starts a container in the background). The first is the primary. It starts with wal_level=replica, which makes it keep in its WAL the detail a replica needs, and max_wal_senders=5, which allows up to five streaming connections to it. The pg_hba.conf line lets a replica connect with a password, and pg_reload_conf() makes the server re-read its configuration so the line takes effect. (docker exec runs a command inside a running container, and the until … pg_isready loop waits for the server to start accepting connections.) The second container is the replica. pg_basebackup copies the primary's data files to give the replica a starting point, and its -R option writes the settings that tell it to keep following the primary afterwards. We start the replica with hot_standby=on, so that it answers queries while it follows, and with recovery_min_apply_delay=3s, which makes it wait three seconds after a record arrives before replaying it. A real replica lags by much less, but it lags the same way, and three seconds makes the gap easy to see. The last lines create the row, change it on the primary, and read it from the replica once a second. (psql -qAt prints only the value, with no headers or borders.)

Predict before you read on

The UPDATE has finished on the primary. You now read the row from the replica once a second for five seconds. What do you see?

Start a primary and a delayed replica, update a row on the primary, and read the replica once a second
shell
Shell
# primary, and a replica that takes a base backup from it and applies WAL 3 s late
docker network create pgnet
docker run -d --name pgp --network pgnet -e POSTGRES_PASSWORD=x postgres:16-alpine \
  -c wal_level=replica -c max_wal_senders=5 -c hot_standby=on
until docker exec pgp pg_isready -U postgres; do sleep 1; done; sleep 2
docker exec pgp sh -c 'echo "host replication all all scram-sha-256" >> /var/lib/postgresql/data/pg_hba.conf'
docker exec pgp psql -U postgres -c "select pg_reload_conf()"
docker run -d --name pgr --network pgnet -e PGPASSWORD=x --entrypoint sleep postgres:16-alpine infinity
docker exec pgr sh -c 'mkdir -p /var/lib/postgresql/data2 && chown postgres /var/lib/postgresql/data2 && chmod 700 /var/lib/postgresql/data2'
docker exec -u postgres -e PGPASSWORD=x pgr pg_basebackup -h pgp -U postgres -D /var/lib/postgresql/data2 -R -X stream
docker exec -u postgres pgr pg_ctl -D /var/lib/postgresql/data2 \
  -o "-c recovery_min_apply_delay=3s -c hot_standby=on -p 5433" -l /tmp/pgr.log start
 
# then: create a row, change it on the primary, and poll the replica
docker exec pgp psql -U postgres -c "CREATE TABLE profile (id int PRIMARY KEY, name text); INSERT INTO profile VALUES (1,'Old Name')"
sleep 5                                          # let the replica replay the setup
docker exec pgp psql -U postgres -c "UPDATE profile SET name='New Name' WHERE id=1"
for i in 1 2 3 4 5; do docker exec pgr psql -U postgres -p 5433 -qAt -c "select name from profile where id=1"; sleep 1; done
output
C++
Old Name
Old Name
Old Name
New Name
New Name

The setup commands print container ids and status lines, and only the poll results are shown. The update committed on the primary at time zero. The five reads happened about 0.1, 1.2, 2.4, 3.5 and 4.6 seconds after it. The first three returned Old Name and the last two returned New Name, so the replica switched somewhere between 2.4 and 3.5 seconds, once its three-second delayed replay caught up. That matches the Scene: the replica had the record long before it replayed it, and for those three seconds every read saw the old row.

1.3Is that a bug?

Not by itself. The primary's answer was correct the whole time, and so was the replica's, for the state it had reached. The problem arises when the application reads from the replica right after writing and assumes it will see its own change. A user who saves their profile, reloads the page and sees the old name has hit exactly that gap.

Whether the gap is acceptable depends on what's being read. A profile page that shows the old name for a few milliseconds is fine. A page that says "this seat is still free" for a few milliseconds is not. So we need a precise way to say how far behind a reader may be, and before that, a fair question: why does anyone accept the gap at all? Why not make the primary wait until the replica has caught up before it says "committed"?

02Why copies fall behind

2.1Four reasons to keep copies

Making the primary wait would close the gap, and section 4 shows what that costs. First, though, it's worth being clear about what the second copy buys. A single server has no gap at all, since every read sees every write that finished before it and there's only one place to look. Replication gives that up on purpose. There are four different reasons to keep more than one copy, and each asks something different of the copies.

ReasonWhat it needs from the copiesExample
DurabilityThe write exists somewhere else before you say "done"A second copy in another zone that the write waits for
AvailabilityAnother copy can take over when one diesPatroni for Postgres or Sentinel for Redis (tools that promote a copy when the original dies)
Read scalingMany copies can answer readsRead replicas behind a connection pooler (a proxy that shares database connections among clients)
LatencyA copy close to the userDynamoDB global tables (Amazon's database, copied to several regions), a CDN (copies of your files at sites near users)

Those rows pull in different directions. Durability wants the copies updated before the write is acknowledged. Latency and availability want each copy to answer without waiting for the others. Every design in this chapter picks a point between those two.

?Why not just copy every write everywhere before replying?

Because then every write waits for the slowest copy, and a write can't finish at all while any copy is unreachable. You'd have built a system that's less available than one server. The whole subject is about waiting for some of the copies, and being precise about what readers can see as a result.

2.2Three ways to arrange the copies

The next decision is which copies may accept a write. Three arrangements are common. Before the table, two words. A leader is a node that accepts writes (the primary of section 1 is one), and the nodes that copy it are its followers. A conflict is two writes to the same key made without either one knowing about the other.

TopologyWho accepts writesHow conflicts ariseReal systems
Single leaderOne node; others replay its logThey don't, until failoverPostgres, MySQL, Kafka partitions, Redis
Multi-leaderSeveral nodes, each replicating to the othersTwo leaders accept writes to the same keyCouchDB, Postgres with bidirectional logical replication (replication of row changes), DynamoDB global tables (default mode)
LeaderlessAny replica; the client or a coordinator writes to severalConcurrent writes land in different orders on different replicasDynamo, Cassandra, Riak

Failover, in the first row, is what happens when the leader dies and a follower is promoted to take its place.

Single leader is the default for a reason: one node orders all writes, so there are no conflicts to resolve. The other two exist because a single leader is one place, in one region, and it has to be up. We'll follow single leader in section 4 and the other two in sections 5 and 6. Whichever arrangement you pick, though, copies disagree for some time after a write, and the next section gives us the words to say exactly how long, and for whom.

03What a reader can be promised

Picture three people reading your profile just after you pressed Save. You reload it yourself. A friend, whom you phoned the moment Save finished, opens it on their own computer. And a stranger on another continent opens it a second later. Each of them could be given a different guarantee, and the consistency model from the opening is a contract about exactly that: which values a read may return. A model is defined by what a client can observe from outside and not by how the system is built, so you can test a system against a model without looking inside. Jepsen does exactly that.

We'll go from the strongest promise to the weakest, and at each step ask what the system gives up in return for dropping a property.

3.1Linearizability: one copy, one moment

The strongest promise we could ask for is that the system behaves as though there were only one copy of your profile, whatever it does underneath. A system is linearizable if every operation appears to take effect at one instant somewhere between the moment it was called and the moment it returned. Herlihy and Wing defined it in 1990.

What you'll feel is recency. Once your Save has returned, or once any reader has seen the new name, every read that starts later must see that name or a newer one. That holds across clients, with no sessions or sticky routing involved.

?Why does "starts later" matter so much?

Because it ties the order of operations to real time. When your Save returns and you phone your friend, who then reloads, your friend must see the new name. Weaker models allow your friend's read to go to a replica that hasn't heard of the change yet. Linearizability is the only model on this list where the phone call is safe.

That's why the things that need a single truth are built on it: leader election, distributed locks, unique usernames, and "is this seat still free". The price comes in section 3.5, and it's paid on every request.

3.2Sequential and causal consistency

Linearizability is expensive, so the next step is to drop one property and see what we keep. The first property to drop is the tie to real time. Sequential consistency (Lamport, 1979) keeps a single order of operations that all clients agree on, and keeps each client's own operations in the order that client issued them. But a read may return an old value. Your friend, after the phone call, might still see Old Name, as long as every client sees the operations in the same order.

Causal consistency drops the single order as well. It only promises that if one operation could have influenced another, everyone sees them in that order. You read a comment and then reply to it, so nobody may see your reply before the comment. Two unrelated writes, say your name change and a stranger's comment elsewhere, can appear in different orders to different readers.

?Why would anyone want causal instead of linearizable?

Because causal consistency can stay available during a network partition, and linearizability can't. A network partition is a failure in which some nodes can't reach others, and a system is available when every request to a live node still gets a real answer. A replica cut off from the others can keep serving causally consistent reads and writes, because it only has to respect the causality it already knows about. Mahajan, Alvisi and Dahlin proved that a slightly stronger variant is the strongest model that can.

3.3Eventual consistency and session guarantees

At the weak end, eventual consistency promises only this: if writes stop, all replicas will converge on the same value. It says nothing about what you read in the meantime. That's all the stranger on another continent needs: at some point they should see the new name. Werner Vogels' Eventually Consistent is the classic description of the model as Amazon used it.

That's too weak for the stale reload, so a middle layer exists. Between eventual and causal sit four session guarantees, from Terry et al.'s Session Guarantees for Weakly Consistent Replicated Data (1994), written for Bayou, a replicated database for mobile devices built at Xerox PARC. They're promises to one client, which makes them cheap to give.

GuaranteeWhat one client is promisedWhat breaks without it
Read your writesYou see your own writesThe profile you just saved shows the old version
Monotonic readsYou never see time go backwardsA comment appears, then disappears on reload
Monotonic writesYour writes apply in the order you made them"Delete draft" is applied before "create draft"
Writes follow readsA write you make after reading X is ordered after XA reply is stored before the post it replies to

Read your writes and monotonic reads cover the stale reload we started with and the flicker between two states. You can often give them to the one user who would notice, without paying for linearizability for everyone.

3.4The models side by side

ModelReal-time orderOne global orderAvailable in a partitionTypical mechanism
LinearizableYesYesNoSingle leader serving the reads, Raft, Paxos
SequentialNoYesNoOne agreed order for all writes, then reads from a local copy
CausalNoNo, only causal pairsYesTracking which writes each write depended on
Session guaranteesNoNoYes, if the client sticks to reachable replicasSession tokens, sticky routing, waiting for a log position
EventualNoNoYesAsynchronous replication, background repair

Raft and Paxos are consensus algorithms, which let several nodes agree on one order of operations even when some fail. Chapter 27 builds them. Sticky routing means always sending one user's reads to the same replica, so that user can't see time go backwards, and a session token is a note the client carries saying how far into the log its own writes and reads have got, so that a replica can wait until it has caught up to that position before answering. Section 4.4 builds that wait on Postgres.

Predict before you read on

A Postgres primary with one async read replica. A user saves a setting (on the primary), then the next page load reads it from the replica and sees the old value. Which guarantee was violated?

3.5What CAP says, and what it leaves out

Gilbert and Lynch's proof of the CAP theorem is narrower than its reputation. It says a system can't be both linearizable and available while the network is partitioned. A node that's cut off from the others must either refuse to answer, because it can't be sure it has the latest value, or answer and risk being stale.

A Venn diagram of three overlapping circles labelled Consistency, Availability and Partition Tolerance, with the overlaps labelled CA, CP and AP
The picture CAP usually gets: three properties, pick any two. The proof is narrower than that. A distributed system doesn't get to opt out of partitions, so the CA overlap isn't a real choice for it. The choice the theorem describes is what a cut-off node does while the partition lasts: refuse to answer (CP) or answer with what it has (AP).Image: JamieMcCarthy, CC BY-SA 4.0, via Wikimedia Commons

That's all it says. It doesn't cover latency, it doesn't cover the weaker models, and it says nothing about the system when there's no partition. Abadi's PACELC fills that gap: if there's a Partition, choose Availability or Consistency; Else, choose Latency or Consistency.

That's the vocabulary. The next step is to see what one real system lets you buy with it. The Postgres pair from section 1 has to answer a single question on every write, and the answer is where the price is set: when does the primary say "committed"?

04One leader: Postgres and Kafka

4.1Five answers to “when is it committed?”

In single-leader replication the primary writes each change to its WAL, as in section 1, and streams the log to its standbys. The choice is how long the primary waits for them before it tells the client. Postgres makes that choice per transaction with synchronous_commit. A synchronous standby is one the primary waits for, and the ones to wait for are named in synchronous_standby_names. The setting has five values, and once that list names at least one standby they differ in how much of the standby's work the primary waits for.

synchronous_commitThe commit returns afterSurvives primary loss?Standby reads see it when the commit returns?
offNothing; a background process flushes the WAL within 3 × wal_writer_delay (200 ms by default)No, even locallyNo
localLocal WAL flushOnly if the primary's disk survivesNo
remote_writeStandby has written WAL to its OS (not fsynced)Unless the standby's OS crashes tooNo
on (default)Standby has flushed WAL to diskYesNo
remote_applyStandby has replayed itYesYes

The last column is the one people miss. The default, on, makes the write durable on the standby, meaning it would survive a crash. It doesn't make the write visible there, because replay is a separate step that runs after the flush. That's the same gap we watched in the Scene in section 1.

Step through one commit with a synchronous standby. The log is numbered by position, and a position in the WAL is called a log sequence number, or LSN. Three Postgres processes take part. The backend serves your connection, the walsender runs on the primary and streams the log, and the walreceiver is its counterpart on the standby.

One COMMIT with synchronous_commit = on
ClientPrimaryStandbyCOMMITWAL up to LSN Xwrite + fsyncflush = XCOMMIT OKreplay
Step 1. The client commits. The primary writes a commit record to its WAL and flushes it locally, as it would with no standby at all.
1 / 6

Every level below remote_apply leaves a window in which the client has heard "committed" and the standby still shows the old row. The next question is what it costs to shrink that window.

4.2What each level costs

The cost is latency, the time one commit takes. We can measure it with pgbench, Postgres's own benchmarking tool. With one client running a single-row UPDATE in autocommit mode (each statement is its own transaction), each run lasted 4 seconds, and seven rounds were interleaved across the five modes. The primary and the standby ran on one host, talking over loopback (the machine's internal network interface, so there's no real network in the way). The ranges are wide and overlap, so read the medians as an ordering and not as exact values.

synchronous_commitMedian latencyRange over 7 runsWhat it waited for
off0.33 ms0.08–0.61 msNothing
local1.74 ms0.66–5.09 msLocal fsync
remote_write2.92 ms0.97–6.36 msLocal fsync + standby write()
on3.17 ms1.60–7.37 msLocal fsync + standby fsync
remote_apply3.74 ms1.86–9.47 ms… + standby replay

On one host the network round trip costs roughly nothing, so what you're seeing is mostly the second fsync and the replay. Across availability zones (separate data centres in one region), add the network round trip to every row from remote_write down. AWS documents inter-AZ latency as single-digit milliseconds, and that goes on top of every commit.

?Why is off safe to use at all?

Because it can lose recent commits, but it can't corrupt anything. The docs say a crash leaves the database "just the same as if those transactions had been aborted cleanly". For a click counter or a log of page views, losing the last few hundred milliseconds in a crash is often a fine trade for a commit about five times faster than local in the table above. You can set it per transaction, so the important writes can keep on.

Waiting for a standby buys durability, and it creates a new way to get stuck. What happens when the standby you're waiting for isn't there?

4.3Trying it: the escape hatch in synchronous replication

Stop the standby, and every commit on the primary waits for it, with no time limit. A CREATE TABLE sent to the primary while its synchronous standby is stopped shows wait_event = SyncRep in the list of sessions, and 80 seconds later it's still there, waiting for a standby that isn't coming back. An operator looking at a hung system will try to cancel the stuck statement, so the question is what the primary does when the wait is cancelled. This is the code that handles it:

src/backend/replication/syncrep.c
postgres/postgres @ REL_16_4 ↗
C
		/*
		 * It's unclear what to do if a query cancel interrupt arrives.  We
		 * can't actually abort at this point, but ignoring the interrupt
		 * altogether is not helpful, so we just terminate the wait with a
		 * suitable warning.
		 */
		if (QueryCancelPending)
		{
			QueryCancelPending = false;
			ereport(WARNING,
					(errmsg("canceling wait for synchronous replication due to user request"),
					 errdetail("The transaction has already committed locally, but might not have been replicated to the standby.")));
			SyncRepCancelWait();
			break;
		}

The comment is candid. The commit record is already in the primary's WAL by the time the backend starts waiting, so it can't be rolled back, and a cancel just stops the waiting. The client gets a warning and a success. We can watch it happen with the containers from section 1, which should still be running. The block below makes every commit wait for a standby (ALTER SYSTEM SET synchronous_standby_names = '*' names any standby, and pg_reload_conf() applies it), creates a table while the replica is still up so that the CREATE TABLE can finish, and then stops the replica with pg_ctl stop. The INSERT runs in the background (the trailing &). Three seconds later the first query looks in pg_stat_activity, Postgres's list of current sessions, to show what the insert is waiting on, and the second sends that session a cancel with pg_cancel_backend. The last line counts the row on the primary.

Cancel a commit that's waiting for a stopped synchronous standby
shell
Shell
# make every commit on the primary wait for a standby, then stop the replica
docker exec pgp psql -U postgres -c "ALTER SYSTEM SET synchronous_standby_names = '*'"
docker exec pgp psql -U postgres -c "select pg_reload_conf()"
docker exec pgp psql -U postgres -c "CREATE TABLE orders (id int)"
docker exec -u postgres pgr pg_ctl -D /var/lib/postgresql/data2 stop
 
# this insert commits on the primary, then waits for a standby that is gone
docker exec pgp psql -U postgres -c "INSERT INTO orders VALUES (42)" &
sleep 3
docker exec pgp psql -U postgres -Atc "select wait_event from pg_stat_activity where query like 'INSERT%'"
docker exec pgp psql -U postgres -Atc "select pg_cancel_backend(pid) from pg_stat_activity where wait_event = 'SyncRep'"
wait
docker exec pgp psql -U postgres -Atc "select count(*) from orders where id = 42"
output
Output
SyncRep
t
INSERT 0 1
WARNING:  canceling wait for synchronous replication due to user request
DETAIL:  The transaction has already committed locally, but might not have been replicated to the standby.
1

The setup commands print their own status lines (ALTER SYSTEM, CREATE TABLE and so on), and only the output of the last four commands is shown. SyncRep is the wait event after three seconds, the backend sitting in the wait we saw in the diagram in 4.1. The t is pg_cancel_backend reporting that it delivered the cancel. Then the insert returned INSERT 0 1, a success, together with the warning that most drivers never surface. (The warning goes to the error stream and the INSERT 0 1 line to standard output, so your terminal may show them in either order.) The final 1 says the row is committed on the primary, where it exists and nowhere else.

Keeping one named standby and accepting that writes block has its own problem, because any restart of that standby becomes a write outage. The usual answer is to let any one of several standbys do the job. With synchronous_standby_names = 'ANY 1 (s1, s2)', a commit finishes once either of two standbys confirms, so losing one blocks nothing. FIRST 1 (s1, s2) instead prefers s1 and falls back to s2, in that order.

4.4Reading your own writes from a replica

Back to the reload that started the chapter. How often would it happen with each setting, and how do we stop it? To find out, a client updated one row on the primary and then immediately read it from the standby, 2,000 times in a row. A read counts as stale, as in section 1, if it returns anything other than the value just written.

Primary's synchronous_commitStale reads (3 runs of 2,000)Median
local (standby is async)1,924 · 1,975 · 1,93196.6%
on (standby flushed)97 · 63 · 383.2%
remote_apply (standby replayed)0 · 0 · 00%

With the standby asynchronous, almost every immediate read is stale, because the read arrives long before the replica replays the record. With remote_apply none is, because the commit doesn't return until replay is done. The middle row is the surprising one.

?Why does on still give stale reads?

Because, as the table in 4.1 said, a flushed record isn't a replayed record. The window between the standby's fsync and its replay is short on an idle server, and it caught roughly 3% of reads here. Under a heavy write load, or while replay waits on a conflicting query, that window grows.

There are three standard fixes, from cheapest to most expensive.

FixHow it worksCost
Pin the session to the primary for a whileAfter a write, send that user's reads to the primary for N secondsPrimary load; N is a guess at lag
Wait for an LSNReturn pg_current_wal_insert_lsn() with the write; before reading, the replica waits until pg_last_wal_replay_lsn() passes itA token to carry and a short wait on the replica
remote_applyEvery commit waits for replay on the synchronous standbysEvery write pays replay latency; only helps reads on those standbys

The LSN approach gives read your writes to one session, built from two functions: the write hands back its position in the log, and the reader refuses to answer until its replica has replayed past that position. It also degrades well. If the replica is far behind, the reader can give up waiting and go to the primary.

Postgres counts the standbys you name. A system built around replication from the start has to decide who counts without a list.

4.5Kafka's version: in-sync replicas

Suppose the profile service also publishes every name change to Kafka, so that other services can hear about it. Kafka stores messages in partitions, each of which is an ordered log kept by several brokers (Kafka servers). One broker is the partition's leader and the others are its followers. A producer, the program sending messages, chooses how long to wait with the setting acks: acks=0 doesn't wait at all, acks=1 waits for the leader alone, and acks=all waits for the followers that are keeping up.

Kafka doesn't count a fixed set of standbys. It keeps a dynamic in-sync replica set, or ISR: the followers that have recently caught up with the leader. With acks=all, a produce waits for every member of the current ISR, and min.insync.replicas sets how small the ISR may shrink before writes are refused. Here is the check, in Kafka's source:

core/src/main/scala/kafka/cluster/Partition.scala
apache/kafka @ 3.7.0 ↗
scala
          val minIsr = effectiveMinIsr(leaderLog)
          val inSyncSize = partitionState.isr.size
 
          // Avoid writing to leader if there are not enough insync replicas to make it safe
          if (inSyncSize < minIsr && requiredAcks == -1) {
            throw new NotEnoughReplicasException(s"The size of the current ISR ${partitionState.isr} " +
              s"is insufficient to satisfy the min.isr requirement of $minIsr for partition $topicPartition")
          }

Note the requiredAcks == -1 condition (-1 is how acks=all is represented): the check applies only to acks=all. A producer with acks=1 writes happily to a leader whose ISR is just itself. The Kafka chapter follows the high watermark (the point up to which messages count as committed) and leader epochs (numbers that identify each leader's term) through a failover.

Postgres and Kafka share one limit. Every write passes through a single leader, which is one machine in one place that has to be up. Dynamo's designers asked what a store would look like if no machine had that job.

05Replication without a leader

Amazon's Dynamo paper (2007) removed the leader. Any replica accepts a write, and the client, or a coordinator acting for it, sends each operation to several replicas and waits for some of them. Cassandra, Riak and ScyllaDB work this way. With no leader there's no single place that knows the newest version of your name, so the client has to ask several replicas and rely on the replicas it asks overlapping with the ones that took the write.

5.1N, W and R

Three numbers describe a leaderless configuration:

  • N: how many replicas hold each key.
  • W: how many must acknowledge a write before it succeeds.
  • R: how many must answer a read before it returns.

Take N = 3 replicas of your profile, with W = 2 and R = 2. A write reached at least two of the three, and a read asks at least two of the three. Two groups of two chosen from three always share at least one replica, so a read always reaches at least one replica that has the latest completed write, and can return it by picking the value with the highest version. In general, if W + R > N, every read set overlaps every write set in at least one replica, and groups of replicas big enough to guarantee that are called quorums. Dynamo's paper gives (3, 2, 2) as the common configuration. A node that handles a request on the client's behalf, sending it to the replicas and collecting their answers, is called the coordinator.

N = 3, W = 2, R = 2 (Cassandra QUORUM both ways)2 + 2 > 3overlap ≥ 1
Replicas that can be down for writesN − W1
Replicas that can be down for readsN − R1
N = 3, W = 1, R = 1 (Cassandra ONE)1 + 1 ≤ 3no overlap
N = 5, W = 3, R = 33 + 3 > 5survives 2 down
Quorum rule of thumb for N replicasW = R = ⌊N/2⌋ + 1

Cassandra lets each request name how many replicas it waits for, and calls that number the request's consistency level. Its QUORUM is exactly that ⌊RF/2⌋ + 1 (RF is its name for N, the replication factor), ONE waits for a single replica, and LOCAL_QUORUM counts only replicas in the coordinator's datacenter, so it doesn't pay a cross-region round trip.

Quorums also tolerate slow nodes well, because the client waits for the fastest W or R replies and ignores the rest. One replica in a garbage-collection pause (a stretch in which its runtime freezes the process to reclaim memory) adds no latency at all to a W = 2 write with N = 3. A single leader can't do that: if the leader is slow, every write is slow.

5.2A quorum read

Step through a read at R = 2 when one replica missed the latest write of your new name.

Quorum read, N = 3, R = 2, one stale replica
ClientCoordinatorReplica AReplica BGET profile 1data readdigest readNew Name (ts 200)digestNew Name
Step 1. The client asks any node. That node becomes the coordinator for this request.
1 / 6

The overlap rule promises that a read sees the latest completed write. The word completed is carrying weight, and the next subsection shows what happens when a write is still in flight.

5.3Why W + R > N isn't linearizable

Is a system with W + R > N linearizable? Let's test it on the profile. A writer is setting the name from Old Name (version 1) to New Name (version 2) with W = 2, and the messages to two of the three replicas are slow. Two readers come along while it's in flight.

Two quorum reads during one slow write
Replica AReplica BReplica CReaderswhat each one was toldnameOld Name · v1nameOld Name · v1nameOld Name · v1reader 1saw New Namereader 2saw Old Name
Step 1. Three replicas hold Old Name, version 1. Writes use W = 2 and reads use R = 2, so any read and any *completed* write share a replica.
1 / 6

Nothing in the overlap rule covers a write in flight, and that gap is enough to break linearizability.

?Why does making readers write back fix it?

Because it makes a reader finish the write it saw. Suppose reader 1, before returning New Name, writes it back to a quorum. Then reader 2's quorum must overlap at least one replica that holds New Name, and reader 2 can't be told the old value. That's the ABD algorithm (Attiya, Bar-Noy and Dolev, 1995). Cassandra does the same thing and calls it read repair. Its default is blocking read repair (read_repair = 'BLOCKING' since Cassandra 4.0): when the replies disagree, the coordinator writes the newest value to the stale replica before answering the client. The Cassandra docs call the resulting property monotonic quorum reads, and setting read_repair = 'NONE' gives it up.

Even ABD only gives a linearizable register, a single value that's read and overwritten. It can't do compare-and-set (change the value only if it still equals what I last read), because two writers can each pass a quorum with different values. For that, Cassandra has lightweight transactions (LWTs), which run Paxos. The consensus chapter covers why that step needs consensus and not just overlap.

5.4Sloppy quorums, hinted handoff and anti-entropy

Quorums assume the key's N home replicas are reachable. What should a write do when some aren't? Dynamo's answer was to keep accepting it. Nodes sit on a circle called the ring, each key lives on N consecutive nodes of it, and the next healthy node along the ring can stand in for a home replica that's down. The stand-in keeps the data together with a hint saying whose it is, and hands it over when the home replica returns. This is hinted handoff, and a write acknowledged partly by stand-ins is a sloppy quorum. (Chapter 29 explains how keys are placed on the ring.)

Eight nodes n1 to n8 on a ring; an arc from token t1 covers n2, n3 and n4 for key foo, and an arc from t2 covers n3, n4 and n5 for key bar
Placement on the ring, as Cassandra draws it, with N = 3. A key's hash lands between two tokens, and its replicas are the next three nodes clockwise: the hash of foo goes to n2, n3 and n4, and the hash of bar to n3, n4 and n5. Those are the key's home replicas. If n3 were down, the next node along, n5 for foo, is the natural stand-in.Image: Apache Cassandra documentation, Apache License 2.0
A sloppy quorum and a hint
Home AHome BHome CStand-in Dnext on the ringnameOld Name · v1nameOld Name · v1nameOld Name · v1hint for BNew Name · v2
Step 1. The profile key lives on home replicas A, B and C, all holding Old Name. Node D is the next one along the ring and holds nothing for this key. Writes use W = 2.
1 / 6

A sloppy quorum breaks the overlap argument. The W acknowledgements came from stand-in nodes, and the R replies come from the home nodes, so the two groups needn't share a replica. A write can succeed and a following quorum read can miss it entirely until the hints are delivered. Sloppy quorums buy write availability and give up the W + R > N argument to get it.

Sequence diagram: a client writes x=1 through a coordinator to three replicas; replica 2 is restarting, replicas 1 and 3 acknowledge, the coordinator acknowledges the client, stores a hint for replica 2, and replays the write when replica 2 reports alive
Cassandra's version of the hint, from its documentation. Replica 2 is restarting when the write arrives. Replicas 1 and 3 acknowledge, which is already a quorum of the home replicas, so the client gets its answer at t1. The coordinator keeps the hint itself and replays the write once replica 2 is seen alive at t3. Here the hint doesn't count towards the quorum; it only makes replica 2 catch up sooner than repair would.Image: Apache Cassandra documentation, Apache License 2.0

Replicas that missed writes and never got hints are fixed by anti-entropy: a background comparison of Merkle trees (trees of hashes), so that two replicas can find their differences by exchanging hashes and not the data itself. In Cassandra that's nodetool repair.

5.5Quorums in a storage engine: Aurora

Quorum arithmetic isn't only for key-value stores. Amazon Aurora (SIGMOD 2017) is a relational database that splits its storage into small pieces and stores six copies of each piece across three availability zones, with a write quorum of 4 and a read quorum of 3. What happens as copies are lost?

FailureCopies lostWrites?Reads?
One node1Yes (5 ≥ 4)Yes
One whole AZ2Yes (4 ≥ 4)Yes
One AZ plus one more node3No (3 below 4)Yes (3 ≥ 3), so it can rebuild

4 + 3 > 6, so reads overlap writes. And 4 + 4 > 6, so two writes overlap too, which stops two writers from both succeeding on disjoint copies. The paper's argument for six copies is exactly the last row: "AZ+1" failures must not lose data.

Leaderless stores and quorums give every replica the power to accept a write. That raises a question that a single leader never had to answer: what if two replicas accept different names for the same user at the same time?

06Two writes at once

6.1Last writer wins, and what it throws away

Say you rename yourself to "Name A" from your phone, connected to a replica in Frankfurt, while a laptop you left logged in renames you to "Name B" through a replica in Oregon. Neither replica has heard of the other write. When they exchange them, something has to decide what your name is afterwards.

The simplest rule is to stamp every write with the time it was made and keep the write with the highest timestamp. This is last-writer-wins, and Cassandra applies it to every column. Last-writer-wins converges, but it silently discards every concurrent write except one. And "last" is decided by wall clocks on different machines, so a node whose clock runs a few seconds fast wins every conflict for as long as it's fast, even against writes made later in real time.

?Why not just use better clocks?

Because two writes that happen at "the same time" on different replicas aren't ordered, whatever the clocks say. Neither saw the other. Picking one by timestamp is a policy and not a fact. (Chapter 26 looks at what clocks can and can't tell you.) What you need first is a way to detect that two writes were concurrent.

6.2Detecting concurrency with version vectors

A version vector keeps one counter per replica. Each replica increments its own entry when it accepts a write, and a write carries the vector it was based on. Compare two vectors entry by entry:

ComparisonMeaningWhat to do
Every entry of A ≤ BA happened before BKeep B
Every entry of B ≤ AB happened before AKeep A
NeitherConcurrentKeep both as siblings, or merge

Try it on "Name A" and "Name B". Both writes started from the same version, {Frankfurt: 0, Oregon: 0}. Frankfurt accepts "Name A" and stamps it {Frankfurt: 1, Oregon: 0}, and Oregon accepts "Name B" and stamps it {Frankfurt: 0, Oregon: 1}. Each vector has one entry bigger than the other's, so neither came before the other, and the store knows the writes were concurrent. It keeps both values, which are called siblings. Dynamo returned the siblings to the application, which merged them. For Amazon's shopping cart, the service Dynamo was built for, the merge was a union of the items. The paper is candid about the consequence: "deleted items can resurface", because a union can't tell a removed item from one the other replica never had.

Merging by hand in every application is a lot to ask. What if the data type merged itself?

6.3CRDTs: data types that merge themselves

A conflict-free replicated data type, or CRDT, is a data structure whose merge function needs no application help and always converges. Shapiro, Preguiça, Baquero and Zawirski formalised them in 2011.

The design requirement is that merge must be commutative, associative and idempotent: merging in a different order, in different groupings, or twice gives the same result. Then it doesn't matter in which order replicas exchange state, how often, or whether messages are duplicated, because they all reach the same value. (Mathematically, the states form a join-semilattice and merge is the join.)

The simplest example is a grow-only counter. Suppose your profile counts page views in two data centres, A and B. Each replica keeps a vector with one slot per replica, and only ever increments its own slot.

A G-counter across a partition
Replica Aowns slot 1NetworkReplica Bowns slot 2A plain number insteadlast writer winsviews[0, 0] · value 0views[0, 0] · value 0views3 (or 2)
Step 1. Both replicas start with the vector [0, 0]. A may only increment the first slot and B the second. The value is the sum of the slots.
1 / 6

Per-entry max is idempotent, since merging twice changes nothing, and commutative, since the order doesn't matter, and those two properties are what make it safe to merge any number of times in any order. Subtraction breaks them, so a counter that can go down is built from two G-counters, one for increments and one for decrements (a PN-counter).

6.4Trying it: an OR-set keeps the list honest

Dynamo's cart problem has a CRDT answer. Suppose your profile has a list of interests, kept on two replicas. The observed-remove set, or OR-set, gives each add a unique tag, and a remove deletes only the tags it has seen. If someone adds "chess" at the same moment someone else removes it, the add carries a tag the remove never saw, so it survives. An add wins over a concurrent remove, and a remove that wasn't concurrent with anything sticks.

The script below builds a G-counter and an OR-set from a few lines each. merge is the only way state travels between replicas. It first counts views on a and b and merges them in both directions and then once more, to show repeated merges change nothing. Then replica x removes "chess" while replica y adds it again, and both merge. The last lines apply last-writer-wins to the same two writes, with a timestamp that puts the remove just ahead because one clock runs a little fast.

G-counter and OR-set merges, against last-writer-wins
python
Python
import uuid
 
class GCounter:
    def __init__(self, me): self.me, self.p = me, {}
    def inc(self, n=1): self.p[self.me] = self.p.get(self.me, 0) + n
    def value(self): return sum(self.p.values())
    def merge(self, o):
        for r, n in o.p.items(): self.p[r] = max(self.p.get(r, 0), n)
 
class ORSet:
    def __init__(self): self.adds, self.removes = set(), set()   # (element, tag)
    def add(self, e): self.adds.add((e, uuid.uuid4().hex[:6]))
    def remove(self, e): self.removes |= {t for t in self.adds if t[0] == e}
    def value(self): return {e for e, t in self.adds - self.removes}
    def merge(self, o): self.adds |= o.adds; self.removes |= o.removes
 
# profile views counted in two regions, each only bumping its own entry
a, b = GCounter("a"), GCounter("b")
a.inc(3); b.inc(2)
a.merge(b); b.merge(a); a.merge(b)          # repeated merges change nothing
print("views a =", a.value(), " views b =", b.value())
 
# the profile's list of interests, edited on two replicas at the same time
x, y = ORSet(), ORSet()
x.add("chess"); y.merge(x)                  # both replicas list chess
x.remove("chess")                           # replica x: the user removes it
y.add("chess")                              # replica y, concurrently: it is added again
x.merge(y); y.merge(x)
print("interests x =", x.value(), " interests y =", y.value())
 
# last-writer-wins, given the same two writes with timestamps
lww = max([(1001, "remove chess"), (1000, "add chess again")])
print("last-writer-wins keeps only:", lww[1])
output
Output
views a = 5  views b = 5
interests x = {'chess'}  interests y = {'chess'}
last-writer-wins keeps only: remove chess

Both counters report 5, the sum of 3 and 2, which no plain number with last-writer-wins would reach. Both replicas end with "chess" in the list, because the concurrent add carried a tag the remove had never seen. Last-writer-wins, given the same two writes and a clock that's one tick ahead on the remover, keeps only the remove, and the add is lost without any error.

The removes set only ever grows in this naive version, because a replica can't know when every other replica has seen a remove. Production CRDTs spend much of their complexity on garbage-collecting that metadata safely.

You'll meet CRDTs in Riak's data types, Redis Enterprise's Active-Active databases, and the collaborative-editing libraries Automerge and Yjs.

We now have a single leader, quorums and CRDTs, each paying for a different promise. The practical question is which to reach for.

07Choosing a consistency level

7.1Match the requirement to the mechanism

Start from what the feature needs a reader to be able to rely on, find the weakest model that gives it, and pay only for that.

You needWeakest model that does itMechanism to reach for
Users see their own editsRead your writesPin to primary after writes, or LSN tokens
A feed that never jumps backwardsMonotonic readsSticky replica per session
Replies after the posts they answerCausalVersion vectors, or one leader per conversation
A uniqueness check, a lock, a leaderLinearizablePrimary-only reads, Raft/etcd, Paxos LWTs
No acknowledged write lost if a node dies(Durability)Sync standby, acks=all + min.insync.replicas, W ≥ 2
A like-counter across regionsEventual, mergedCRDT counter

?Why not just make everything linearizable to be safe?

Because you pay for it on every request. Linearizable reads have to reach the leader, or a quorum, or wait out a lease (a time-limited right to act as the leader), and in a multi-region deployment that's a cross-region round trip per read. And the price in DynamoDB is literal: eventually consistent reads are half the cost of strongly consistent ones, and reads from global secondary indexes (GSIs) can only be eventually consistent.

(In the table, etcd is a small key-value store built on Raft, often used to hold locks and leader elections for other systems.) Whichever row you pick, its gap is still there in production: replicas lag, standbys go away, hints pile up. The last section is about watching for those gaps on a running system.

08Running replication in production

8.1Watching the gap

Each question the chapter raised has something you can query on a running system.

SQL
-- How far behind is each standby, in time and in bytes? (sections 1 and 4.4)
-- Postgres, on the primary: one row per standby
select application_name, state, sync_state,
       write_lag, flush_lag, replay_lag,
       pg_wal_lsn_diff(pg_current_wal_lsn(), replay_lsn) as replay_bytes_behind
from pg_stat_replication;
 
-- How stale does this standby look? (section 1)
-- Postgres, on a standby
select now() - pg_last_xact_replay_timestamp() as apparent_lag;
 
-- Which commits are waiting on a synchronous standby right now? (section 4.3)
select pid, now() - query_start as waiting, query
from pg_stat_activity where wait_event = 'SyncRep';

pg_last_xact_replay_timestamp() is the time of the last replayed commit, so on an idle primary the "lag" grows even though the standby is fully caught up. Measure lag in bytes, or write a heartbeat row every second and measure its age.

For Kafka (section 4.5), the two to alert on are UnderReplicatedPartitions (followers falling out of the ISR) and UnderMinIsrPartitionCount (partitions now refusing acks=all writes). For Cassandra (section 5.4), watch hint backlog and when each table was last repaired.

8.2Rules that hold up

  1. Choose the promise per feature. Use the table in 7.1, and don't make everything linearizable.
  2. Give one user read-your-writes cheaply. Pin that user to the primary after a write, or have the reader wait for the write's LSN, before reaching for remote_apply for everyone.
  3. If you rely on a synchronous standby for durability, use a quorum of standbys (ANY 1 (s1, s2)), alert on sessions waiting in SyncRep, and never cancel them.
  4. For Kafka with replication factor 3, use acks=all with min.insync.replicas=2.
  5. Measure lag in bytes or with a heartbeat row, and not with the replay timestamp on a quiet primary.
  6. Run Cassandra repair well inside gc_grace_seconds if you delete data.
  7. Don't use last-writer-wins where a lost concurrent write matters. Use version vectors with a merge, a CRDT, or conditional writes.

8.3What you trade for what

You getYou payWhen the bill arrives
Fast commits from asynchronous replicationA write the client was told about can be missing on the replicaAfter a failover
synchronous_commit = on: the write is durable on a standbyA second fsync on every commit, and standby reads can still be staleAs commit latency, and as the roughly 3% stale reads in 4.4
remote_apply: standby reads see the writeEvery write waits for replayAs commit latency, on every write
Quorums that ignore a slow nodeW + R > N isn't linearizable without read repairAs a reader that sees an older value than an earlier reader
Sloppy quorums: writes accepted while home replicas are downReads can miss acknowledged writesUntil hints are handed back
Last-writer-wins: convergence with no coordinationConcurrent writes silently discardedWhen two users edit the same thing
CRDTs: automatic mergingNo uniqueness or boundsWhen you need "only one" or "never below zero"

8.4Symptom, cause, fix

SymptomLikely causeFix
A user's edit "disappears" on reloadRead went to a lagging replicaPin to primary after write, or LSN wait
Items flicker between two statesReads alternate between replicas at different positionsSticky routing per session
All writes hang; primary is healthySynchronous standby downANY 1 (s1, s2) quorum; alert on SyncRep waits
Acknowledged writes missing after failoverAsync replication, or a cancelled sync waitSync standby, acks=all with min.insync.replicas=2, don't kill SyncRep sessions
Deleted rows come back (Cassandra)Repair not run within gc_grace_secondsSchedule repair well inside the window
A concurrent update silently lostLast-writer-wins on a multi-leader or leaderless storeVersion vectors with merge, CRDTs, or conditional writes
NotEnoughReplicasException from producersISR shrank below min.insync.replicasFind the slow broker: disk, network, GC

09Summary

  1. A reload can show the old name because a replica replays the log after the primary commits. The replica had the record and the old row at once, and the gap between them is its lag.
  2. Replication trades freshness for latency and availability. Every design waits for some copies and not others, and the consistency model says what readers see as a result.
  3. Linearizability means one copy and real-time order. It's the only model where "I wrote it, then told you, and you read it" is safe, and it can't stay available in a partition.
  4. Session guarantees fix most replica-lag bugs cheaply. Read your writes and monotonic reads are promises to one client, so you pay for them only where that client reads.
  5. Postgres's default synchronous commit makes a write durable on the standby before it's visible there. In the 2,000-read test, about 3% of immediate standby reads were stale with on, and none with remote_apply.
  6. A cancelled synchronous commit is still committed. Postgres returns success with a warning, and the row exists only on the primary.
  7. Kafka waits for whichever replicas are currently in sync. acks=all means "the current ISR", which can shrink to the leader alone, so min.insync.replicas is what sets a floor.
  8. W + R > N guarantees that a read overlaps every completed write. That still isn't linearizable, because an in-flight write can let a later read see an older value, unless readers write back (ABD, blocking read repair).
  9. Sloppy quorums give up the overlap. They keep writes available and let reads miss them until hints are delivered.
  10. Last-writer-wins converges by throwing writes away. Version vectors detect concurrency, and CRDTs resolve it without the application.
  11. CRDTs can't protect invariants. Anything with a uniqueness or a bound needs a leader or consensus.

10Build this

A consistency checker for your own database.

  • Start a Postgres primary and one standby (pg_basebackup -R does most of the work) and reproduce the stale-read table in section 4.4. Then add pg_current_wal_insert_lsn() to every write, make the reader wait for pg_last_wal_replay_lsn() to pass it, and watch the stale count go to zero with synchronous_commit = local.
  • Write a three-replica in-memory key-value store in one process, with configurable N, W, R and a per-message random delay. Record every operation's start time, end time and result, then check the history for a read that returned an older value than a read that finished before it started. Find the anomaly from section 5.3, then add read repair and watch it disappear.
  • Implement a PN-counter and an OR-set, and write a randomized test that merges replicas in random orders with duplicated messages. Assert that they all converge.

11Interview questions

beginnerWhat's the difference between synchronous and asynchronous replication?›

With asynchronous replication the primary acknowledges a commit once it's durable locally, then streams it to replicas, so a primary crash can lose recent acknowledged writes. With synchronous replication the commit waits for at least one replica to confirm. Postgres lets you choose what "confirm" means: written, flushed or replayed (remote_write, on, remote_apply). Only the last makes the write visible to queries on the standby.

beginnerA user updates their profile and the next page shows the old one. What happened and how do you fix it?›

A read-your-writes violation: the read went to a replica that hadn't replayed the write. Fixes, cheapest first: route that user's reads to the primary for a short window after a write; carry the write's LSN and have the replica wait until it has replayed past it; or use remote_apply so commits wait for replay on the synchronous standbys.

intermediateExplain linearizability versus sequential consistency.›

Both require a single total order of operations consistent with each client's program order. Linearizability also requires that order to respect real time: if operation A finished before B started, A comes first. Sequential consistency allows a read to return an older value, provided all clients agree on the history. The practical difference is whether an out-of-band message ("I just saved it") guarantees the other party will see the write.

intermediateN = 3, W = 2, R = 2. Is this system linearizable?›

Not by itself. Every read overlaps every completed write, but a write in progress can be seen by one read and missed by a later one on a different quorum. Readers that write the value back to a quorum before returning (ABD, Cassandra's blocking read repair) fix that for reads and writes of a single register. Compare-and-set still needs consensus, such as Cassandra's Paxos lightweight transactions.

intermediateWhat does min.insync.replicas do in Kafka, and when does it apply?›

It sets the smallest ISR a partition may have and still accept acks=all writes; below it the leader throws NotEnoughReplicasException. It doesn't affect acks=0 or acks=1 producers at all. With replication factor 3, acks=all plus min.insync.replicas=2 means any one broker can fail without losing an acknowledged write or stopping producers.

deepYour Postgres uses synchronous replication, yet after a failover an acknowledged order is missing. How?›

The commit was written locally first, and the backend then waited for the standby. If that wait was cancelled (statement cancel, pooler timeout, pg_cancel_backend, even a terminated connection), Postgres returns success with a WARNING that the transaction "has already committed locally, but might not have been replicated". A failover to that standby loses it. Other paths: a standby removed from synchronous_standby_names, or synchronous_commit set to local or off for that session.

deepWhy can't a CRDT enforce 'account balance never goes negative'?›

CRDT replicas accept updates without coordinating, and merge afterwards. Two partitioned replicas can each approve a withdrawal that's fine locally and overdraws the account together, and no merge can undo that without losing a withdrawal. Invariants over the global state need coordination at write time: a single leader, consensus, or escrow (splitting the balance into per-replica allowances in advance).

deepWhat does a sloppy quorum give up, and why would Dynamo accept that?›

The overlap argument. When home replicas are down, writes go to stand-in nodes that hold hints, so a later quorum read of the home replicas can miss an acknowledged write until handoff completes. Dynamo's priority was that a shopping-cart write should never be refused, and a stale read, reconciled later by merging, was the acceptable cost.

12Go deeper

check yourself
synchronous_commit = on, and the commit returns. Where is the row visible?›

On the primary. On the synchronous standby it's durable but may not be replayed yet, so standby queries can miss it. Only remote_apply waits for replay.

N = 5. What's the smallest W with R = 2 that still guarantees overlap?›

W = 4, since W + R must exceed 5. That write tolerates only one replica down. W = R = 3 tolerates two, so it's the usual choice.

Why is per-entry max a valid merge for a G-counter but plain max of totals isn't?›

Max of totals loses increments made on different replicas: 3 and 2 merge to 3, not 5. Per-entry max keeps each replica's own count, so the sum is exact, and it's still idempotent and commutative.

Which is cheaper in DynamoDB, and by how much: an eventually consistent or a strongly consistent read?›

Eventually consistent reads cost half as much. GSI reads can only be eventually consistent.

Jepsen: Consistency Models

A clickable map of every model in this chapter and the ones between them, with precise definitions and which are available in a partition. jepsen.io/consistency

DeCandia et al., Dynamo (SOSP 2007)

Consistent hashing, sloppy quorums, hinted handoff, vector clocks and the shopping cart, in one paper. PDF

Shapiro et al., Conflict-free Replicated Data Types

The formal treatment: state-based and operation-based CRDTs, and a catalogue of counters, registers, sets and graphs. Comprehensive study

src/backend/replication/syncrep.c

Postgres synchronous replication in about a thousand lines, including the comments on why a cancelled wait can't abort. REL_16_4

Bailis et al., Probabilistically Bounded Staleness (VLDB 2012)

How stale a partial quorum is in practice, modelled from production latency distributions. PDF

Verbitski et al., Amazon Aurora (SIGMOD 2017)

Six copies, a 4/6 write quorum and a 3/6 read quorum, and the AZ+1 argument behind them. Paper

Consensus: Raft, Paxos & Leases

How a replicated log gets linearizable writes and reads, and what a quorum read costs. Chapter 27.

Failure Detection & Membership

When a replica counts as down, and why failovers lose writes when that guess is wrong. Chapter 30.

PostgreSQL Deep Dive

WAL, the standby's replay loop and replication slots in more depth. Chapter 21.

Kafka & the Log as a Primitive

The ISR, the high watermark and leader epochs through a broker failure. Chapter 23.

Redis Internals

Asynchronous replication and the acknowledged write that vanishes in a failover. Chapter 22.

Time, Clocks & Ordering

Why wall-clock timestamps can't order concurrent writes, the very thing last-writer-wins relies on. Chapter 26.

Filesystems & the Page Cache

What a flush (fsync) waits for, which every durable commit in this chapter depends on. Chapter 08.