10.1 Linearizability
Consistency models
- yes
- recency
- yes
- total order
Linearizable. Every read returns the most recent write, as if there were one copy. Requires cross-node coordination on every operation, and by the CAP result it must give up availability under a partition.
Other names for the same thing: atomic consistency, strong consistency, immediate consistency, external consistency.
The basic idea: make a system APPEAR AS IF THERE IS ONLY ONE COPY OF THE DATA, and all operations on it are ATOMIC.
In a linearizable system, AS SOON AS ONE CLIENT SUCCESSFULLY COMPLETES A WRITE, ALL CLIENTS READING FROM THE DATABASE MUST BE ABLE TO SEE THE VALUE JUST WRITTEN.
Linearizability is a RECENCY GUARANTEE: the value read is the most recent, up-to-date value, not from a stale cache or replica.
1.1 The sports website — why it's a violation
“If they had hit reload at the same time, it would have been less surprising to get different results, because they wouldn't know exactly when their requests were processed. However, Bryce knows he hit reload after hearing the score, and therefore expects his result to be at least as recent as Aaliyah's.”
1.2 Building up the definition
Vocabulary: a register = one key in a KV store, one row, one document. Operations:
Read(x) ⇒ vWrite(x, v) ⇒ r(r = OK or Error)CAS(x, v_old, v_new) ⇒ r— atomically set x to v_new only if it currently equals v_old
Each bar is a request. The start is when the client SENT it; the end is when the client RECEIVED the response. The client doesn't know exactly when the database processed it — only that it happened somewhere in between.
Step 1 — what concurrent operations may return:
A read that completed before the write must return 0. A read that overlaps the write is concurrent with it, so it may return either 0 or 1. A read that begins after the write completed must return 1.
Step 2 — the extra constraint that makes it linearizable:
If reads concurrent with a write could freely return either value, readers could see the value FLIP BACK AND FORTH several times while a write is going on. THAT IS NOT WHAT WE EXPECT OF A SYSTEM THAT EMULATES "A SINGLE COPY OF THE DATA."
We imagine there must be SOME POINT IN TIME (between the start and end of the write) at which the value ATOMICALLY FLIPS from 0 to 1. Thus, IF ONE CLIENT'S READ RETURNS THE NEW VALUE, ALL SUBSEQUENT READS MUST ALSO RETURN THE NEW VALUE, even if the write operation has not yet completed.
A is the first to read the new value; B begins strictly after A's read, so it must also return 1. We imagine there is some point in time between the start and end of the write at which the value atomically flips from 0 to 1 — so if one client's read returns the new value, all subsequent reads must also return the new value, even if the write operation has not yet completed.
Step 3 — the full picture with CAS:
Each operation is marked with a vertical line at the time we think it took effect. Those markers are joined in a sequential order, and the result MUST BE A VALID SEQUENCE OF READS AND WRITES FOR A REGISTER (every read returns the value set by the most recent write).
The requirement of linearizability is that THE LINES JOINING UP THE OPERATION MARKERS ALWAYS MOVE FORWARD IN TIME, NEVER BACKWARD.
Four details from the worked example worth internalizing:
- B sent a read first, then D sent Write(x,0), then A sent Write(x,1) — and B's read returned 1. ✔ OK: the database processed D's write, then A's write, then B's read. Not the order they were sent, but an acceptable order, because the three requests are CONCURRENT. (Perhaps B's read was delayed in the network.)
- B's read returned 1 BEFORE A received its "write succeeded" response. ✔ OK — the OK response to A was slightly delayed in the network.
- This model assumes NO transaction isolation; another client may change a value at any time. C reads 1 then reads 2 because B changed it in between. An atomic CAS can be used to check the value hasn't been concurrently changed.
- The final read by B is NOT linearizable. It's concurrent with C's CAS (2→4). In the absence of other requests, returning 2 would be fine. However, A had ALREADY READ 4 before B's read started, so B IS NOT ALLOWED TO READ AN OLDER VALUE THAN A.
It is possible (though COMPUTATIONALLY EXPENSIVE) to TEST whether a system's behavior is linearizable, by recording the timings of all requests and responses and checking whether they can be arranged into a valid sequential order. (This is what Jepsen's Knossos checker does — Ch 9.)
Linearizability is the STRONGEST consistency model in common use. It includes read-after-write consistency, monotonic reads, and consistent prefix reads (Ch 6) and more.
1.3 Linearizability vs Serializability — the confusion that must be cleared
| Serializability | Linearizability |
|---|---|
| An isolation level of transactions — multiple objects. | A guarantee on reads and writes of a register — an individual object. |
| Transactions behave as if executed in some serial order. “It is OK for that serial order to be different from the order in which the transactions were actually run.” | Does not group operations into transactions, so it does not prevent write skew. |
| ⇒ No recency requirement. Stale reads are allowed by serializability. | Is a recency guarantee: if one operation finishes before another starts, the later one must observe a state at least as new. |
Both together = strict serializability (strong one-copy serializability, strong-1SR).
| System | What it provides |
|---|---|
| Single-node databases | typically linearizable |
| CockroachDB | serializability and SOME recency guarantees, but NOT strict serializability — because that would require expensive coordination between transactions |
| Spanner, FoundationDB | strict serializability |
The consistency model and isolation level can be chosen LARGELY INDEPENDENTLY from each other.
1.4 Where linearizability is genuinely required
① Locking and leader election.
A system using single-leader replication must ensure there is indeed only ONE leader, not several (split brain). No matter how the lease mechanism is implemented, IT MUST BE LINEARIZABLE. It shouldn't be possible for two nodes to acquire the lease at the same time.
(ZooKeeper and etcd implement this. Note: strictly speaking, ZooKeeper provides linearizable WRITES, but READS may be stale, since there is no guarantee they are served from the current leader. etcd since v3 provides linearizable reads by default. Libraries like Apache Curator provide higher-level recipes — and many subtle details are involved, e.g. the fencing issue of Ch 9.)
Granular case: Oracle RAC uses a lock per disk page, with multiple nodes sharing disk storage. Since these linearizable locks are ON THE CRITICAL PATH of transaction execution, RAC deployments usually have A DEDICATED CLUSTER INTERCONNECT NETWORK.
② Constraints and uniqueness guarantees.
A username or email must uniquely identify one user; a file storage service can't have two files with the same path. To enforce this AS THE DATA IS WRITTEN, YOU NEED LINEARIZABILITY.
It's similar to a lock: registering a username is like ACQUIRING A LOCK ON IT — very similar to an atomic CAS, setting the username to the user's ID provided it is not already taken.
Same for: a bank balance never going negative, not selling more items than are in stock, two people not booking the same seat. All require A SINGLE UP-TO-DATE VALUE THAT ALL NODES AGREE ON.
But the practical escape hatch: "In real applications it is sometimes acceptable to treat such constraints LOOSELY — if a flight is overbooked, you can move customers to a different flight and OFFER COMPENSATION. In such cases linearizability may not be needed." (Ch 13.)
A HARD uniqueness constraint requires linearizability. Other kinds of constraints — foreign-key or attribute constraints — CAN be implemented without it.
③ Cross-channel timing dependencies — the subtle one.
Notice: if Aaliyah hadn't exclaimed the score, Bryce WOULDN'T HAVE KNOWN his result was stale. The violation was noticed ONLY BECAUSE THERE WAS AN ADDITIONAL COMMUNICATION CHANNEL (Aaliyah's voice to Bryce's ears).
If file storage is not linearizable, the message queue (③④) may be faster than the storage service's internal replication, so at ⑤ the transcoder sees an old version of the file or nothing at all — and “the original and transcoded videos become permanently inconsistent.”
The cause is two communication channels between the web server and the transcoder: the file storage and the message queue. Exactly analogous to database replication plus the real-life audio channel between Aaliyah's mouth and Bryce's ears.
(The same race occurs with push notifications: the notification arrives quickly, but the subsequent data fetch goes to a lagging replica and doesn't see the data the notification was about.)
Linearizability is not the only way to avoid this, but it's the simplest to understand. IF YOU CONTROL THE ADDITIONAL COMMUNICATION CHANNEL (as with the message queue, but not with Aaliyah and Bryce) you can use alternatives similar to "reading your own writes" (Ch 6), AT THE COST OF ADDITIONAL COMPLEXITY.
1.5 Which replication methods are linearizable?
| Method | Verdict | Detail |
|---|---|---|
| Single-leader | POTENTIALLY | Linearizable as long as all reads and writes go to the leader — AND YOU KNOW FOR SURE WHO THE LEADER IS. A node may think it's the leader when it isn't, and a delusional leader that keeps serving requests is likely to violate linearizability. With asynchronous replication, failover may result in COMMITTED WRITES BEING LOST, violating both durability and linearizability. (Sharding doesn't affect it — it's a single-object guarantee.) |
| Consensus algorithms | LIKELY | Essentially single-leader replication with AUTOMATIC LEADER ELECTION AND FAILOVER, carefully designed to prevent split brain. ZooKeeper uses Zab, etcd uses Raft. ⚠️ But using consensus does not GUARANTEE all operations are linearizable: if it allows reads on a node WITHOUT CHECKING IT IS STILL THE LEADER, results may be stale |
| Multi-leader | NOT | Concurrently processes writes on multiple nodes and asynchronously replicates, producing conflicting writes |
| Leaderless | PROBABLY NOT | See below |
Why quorums are NOT automatically linearizable — the counterexample:
n = 3, w = 3, r = 2. Initial x = 0. A writer sets x = 1, sending to all three replicas.
How to fix it — and what it costs:
- The reader must perform READ REPAIR SYNCHRONOUSLY before returning results
- Before writing, the writer must READ THE LATEST STATE OF A QUORUM to fetch the greatest prior timestamp and ensure the new write has a greater one
Riak does NOT perform synchronous read repair because of the performance penalty. Cassandra DOES wait for read repair on quorum reads — but IT LOSES LINEARIZABILITY BECAUSE OF ITS USE OF TIME-OF-DAY CLOCKS for timestamps (Ch 9 §3.4).
What's more, ONLY LINEARIZABLE READS AND WRITES can be implemented this way; A LINEARIZABLE CAS CANNOT, BECAUSE IT REQUIRES A CONSENSUS ALGORITHM.
It is safest to assume that a leaderless Dynamo-style system does NOT provide linearizability, even with quorum reads and writes.
1.6 The cost of linearizability
| Replication | What happens during the partition |
|---|---|
| Multi-leader | Each region continues operating normally. Writes are simply queued up and exchanged when connectivity is restored. |
| Single-leader | The leader must be in one region. Clients in the follower region cannot make any writes nor any linearizable reads. They can read from the follower, but the results might be stale. ⇒ The application becomes unavailable in regions that can't reach the leader. |
The CAP theorem, stated honestly:
| Choice | Meaning |
|---|---|
| CP — consistent under network partitions | If your application requires linearizability and some replicas are disconnected, THOSE REPLICAS WILL BE TEMPORARILY UNABLE TO PROCESS REQUESTS: they must either wait or return an error — EITHER WAY, THEY BECOME UNAVAILABLE |
| AP — available under network partitions | If linearizability isn't required, each replica can process requests INDEPENDENTLY even when disconnected. The application REMAINS AVAILABLE, but ITS BEHAVIOR IS NOT LINEARIZABLE |
The unhelpful framing: "pick two out of three" is MISLEADING. Network partitions are A KIND OF FAULT — they aren't something you CHOOSE, they WILL HAPPEN whether you like it or not. The only way to guarantee no partitions is to have NO NETWORK — that is, only one replica — but then you don't have high availability either.
A better phrasing: EITHER CONSISTENT OR AVAILABLE WHEN PARTITIONED. A more reliable network makes this choice less often, but at some point THE CHOICE IS INEVITABLE.
The book's verdict on CAP — read this carefully:
CAP as formally defined is OF VERY NARROW SCOPE. It considers only ONE consistency model (linearizability) and ONE kind of fault (network partitions, which per Google data cause LESS THAN 8% OF INCIDENTS). It says nothing about network delays, dead nodes, or other trade-offs.
The formalization of AVAILABILITY does not match the usual meaning of the term. MANY HIGHLY AVAILABLE (fault-tolerant) SYSTEMS DO NOT MEET CAP'S IDIOSYNCRATIC DEFINITION OF AVAILABILITY. Moreover, some system designers choose (WITH GOOD REASON) to provide NEITHER linearizability NOR CAP's form of availability — SO THOSE SYSTEMS ARE NEITHER CP NOR AP.
CAP DESERVES CREDIT for a culture shift — it helped trigger the NoSQL movement — but IT HAS LITTLE PRACTICAL VALUE FOR DESIGNING SYSTEMS TODAY, and has been SUPERSEDED BY MORE PRECISE RESULTS. There is a lot of misunderstanding and confusion around CAP, and IT DOES NOT HELP US UNDERSTAND SYSTEMS BETTER, SO IT'S BEST NOT TO DWELL ON IT.
(PACELC generalizes it: during a Partition choose A or C; Else, when there's no partition, choose Latency or Consistency. But it inherits several of CAP's problems, such as the counterintuitive definitions.)
1.7 The REAL reason linearizability is rare: latency, not fault tolerance
Surprisingly few systems are linearizable in practice. EVEN RAM ON A MODERN MULTI-CORE CPU IS NOT LINEARIZABLE. If a thread on one core writes to a memory address and a thread on another core reads it shortly after, it is not guaranteed to read the value written — unless a memory barrier / fence is used.
Why? Every CPU core has its OWN CACHE AND STORE BUFFER. Reads are served from the cache; changes are asynchronously written to main memory. There are now MULTIPLE COPIES OF THE DATA, asynchronously updated ⇒ LINEARIZABILITY IS LOST.
It makes NO SENSE to use the CAP theorem to justify the multi-core memory consistency model. Within one computer we assume reliable communication, and we don't expect one CPU core to keep operating if disconnected from the rest. THE REASON FOR DROPPING LINEARIZABILITY IS PERFORMANCE, NOT FAULT TOLERANCE.
The same is true of many distributed databases: they drop linearizable guarantees PRIMARILY TO INCREASE PERFORMANCE, not so much for fault tolerance. LINEARIZABLE SYSTEMS TEND TO BE HIGHER LATENCY — ALL THE TIME, NOT ONLY DURING A NETWORK FAULT.
And there is no clever way out:
Attiya and Welch prove that IF YOU WANT LINEARIZABILITY, THE RESPONSE TIME OF READ AND WRITE REQUESTS IS AT LEAST PROPORTIONAL TO THE UNCERTAINTY OF DELAYS IN THE NETWORK. In a network with highly variable delays — most computer networks — the response time of linearizable reads and writes is INEVITABLY GOING TO BE HIGH. A FASTER ALGORITHM FOR LINEARIZABILITY DOES NOT EXIST.