💻 Computer Science · Graduate · CS 480

Distributed Systems

A graduate course in distributed systems for readers who can already build a service that works on one machine. The organizing fact is partial failure: some components fail while others keep running, and no participant can distinguish a crashed peer from a slow one. Everything else follows. You will work through system models and their synchrony assumptions, logical clocks and causality,…

Start the interactive course (quizzes, progress, videos) →

Free forever. No sign-up, no ads. 17 lessons. The full lesson text is below so you can read it right here.

Module 1: Partial Failure and System Models

What actually distinguishes a distributed system, the models used to reason about one, and why the clock on the wall cannot order its events.

Partial Failure: The One Thing That Is Different

  • Enumerate the outcomes of a remote call and explain why the caller cannot distinguish among several of them.
  • Prove the Two Generals impossibility and say what it implies for acknowledgement protocols.
  • Distinguish at-most-once, at-least-once and exactly-once semantics, and explain what idempotence actually buys.

In 1994 four engineers at Sun Microsystems, Jim Waldo, Geoff Wyant, Ann Wollrath and Sam Kendall, published a paper attacking the most attractive idea in their industry. The idea was transparency: a call to an object on another machine should look, in your source code, exactly like a call to an object in your own address space. Their argument was that the differences cannot be hidden, and that one of them is not a matter of degree at all. Latency is a thousand times worse and memory access is different in kind, but the difference that breaks the abstraction outright is partial failure.

What a local call can do, and what a remote call can do

Call a function in your own process and there are two outcomes: it returns, or it raises. Either way you know which. If the machine dies, your caller dies with it, so there is no case where the caller is alive and uncertain.

Now call a service over a network. Enumerate honestly:

What happenedWhat the caller observes
Request lost in the networkTimeout
Server crashed before processingTimeout
Server processed it, then crashed before replyingTimeout
Server processed it and replied, response lostTimeout
Server is alive but paused for 40 seconds by garbage collectionTimeout, then possibly a late reply
Everything worked, the reply is simply slowTimeout, then a reply

Six situations, one observation. The caller cannot tell whether the work was done. That is partial failure: some components fail while others keep running, and the survivors cannot determine which is which. Every hard problem in this course descends from that sentence.

Notice in particular the fifth row. A process that pauses for tens of seconds is indistinguishable from a dead one, and such pauses are not hypothetical: a stop-the-world garbage collection on a large heap, a virtual machine live-migrating, a laptop lid closing, a disk that is slow rather than broken, or a thread descheduled under load all produce them. Any protocol whose correctness depends on "the node did not respond, therefore it is dead" is wrong, and the fix, fencing tokens, appears in a later lesson.

Why this matters: the interesting question in a distributed system is almost never "what if a machine fails". It is "what if a machine fails and nobody can agree on whether it did".

The Two Generals, proved

Two armies camp on hills either side of a valley holding the enemy. They will win if they attack simultaneously and lose if only one attacks. Their only way to communicate is a messenger crossing the valley, who may be captured. Can they agree on a time?

No, and the proof is short. Suppose some protocol guarantees agreement using at most k messages in the worst case, and consider the shortest such protocol. Look at its final message. Since messages can be lost, the protocol must still work when that message is lost, which means the outcome does not depend on its arrival. Therefore the sender's decision cannot depend on it having been received, and the recipient's decision cannot depend on it at all. So the message is unnecessary, and deleting it gives a shorter correct protocol, contradicting minimality. Peel messages off one at a time and you reach a protocol with no messages, which obviously cannot coordinate two parties. Hence no finite protocol exists.

A -> B  "attack at dawn"          B does not know A knows it arrived
B -> A  "acknowledged"            A does not know B knows the ack arrived
A -> B  "acknowledging your ack"  and so on, forever

What the impossibility rules out is certainty. What engineers actually do is accept a probability: retry until an acknowledgement arrives, and accept that the last acknowledgement in any exchange is unconfirmed. TCP does exactly this and nothing more. The result is not a curiosity; it is why a payment API cannot promise you that the charge either definitely happened or definitely did not, and why it gives you an idempotency key instead.

At most once, at least once, and the phrase that is always wrong

Given that a timeout tells you nothing, you must choose what to do next, and the choice defines your delivery semantics.

PolicyOn timeoutRiskSuitable for
At most onceGive upThe operation may never happenBest-effort telemetry, cache invalidation hints
At least onceRetry until acknowledgedThe operation may happen more than onceAlmost everything, combined with idempotence
Exactly onceNot achievable as a delivery guaranteeBelieving you have itNothing

The third row deserves care, because the phrase is used constantly and is not simply marketing. Exactly-once delivery is impossible, by the Two Generals argument: the sender can never know a message arrived, so it must either risk not sending again or risk sending twice. Exactly-once processing is achievable, and it works by combining at-least-once delivery with a receiver that discards duplicates:

client:  generate a unique request id once, reuse it on every retry
server:  in the SAME transaction as the effect,
             insert the request id into a processed-requests table
         if the insert conflicts, the work was already done: return the stored result

The atomicity of that pair is the whole mechanism. Deduplicating outside the transaction that performs the effect leaves a window in which the effect happened and the record of it did not.

The cheaper cousin is an idempotent operation, one whose repetition changes nothing: setting a value, deleting by key, or moving a state machine to a named state. Incrementing a counter, appending to a list and charging a card are not idempotent, and turning them into idempotent operations, usually by having the client supply the identity of the effect, is one of the most reliable design moves in this subject.

Remember: retries are not a detail of the transport layer. Choosing at-least-once forces a decision about duplicate handling all the way up in the application's data model.

The eight fallacies, with their consequences

Around 1994, Peter Deutsch at Sun listed the assumptions that new distributed programmers make and that later cost them; James Gosling added the eighth. They are worth reading as a list of specific bugs rather than as a slogan.

FallacyThe failure it produces
The network is reliableNo retry path; a dropped packet becomes lost data
Latency is zeroA loop making one remote call per item, fine in testing on localhost
Bandwidth is infiniteAn API that returns a whole table because it did so when the table was small
The network is secureUnauthenticated internal traffic, trusted because it is internal
Topology does not changeHardcoded addresses, and no reaction to a rebalanced cluster
There is one administratorAn upgrade that assumes every node runs the same version
Transport cost is zeroSerialization and cross-zone traffic dominating the bill
The network is homogeneousProtocol assumptions that hold on one link type and not another

Latency, in numbers

The second fallacy is worth quantifying, because the gaps are larger than intuition suggests. These are orders of magnitude rather than measurements, and they shift with hardware, but the ratios are stable.

OperationApproximate timeRelative
L1 cache reference1 nanosecond1
Main memory reference100 nanoseconds100
Solid state disk random read100 microseconds100,000
Round trip within a datacenter0.5 milliseconds500,000
Round trip from California to Europe150 milliseconds150,000,000

Two consequences. First, the last row is bounded below by physics: light in fibre covers roughly 200 kilometres per millisecond, so a transatlantic round trip cannot go below about 60 milliseconds no matter what you buy. Any design that requires a synchronous cross-continent round trip on the critical path has accepted that floor. Second, the gap between the datacenter row and the memory row is why a system that makes one remote call per item processed will be five orders of magnitude slower than the same loop in one process, and why batching is not an optimization but a design requirement.

What you can and cannot have

  • You cannot detect failure. You can only detect the absence of a response within a timeout you chose. Every failure detector is a guess with a tunable error rate, which the next lesson formalizes.
  • You cannot observe a global state. There is no instant at which you can read every machine's memory. What you can construct is a consistent snapshot, and that takes an algorithm.
  • You cannot rely on message order. Two messages sent in sequence over different paths may arrive in either order, and a retried message can arrive after its own duplicate.
  • You can bound uncertainty. With enough assumptions, about clock error, message delay, or the number of failures, real guarantees follow. Naming those assumptions is the subject of the next lesson.

Common misconceptions

  • "Distributed systems are hard because they are concurrent." Concurrency is hard on one machine too, and there the tools are good. What is new is that a participant can vanish mid-operation while the rest continue, and nobody can confirm it.
  • "A timeout means the server is down." It means you did not get a response in the window you chose. The server may be running, may have completed your request, and may reply after you have already retried.
  • "Exactly-once delivery is available if you buy the right queue." Delivery cannot be exactly once. Processing can be, through deduplication tied atomically to the effect, and any product claiming otherwise is describing that mechanism.
  • "Making calls transparent is good abstraction." The Waldo argument is that this abstraction leaks in exactly the situations you most need to handle. A remote call should look remote in the source, so that its failure modes are visible to the person writing the error handling.
  • "Networks in a modern cloud are reliable enough to ignore." Partitions, asymmetric reachability where A can reach B but not the reverse, and gray failures where a link is up but drops most packets are all routinely reported in production post-mortems. Reliability improves the odds; it does not change the model.

Where this leaves us

  • Partial failure, not latency or scale, is what makes a distributed system a different kind of object: components fail independently and survivors cannot tell what happened.
  • Six distinct outcomes of a remote call collapse into one observation, a timeout, so the caller never learns whether the work was done.
  • The Two Generals result proves no finite protocol achieves certain agreement over a lossy channel, which is why the final acknowledgement in any exchange is always unconfirmed.
  • At-least-once plus deduplication tied atomically to the effect is what people mean by exactly-once processing; exactly-once delivery does not exist.
  • The eight fallacies are a list of specific bugs, and the latency table shows why batching is a design requirement rather than an optimization.
  • You cannot detect failure, observe global state, or rely on ordering for free; you can buy each of them with explicit assumptions, which is what the next lesson names.

Everything from here on is an answer to one question: given that you cannot tell a dead peer from a slow one, what can still be guaranteed?

Sources

  1. Waldo, J., Wyant, G., Wollrath, A., & Kendall, S. (1994). A note on distributed computing (Technical Report SMLI TR-94-29). Sun Microsystems Laboratories. scholar.harvard.edu
  2. Kleppmann, M. (2017). The trouble with distributed systems. In Designing data-intensive applications (ch. 8). O'Reilly Media.
  3. Tanenbaum, A. S., & van Steen, M. (2017). Introduction to distributed systems. In Distributed systems: principles and paradigms (3rd ed., ch. 1). Pearson.
  4. Wikipedia contributors. (n.d.). Fallacies of distributed computing. en.wikipedia.org
  5. Wikipedia contributors. (n.d.). Two Generals' Problem. en.wikipedia.org
  6. Wikipedia contributors. (n.d.). Idempotence, including its use in request retries. en.wikipedia.org
Key terms
Partial failure
The condition in which some components of a system fail while others continue, and the survivors cannot determine which is which.
Two Generals problem
The proof that two parties communicating over a lossy channel cannot reach certain agreement with any finite protocol.
At-least-once delivery
Retrying until acknowledged, which guarantees the message arrives but permits duplicates.
Exactly-once processing
At-least-once delivery combined with duplicate suppression committed atomically with the effect; distinct from exactly-once delivery, which is impossible.
Idempotent operation
One whose repeated application has the same result as a single application, making retries safe.
Gray failure
A component that is neither up nor down, such as a link that stays connected while dropping most packets.
Fallacies of distributed computing
The eight assumptions, listed by Deutsch and Gosling, that produce predictable classes of bugs when left implicit.

System Models: Synchrony Assumptions and Failure Models

  • Distinguish the synchronous, asynchronous and partially synchronous timing models and say what each permits.
  • Order the failure models from crash-stop to Byzantine and state the replica counts each requires.
  • Separate safety from liveness and explain which one a real system is allowed to suspend.

In 1980 Marshall Pease, Robert Shostak and Leslie Lamport proved a result with an uncomfortably small number in it. If processes can behave arbitrarily, including lying differently to different peers, then three processes cannot reach agreement when one of them is faulty. Not "it is difficult": it is impossible, and you need four. The general bound is 3f + 1 processes to tolerate f arbitrary faults. That result is meaningless without a stated model, and the habit of stating the model before claiming anything is what this lesson is for.

Why a model comes first

"This algorithm is correct" is not a complete sentence. Correct under what assumptions about message delay, clock accuracy, and the ways a process can misbehave? Change any of those and the same code goes from provably correct to provably impossible. Two dimensions carry almost all of the weight: the timing model and the failure model.

Timing: three models, in order of strength

ModelAssumesConsequenceRealism
SynchronousKnown upper bounds on message delay, processing time, and clock driftTimeouts are perfect failure detectors, since exceeding the bound proves failureRare; some real-time and hardware-controlled networks
AsynchronousNo bounds at all; messages are eventually delivered but may take arbitrarily longA slow process is indistinguishable from a dead one, foreverSafe but pessimistic; consensus is impossible here, as the FLP lesson shows
Partially synchronousBounds exist but are unknown, or hold only after some unknown global stabilization timeSafety can hold always, liveness once the system stabilizesThe model real systems are designed for

The middle row is where the theory bites. If you assume nothing about timing, you cannot write a deterministic protocol that always terminates for consensus. The bottom row is Dwork, Lynch and Stockmeyer's 1988 repair, and it is not a fudge: it captures the actual behaviour of networks, which are usually well behaved and occasionally, unpredictably, not. Paxos and Raft are both designed for it, which is why both are always safe and only guaranteed to make progress during a period of stability.

What matters here: the partially synchronous model is the reason a well-built system never returns a wrong answer during a network problem, and may return no answer at all.

Failure: a ladder, not a switch

ModelA faulty process mayReplicas to tolerate f faults
Fail-stopHalt, and the halt is reliably detectable by othersf + 1 for availability
Crash-stopHalt permanently, with no announcement2f + 1 for consensus
Crash-recoveryHalt and later restart, losing anything not written to stable storage2f + 1, plus durable state on disk
OmissionDrop individual messages it should have sent or received2f + 1, with retransmission
Byzantine (arbitrary)Do anything: send wrong values, send different values to different peers, collude3f + 1

Two rows deserve comment. Fail-stop is the model beginners implicitly assume and the one the real world almost never provides, because detecting the halt reliably requires the synchronous timing model. Crash-recovery is the one production systems actually live in: machines reboot, and a replica that forgot what it promised is more dangerous than one that stayed down, which is why consensus implementations insist on flushing state to disk before replying.

The jump from 2f + 1 to 3f + 1 in the last row is worth understanding rather than memorizing. With crash faults, a majority quorum works because any two majorities intersect, and the shared member remembers the truth. With Byzantine faults the shared member might lie, so you need the intersection to contain at least one honest process: with n processes and f faulty, quorums of size n - f intersect in n - 2f processes, and requiring that to exceed f gives n greater than 3f.

A worked instance of the Byzantine bound

Three generals, one a traitor, must agree to attack or retreat, and loyal generals must reach the same decision while following a loyal commander's order.

Case A: the commander is loyal and orders ATTACK.
  Loyal lieutenant L1 receives ATTACK from the commander.
  Traitor L2 tells L1 that it was told RETREAT.
  L1 sees: commander says ATTACK, peer says the commander said RETREAT.

Case B: the commander is the traitor.
  It tells L1 ATTACK and L2 RETREAT.
  L1 sees: commander says ATTACK, peer says the commander said RETREAT.

L1's observations are IDENTICAL in both cases, yet in case A it must
attack (to follow a loyal commander) and in case B it must match L2.
No decision rule can satisfy both.

Adding a fourth process breaks the symmetry, because a majority vote among three lieutenants can outvote one liar. That is the whole content of 3f + 1, and it is why Byzantine fault tolerance is expensive enough that it is reserved for settings with genuinely untrusted participants, such as blockchains, or for the rare hardware-critical system.

Failure detectors: making the assumption explicit

Rather than assuming timing bounds directly, Chandra and Toueg proposed packaging the uncertainty into a module. A failure detector is an oracle at each process that outputs a set of processes it currently suspects, and it is classified by two properties:

  • Completeness: every crashed process is eventually suspected. This is the easy half; suspecting everyone achieves it.
  • Accuracy: correct processes are not suspected, in some form. This is the hard half, and it is where the timing assumption hides.

A perfect detector never suspects a live process and requires synchrony. An eventually perfect detector may make mistakes for a while but eventually stops, which is exactly what a timeout with exponential backoff provides in a partially synchronous network. Chandra, Hadzilacos and Toueg showed that the weakest detector sufficient for consensus is Omega, which eventually outputs the same correct process as leader at every correct process. That is a precise statement of a familiar engineering intuition: consensus is exactly as hard as electing a stable leader, and no easier.

In practice a failure detector is a heartbeat with a timeout. Phi-accrual detectors, used in Cassandra and Akka, refine this by reporting a continuously varying suspicion level derived from the observed distribution of heartbeat arrival times, so a link that is slow but consistent is not mistaken for a dead one.

Safety and liveness

Every property in this course is one of two kinds, and confusing them is the most common source of muddled requirements.

SafetyLiveness
InformallyNothing bad happensSomething good eventually happens
Violated byA finite prefix of an execution; once broken, it stays brokenOnly an infinite execution; at any finite point it may still be satisfied
ExamplesTwo replicas never commit different values at the same log position; a committed transaction is never lostEvery request eventually receives a response; a leader is eventually elected
Under a partitionMust be preservedMay be suspended

Alpern and Schneider showed in 1985 that every property is an intersection of a safety property and a liveness property, so the classification is exhaustive rather than a rule of thumb. The design consequence is blunt: when the network misbehaves, give up liveness, never safety. A system that stops responding during a partition is behaving correctly. A system that keeps responding by committing conflicting values is broken, and the damage outlives the partition.

The upshot: an impossibility result in this field is almost always about liveness. FLP does not say a protocol can be wrong; it says no deterministic protocol can guarantee it will always finish.

Assembling a model

Choosing one row from each table gives you a setting, and the settings differ sharply in what is achievable.

TimingFailuresConsensus?Representative system
AsynchronousCrash-stop, one faultImpossible deterministically (FLP)The theoretical worst case
Asynchronous plus randomizationCrash-stopPossible with probability 1Ben-Or style protocols
Partially synchronousCrash-recovery, f of 2f + 1Yes, safe always and live after stabilizationPaxos, Raft, ZooKeeper, etcd
Partially synchronousByzantine, f of 3f + 1Yes, with more messages and cryptographyPBFT and its descendants
SynchronousCrash-stopYes, and simply, since timeouts detect failure exactlySome hardware and avionics buses

Common misconceptions

  • "The asynchronous model means the network is slow." It means there is no bound on delay, so no timeout can be justified. It is a statement about knowledge, not about speed.
  • "Real networks are synchronous enough, so the theory is academic." A network that is usually fast and occasionally arbitrarily slow is exactly the partially synchronous model, and the theory for it is what tells you which guarantees survive the occasional part.
  • "Byzantine tolerance is just extra checksums." A Byzantine process can send well-formed, correctly signed, mutually inconsistent messages. The defence is the 3f + 1 quorum structure, not integrity checks.
  • "A failure detector detects failures." It suspects. Its usefulness is characterized by how wrong it is allowed to be and for how long, and every implementation is a timeout with a policy.
  • "Safety and liveness are both negotiable under load." Liveness is negotiable and often deliberately sacrificed. Trading away safety converts a temporary outage into permanent data corruption.
  • "Crash-stop and crash-recovery are nearly the same." A recovered node with amnesia can contradict a promise it already made, which is why consensus protocols require durable writes before responding.

Putting it together

  • A correctness claim is meaningless without a timing model and a failure model, because the same protocol can be correct in one and impossible in another.
  • Synchronous assumes known bounds, asynchronous assumes none, and partially synchronous assumes bounds that hold eventually; the last is the model production systems target.
  • The failure ladder runs from fail-stop through crash-stop and crash-recovery to Byzantine, requiring 2f + 1 replicas for crash faults and 3f + 1 for arbitrary ones.
  • The 3f + 1 bound comes from needing quorum intersections to contain an honest process, and the three-general case shows why three is not enough.
  • Failure detectors package the timing assumption as completeness plus accuracy, and Omega, eventual agreement on a leader, is the weakest one sufficient for consensus.
  • Safety is violated by a finite prefix and must never be given up; liveness may be suspended, and impossibility results in this field are almost always about liveness.

One assumption remains unexamined, and it is the one most systems quietly rely on: that the clocks agree. The next lesson takes it apart.

Sources

  1. Lamport, L., Shostak, R., & Pease, M. (1982). The Byzantine generals problem. ACM Transactions on Programming Languages and Systems, 4(3), 382-401. lamport.azurewebsites.net
  2. Dwork, C., Lynch, N., & Stockmeyer, L. (1988). Consensus in the presence of partial synchrony. Journal of the ACM, 35(2), 288-323.
  3. Chandra, T. D., & Toueg, S. (1996). Unreliable failure detectors for reliable distributed systems. Journal of the ACM, 43(2), 225-267.
  4. Alpern, B., & Schneider, F. B. (1985). Defining liveness. Information Processing Letters, 21(4), 181-185.
  5. Cachin, C., Guerraoui, R., & Rodrigues, L. (2011). Introduction to reliable and secure distributed programming (2nd ed.). Springer.
  6. Wikipedia contributors. (n.d.). Byzantine fault, including the 3f + 1 bound. en.wikipedia.org
  7. Wikipedia contributors. (n.d.). Failure detector. en.wikipedia.org
Key terms
Synchronous model
Known upper bounds on message delay, processing time and clock drift, which makes a timeout a perfect failure detector.
Asynchronous model
No timing bounds at all, so a slow process cannot be distinguished from a crashed one.
Partial synchrony
Bounds that are unknown or hold only after an unknown stabilization time; the model real consensus systems target.
Crash-recovery
A failure model in which a halted process may restart, retaining only what it wrote to stable storage.
Byzantine fault
Arbitrary misbehaviour, including sending different values to different peers, requiring 3f + 1 processes to tolerate f faults.
Failure detector
A module that suspects processes, classified by completeness and accuracy; Omega is the weakest one sufficient for consensus.
Safety property
A property violated by a finite prefix of an execution, such as never committing two different values at one log position.
Liveness property
A property violated only by an infinite execution, such as eventually responding to every request; it may be suspended during a partition.

Physical Time, and Why It Cannot Order Your Events

  • Explain clock drift, the difference between wall-clock and monotonic clocks, and what NTP does and does not guarantee.
  • Show how clock skew silently loses writes under last-write-wins conflict resolution.
  • Describe TrueTime's uncertainty interval and commit-wait, and state the latency it costs.

At midnight UTC on 31 December 2016 a leap second was inserted: the last minute of the year had 61 seconds. Some of Cloudflare's servers, asked how much time had elapsed between two moments, returned a negative number. The code receiving that number had been written on the assumption that time moves forward, and a portion of DNS queries began to fail. The bug was not in the clock. It was in the belief that a timestamp is a number you can subtract.

Three ways a clock lies

Drift. A server's clock is a quartz oscillator, and its frequency depends on temperature and manufacture. Typical drift is tens of parts per million, which works out to seconds per day. Two machines racked side by side will disagree by seconds within a week if left alone.

Correction. Because of drift, machines run NTP, which compares against a hierarchy of reference clocks and adjusts. The adjustment comes in two flavours, and the difference matters enormously. Slewing speeds the clock up or slows it down slightly until it converges, so time still moves forward monotonically. Stepping jumps it, possibly backwards. A clock that jumps backwards makes an event that happened later carry an earlier timestamp than one that happened before it.

Leap seconds. Earth's rotation is not exactly 86,400 seconds per day, so occasional leap seconds are inserted, twenty-seven of them since 1972. The operating system has to represent 23:59:60, and historically several did so by repeating 23:59:59, which means the same second occurs twice. Google, Amazon and others avoid the problem by smearing: spreading the extra second across a 24-hour window as a tiny frequency change, so no repeated or reversed second is ever visible to applications. In 2022 the international body responsible voted to stop introducing leap seconds by 2035, which fixes the problem forty years after it started causing outages.

The point: a wall-clock timestamp is a best guess about a global quantity, produced by a local device that drifts and is periodically corrected by a process that can move it in either direction.

Two clocks, two jobs

Every modern platform exposes two clocks, and choosing wrongly is the most common time bug in production code.

Wall-clock (real-time) clockMonotonic clock
AnswersWhat time is it?How much time has passed?
Can jumpYes, forwards or backwards, whenever NTP steps itNo; it only ever increases
Comparable across machinesLoosely, within the synchronization errorNever; its zero point is arbitrary
Use forLog timestamps, expiry dates shown to users, anything a human readsTimeouts, retry backoff, latency measurement, lease durations

The Cloudflare incident is exactly this table's third row applied wrongly: an elapsed duration computed from wall-clock readings. Any subtraction of two wall-clock readings is a bug waiting for the clock to be adjusted between them.

How good is NTP, actually?

Over the public internet, a well-configured NTP client typically lands within tens of milliseconds of true time; on a local network with good servers, under a millisecond is achievable; and the Precision Time Protocol, with hardware timestamping in the switches, reaches sub-microsecond accuracy. Those are the good cases. The bad cases are what protocols must survive: a misconfigured server, an asymmetric network path that biases the round-trip estimate, a virtual machine whose clock stalls while it is descheduled, or a firewall silently blocking NTP so a machine drifts freely for months.

The important property is not the typical accuracy. It is that NTP gives you no bound you can rely on. The client's estimate of its own error is itself an estimate. This is precisely the difference between the synchronous and partially synchronous models from the previous lesson, appearing now as a concrete engineering fact.

Where this breaks a real system: last write wins

Consider two replicas of a key-value store, each accepting writes, resolving conflicts by keeping the value with the larger timestamp. This scheme is used widely, and its failure mode is silent.

real time    node A (clock 100 ms fast)     node B (clock accurate)
--------------------------------------------------------------------
12:00:00.000  client 1 writes x = "red"
              A stamps it 12:00:00.100
12:00:00.050                                client 2 writes x = "blue"
                                            B stamps it 12:00:00.050

replication: both values meet. 12:00:00.100 > 12:00:00.050, so "red" wins.

But "blue" was written LATER in real time. Client 2's write is
discarded, with no error, no conflict report, and no log entry.

Nothing here is a bug in the store. The store did what it promised. The lost write is a direct consequence of using a quantity that is only approximately global to make a decision that requires an exact ordering. And the loss is undetectable after the fact, because the discarded value leaves no trace.

The problem generalizes to any use of timestamps for correctness: expiring a lease by comparing wall-clock times across machines, ordering entries in a merged log, deciding which of two cached values is newer. Each is the same mistake.

Remember: use physical time for anything a human reads, and never for a decision that requires knowing which of two events happened first.

The exception that proves the rule: TrueTime

Google's Spanner does use physical time for correctness, and the way it does so shows exactly what the ordinary approach is missing. Instead of returning a timestamp, its TrueTime API returns an interval: a pair of values with the guarantee that the true time lies between them. The interval width is the accumulated uncertainty, and Spanner's published measurements report it as a few milliseconds, under ten in the reported case, maintained by GPS receivers and atomic clocks in every datacenter.

Having an explicit uncertainty makes a new move available. When a transaction commits at timestamp T, Spanner deliberately waits until TrueTime says T is definitely in the past before releasing its locks and acknowledging.

choose commit timestamp T = TT.now().latest
... do the commit work ...
wait until TT.now().earliest > T          <-- commit-wait, about 2 * epsilon
then release locks and reply

The effect is that no transaction can start after this one committed and receive a smaller timestamp, so timestamp order matches real-time order. That property, called external consistency, is what lets Spanner offer globally consistent reads at a timestamp without contacting every replica. The price is written on the label: every commit pays roughly twice the uncertainty in added latency, which is why Google invested in hardware to make the uncertainty small. Most systems have no such hardware, which is why most systems must not do this.

What to use instead

If physical time cannot order events, something else must. The candidates, in ascending order of cost:

MechanismOrdersCost
Sequence number from a single leaderEverything that leader seesA leader, and a failover story
Lamport timestampsCausally related events, consistentlyOne counter per process; cannot detect concurrency
Vector clocksCausally related events, and detects concurrencyOne entry per process, so metadata grows with the cluster
Hybrid logical clocksCausality, with values close to physical timeSlightly more state; useful when humans must read the timestamps
Consensus on a total orderEverything, absolutelyA round trip to a quorum per operation

The next lesson develops the middle rows, which are the ones that require no special hardware and no coordination at all.

Common misconceptions

  • "NTP keeps our servers within a few milliseconds, so timestamps are fine." Usually true, occasionally not, and there is no bound. Correctness cannot rest on a property that holds most of the time and fails without warning.
  • "Clocks only drift forward." Corrections move them in both directions, and a backwards step is what turns an elapsed-time calculation negative.
  • "Use UTC everywhere and the problem is solved." Time zones are a formatting concern. Drift, stepping and leap seconds are all present in UTC.
  • "Last write wins is a conflict resolution strategy." It is a data loss strategy with a schedule you do not control. It is acceptable only when losing a concurrent write is genuinely harmless, and that decision belongs to the application, not the database.
  • "Spanner proves you can just use clocks." Spanner uses clocks by making the error explicit and then waiting it out, at a real latency cost, with hardware most operators do not have. It is the counterexample to using clocks naively, not a licence to do so.
  • "A monotonic clock can be compared across machines." Its origin is arbitrary, often the machine's boot time. Two monotonic readings from different hosts have no relationship at all.

What to carry forward

  • Quartz drift, NTP stepping, and leap seconds each break the assumption that a wall-clock reading is a well-behaved number.
  • Use the monotonic clock for durations, timeouts and leases, and the wall clock only for values a human will read; subtracting two wall-clock readings is a bug.
  • NTP typically achieves tens of milliseconds over the internet and better on a LAN, but it provides no guaranteed bound, which is the partially synchronous model in concrete form.
  • Last-write-wins with physical timestamps silently discards the genuinely later write whenever clock skew exceeds the gap between two writes, and leaves no evidence.
  • TrueTime returns an interval rather than an instant, and Spanner waits out the uncertainty before acknowledging a commit, buying real-time-consistent timestamps at roughly twice the uncertainty in latency.
  • Ordering without special hardware comes from logical mechanisms: leader sequence numbers, Lamport and vector clocks, or consensus.

The next lesson builds the ordering that physical time cannot provide, out of nothing but message passing and counters.

Sources

  1. Cloudflare. (2017). How and why the leap second affected Cloudflare DNS. Cloudflare blog. blog.cloudflare.com
  2. Corbett, J. C., et al. (2012). Spanner: Google's globally-distributed database. 10th USENIX Symposium on Operating Systems Design and Implementation. research.google.com
  3. Google. (n.d.). TrueTime and external consistency. Google Cloud Spanner documentation. cloud.google.com
  4. Google. (n.d.). Leap smear. Google Public NTP documentation. developers.google.com
  5. National Institute of Standards and Technology. (n.d.). Time and frequency division. nist.gov
  6. Kleppmann, M. (2017). Unreliable clocks. In Designing data-intensive applications (ch. 8). O'Reilly Media.
  7. Wikipedia contributors. (n.d.). Leap second. en.wikipedia.org
  8. Wikipedia contributors. (n.d.). Network Time Protocol. en.wikipedia.org
Key terms
Clock drift
The rate at which a local oscillator diverges from true time, typically tens of parts per million for server quartz.
Slewing and stepping
The two ways NTP corrects a clock: gradually changing its rate, or jumping it, possibly backwards.
Monotonic clock
A clock that only increases and has an arbitrary origin, suitable for durations and useless for cross-machine comparison.
Leap second
An extra second inserted to keep civil time aligned with Earth's rotation, which repeats or skips a second unless smeared.
Leap smear
Absorbing a leap second as a small frequency change spread over a long window so applications never see a repeated second.
Last write wins
Conflict resolution by keeping the larger timestamp, which silently discards the genuinely later write under clock skew.
TrueTime
Spanner's clock API returning an interval guaranteed to contain the true time, with uncertainty of a few milliseconds.
Commit-wait
Delaying a commit acknowledgement until the chosen timestamp is definitely past, so timestamp order matches real-time order.

Module 2: Causality

Ordering events without a clock, deciding when two events are genuinely concurrent, and photographing a system that never stops moving.

Lamport Clocks and Vector Clocks

  • Define the happens-before relation and identify concurrent events on a space-time diagram.
  • Compute Lamport timestamps and explain why their converse property fails.
  • Compute vector clocks and use them to distinguish causal ordering from concurrency.

A user posts a message, then three seconds later posts a reply to it. A replica in another region receives the reply first, because the two writes travelled different paths. It now holds a reply to a message that does not exist. No packet was lost, no machine failed, and no clock was consulted. Leslie Lamport gave the vocabulary for exactly this situation in a five-page paper in 1978, which became the most cited work of his career, and the vocabulary is more useful than any single algorithm in this course.

Happens-before

Define a relation on events, written a arrow b and read "a happens before b". It holds when:

  1. a and b occur in the same process and a comes first, or
  2. a is the sending of a message and b is the receipt of that same message, or
  3. there is some c with a arrow c and c arrow b (transitivity).

If neither a arrow b nor b arrow a, the events are concurrent, written a parallel b. Concurrent does not mean simultaneous. It means no information could have flowed between them, so no observer can justify ordering them, and the system is free to process them in either order.

That is the whole idea, and it is worth pausing on how little it assumes. There is no clock. The relation is defined purely by process order and message passing, which are the only things a distributed system actually knows about. Lamport drew the analogy with special relativity explicitly: happens-before is the same shape as the light cone, and concurrency is the same as spacelike separation.

Key idea: causality is not about time. It is about which events could possibly have influenced which others.

Reading a space-time diagram

Three processes, time running left to right, with two messages.

P1:   a ----- b(send m1) ------------------------------
                  \
P2:   ------------- c(recv m1) -- d -- e(send m2) ------
                                          \
P3:   f ---------------------------------- g(recv m2) -- h

Work out some relationships. a arrow b by rule 1. b arrow c by rule 2. So a arrow c, and by transitivity a arrow d, a arrow e, a arrow g, a arrow h. What about f and b? There is no chain of process order and messages connecting them in either direction, so f and b are concurrent, even though on a wall clock f may well have happened first. And f arrow g, since they are in the same process.

Lamport timestamps, and their one weakness

Give every process an integer counter L, starting at 0, and three rules:

before any local event:      L = L + 1
when sending a message:      L = L + 1;  attach L to the message
when receiving a message:    L = max(L, L_message) + 1

Apply it to the diagram:

Eventabcdefgh
Lamport L12345167

The guarantee is one-directional:

if a arrow b   then   L(a) < L(b)          TRUE, always
if L(a) < L(b) then   a arrow b            FALSE

The counterexample is already in the table. L(f) = 1 and L(b) = 2, so the timestamps suggest f came first, yet f and b are concurrent and neither influenced the other. Lamport clocks can therefore certify that something is not causally related, by finding a timestamp inversion, but they can never certify that two events are concurrent.

What they are good for is imposing a total order that is consistent with causality. Break ties by process identifier and you have a rule every process computes identically:

a comes before b  iff  ( L(a), pid(a) )  <  ( L(b), pid(b) )   lexicographically

Any such order respects causality, and different orders are all equally valid, which is exactly the freedom concurrency gives you. Lamport used this in the same paper to build distributed mutual exclusion, and the same trick orders operations in a replicated state machine.

Vector clocks: paying for the converse

To detect concurrency you need more state. Give each of n processes a vector of n counters, and change two rules:

process i, local event:      V[i] = V[i] + 1
send:                        V[i] = V[i] + 1;  attach the whole vector
receive at i:                V = elementwise max(V, V_message);  then V[i] = V[i] + 1

Compare two vectors componentwise. V(a) is less than V(b) when every component of V(a) is at most the corresponding component of V(b) and at least one is strictly smaller. If neither vector is less than the other, the events are concurrent.

On the same diagram, with the order (P1, P2, P3):

Eventabcdefgh
Vector V(1,0,0)(2,0,0)(2,1,0)(2,2,0)(2,3,0)(0,0,1)(2,3,2)(2,3,3)

Now recheck the pair Lamport clocks got wrong. V(f) = (0,0,1) and V(b) = (2,0,0). Is one at most the other? f has a larger third component and b has a larger first, so neither dominates: concurrent, which is the correct answer. And V(a) = (1,0,0) is componentwise at most V(g) = (2,3,2) with a strict decrease, so a arrow g, also correct.

The property that makes this work is an equivalence rather than an implication:

a arrow b     if and only if     V(a) < V(b)
a parallel b  if and only if     V(a) and V(b) are incomparable

The core of it: a Lamport clock summarizes causal history into one number and loses the ability to tell concurrency from ordering. A vector clock keeps one number per process and recovers it exactly. You are paying n integers for a complete answer.

What the metadata costs

The n in that cost is real. A cluster of 500 nodes attaches 500 counters to every message, and the vector must persist alongside stored values, not merely travel with messages. Three responses are used in practice:

  • Version vectors. Track one entry per replica rather than per client, which is a much smaller number and is what Dynamo-style stores do. The vector then answers "is this version a descendant of that one, or are they siblings?"
  • Pruning. Drop the oldest entries when the vector exceeds a size limit, accepting occasional false reports of concurrency. This is a deliberate trade of accuracy for space, and it must be a conservative one: reporting a genuine ordering as concurrency creates a spurious conflict, which is recoverable, while the reverse would lose data.
  • Hybrid logical clocks. Combine a physical timestamp with a logical counter so the value is both causally correct and close to wall-clock time, which makes it readable by humans and usable for snapshot reads. Kulkarni and co-authors described the construction in 2014, and it is used in CockroachDB and MongoDB.

Where the concurrency answer goes

Detecting concurrency is only useful if something acts on it. In a Dynamo-style store, a write carries the version vector it read; if the new write is a descendant of the stored version, it replaces it; if the two are incomparable, the store keeps both as siblings and returns them together on the next read, so the application resolves the conflict with knowledge the database does not have. The canonical example from the Dynamo paper is a shopping cart, where the merge rule is a union of the items, which is why an item deleted during a partition could reappear. That is a deliberate choice: for a cart, a resurrected item is better than a lost purchase.

Compare that with the last-write-wins scheme from the previous lesson, which resolves the same situation by discarding one write according to two clocks that disagree. The difference is not a detail of implementation. It is the difference between reporting a conflict to whoever can resolve it and silently deciding it wrong.

In short: vector clocks do not resolve conflicts. They detect them reliably, which is the part a database can actually do correctly.

Common misconceptions

  • "Concurrent means at the same time." It means neither event could have influenced the other. Two events an hour apart on separate machines with no communication between them are concurrent.
  • "Lamport timestamps can tell you whether two events are related." They cannot. A smaller timestamp is consistent with either causal precedence or concurrency, and only vector clocks distinguish them.
  • "Equal Lamport timestamps mean concurrency." Equal timestamps do imply concurrency, but the converse fails, and it is the converse you usually need.
  • "Vector clocks resolve write conflicts." They identify conflicts. Resolution is a semantic decision made by the application, unless the data type is a CRDT with a built-in merge.
  • "Vector clocks are too expensive to use." Per-process vectors are, in a large cluster. Per-replica version vectors are small, and that is what production systems actually attach.
  • "A total order from Lamport timestamps is the true order." It is a total order consistent with causality. Several different total orders are equally correct, and choosing among them is a decision, not a discovery.

Pulling it together

  • Happens-before is defined by process order, message send and receive, and transitivity; events unrelated by it are concurrent, meaning no information could have flowed between them.
  • Lamport timestamps increment locally and take the maximum on receipt, guaranteeing that a arrow b implies L(a) is less than L(b), with no valid converse.
  • Breaking ties by process identifier turns Lamport timestamps into a total order consistent with causality, which is enough for mutual exclusion and state machine replication.
  • Vector clocks keep one counter per process, and comparison is componentwise: strictly smaller means causal precedence, incomparable means genuine concurrency.
  • The metadata cost is one entry per process, mitigated by per-replica version vectors, pruning, and hybrid logical clocks.
  • Detecting concurrency lets a store keep siblings and hand the conflict to the application, which is strictly more honest than resolving it with clock timestamps.

The next lesson asks a harder question with the same tools: how do you record the state of every process at once, when there is no such thing as at once?

Sources

  1. Lamport, L. (1978). Time, clocks, and the ordering of events in a distributed system. Communications of the ACM, 21(7), 558-565. lamport.azurewebsites.net
  2. Fidge, C. J. (1988). Timestamps in message-passing systems that preserve the partial ordering. In Proceedings of the 11th Australian Computer Science Conference, 56-66.
  3. Mattern, F. (1989). Virtual time and global states of distributed systems. In Parallel and Distributed Algorithms, 215-226. North-Holland.
  4. Kulkarni, S., Demirbas, M., Madappa, D., Avva, B., & Leone, M. (2014). Logical physical clocks and consistent snapshots in globally distributed databases (Technical Report 2014-04). University at Buffalo. cse.buffalo.edu
  5. DeCandia, G., et al. (2007). Dynamo: Amazon's highly available key-value store. 21st ACM Symposium on Operating Systems Principles. allthingsdistributed.com
  6. Wikipedia contributors. (n.d.). Vector clock. en.wikipedia.org
  7. Wikipedia contributors. (n.d.). Happened-before. en.wikipedia.org
Key terms
Happens-before
The partial order generated by process order and message delivery, capturing which events could have influenced which.
Concurrent events
A pair unrelated by happens-before, meaning no information could have flowed between them in either direction.
Lamport timestamp
A single counter per process that guarantees causal precedence implies a smaller value, but not the converse.
Total order consistent with causality
Any tie-broken linearization of Lamport timestamps; several distinct ones are equally valid.
Vector clock
One counter per process, compared componentwise, which detects causal precedence and concurrency exactly.
Version vector
A vector clock with one entry per replica rather than per client, used to identify sibling versions of a stored value.
Sibling versions
Two concurrent versions of a value that a store keeps together and returns to the application to merge.
Hybrid logical clock
A timestamp combining physical time with a logical counter, causally correct and close to wall-clock readings.

Consistent Snapshots of a System That Never Stops

  • Define a consistent cut and identify the inconsistency that produces impossible global states.
  • Run the Chandy-Lamport marker algorithm on a two-process example, including the channel states.
  • Explain the sense in which a recorded snapshot is valid even if the system was never in that state.

By Leslie Lamport's own account, Mani Chandy posed him a problem over dinner: how do you determine the global state of a running distributed system? They had both had too much wine to think about it. He worked it out in the shower the next morning, and when he arrived at Chandy's office, Chandy had the same solution. The resulting algorithm is a page long, it is still the standard answer forty years later, and it is the backbone of checkpointing in modern stream processors.

Why you cannot just ask

Suppose a system moves money between accounts, and you want to check that the total is conserved. Ask each process for its balance and add the results. The answers arrive from different moments in different processes, and money that was in flight belongs to nobody.

P holds 100, Q holds 50.  Total should always be 150.

t=1   Q sends 30 to P     Q now holds 20, the 30 is in the network
t=2   you ask P           P has not received it yet, answers 100
t=3   the 30 arrives      P now holds 130
t=4   you ask Q           Q answers 20

your total: 100 + 20 = 120.  Thirty units have vanished.

Ask in the other order and you can get 130 + 50 = 180, inventing thirty.

Nothing failed. The naive procedure simply combined local states from moments that do not form a coherent global state.

Cuts, consistent and otherwise

Draw the space-time diagram from the previous lesson and cut it with a line crossing every process once. Everything to the left is "in the snapshot". Formally, a cut is a set of events closed under the same-process prefix: if an event is in, everything earlier in its process is in.

A cut is consistent when it is also closed under happens-before: if a receive event is in the cut, the matching send is in the cut too. The forbidden configuration is a message that has been received but not yet sent, which is not a state any execution could have reached.

CONSISTENT                       INCONSISTENT
P: --a--|--send-->               P: --a-----send---->
        |          \                        |     \
Q: --b--|-----------recv--       Q: --b-----|------recv--
    the cut precedes both              the cut has the receive
                                       but not the send

Note what a consistent cut does not require: it need not correspond to any real instant. Because concurrent events can be placed on either side, a single execution has many consistent cuts, and typically none of them is a photograph of a moment.

Why this matters: the goal is not to capture the state at an instant, which is meaningless without a global clock. It is to capture a state the system could plausibly have been in, which is enough to check any property that would have held.

The marker algorithm

Chandy and Lamport's algorithm assumes reliable, FIFO channels and no failures. Each process records its own state and the state of each incoming channel, meaning the messages that were in flight.

ANY process may initiate:
    record my own local state
    send a MARKER on every outgoing channel, before sending anything else
    begin recording incoming messages on every incoming channel

On receiving a MARKER on channel c:
    if this is the FIRST marker I have seen:
        record my own local state
        record channel c's state as EMPTY
        send a MARKER on every outgoing channel, before anything else
        begin recording on every OTHER incoming channel
    else:
        stop recording on channel c
        channel c's state = the messages recorded on it since I started

The snapshot is complete when every process has received a
marker on every one of its incoming channels.

Two design choices carry the correctness. Sending the marker before any subsequent application message means the marker acts as a boundary in the FIFO stream: everything ahead of it was sent before the snapshot and everything behind it after. And recording arrivals until the marker comes back captures exactly the messages that were in flight across the cut.

Running it, with money in flight

Two processes, P and Q, with FIFO channels in both directions. P holds 100, Q holds 50.

step 1   Q sends 30 to P.       Q's balance = 20.  The 30 is in the channel Q->P.

step 2   P initiates a snapshot.
         P records its own state: 100.   (it has not received the 30 yet)
         P sends a MARKER on P->Q.
         P begins recording arrivals on Q->P.

step 3   The 30 arrives at P.
         P's balance becomes 130, but the snapshot already fixed P at 100,
         and the 30 is recorded as the state of channel Q->P.

step 4   Q receives the MARKER on P->Q. It is Q's first marker, so:
         Q records its own state: 20.
         Q records channel P->Q as EMPTY.
         Q sends a MARKER on Q->P.

step 5   P receives the MARKER on Q->P. Not its first, so:
         P stops recording that channel. Its recorded state is { 30 }.

SNAPSHOT:  P = 100,  Q = 20,  channel Q->P = { 30 },  channel P->Q = { }
TOTAL   :  100 + 20 + 30 = 150.  Conserved.

The in-flight money is counted exactly once, in the channel rather than in either process. Skip the channel state and you get 120, which is the vanishing-money bug from the opening.

The snapshot may never have happened

Here is the subtlety that makes the algorithm intellectually interesting. The recorded state, P holding 100 while Q holds 20, was never simultaneously true by any wall clock: at the moment Q dropped to 20, P still held 100 only briefly, and events on other processes may have been interleaved arbitrarily.

What Chandy and Lamport proved is a precise substitute. Let S_start be the global state when the snapshot began and S_end the state when it finished. The recorded state S_snap satisfies:

  • S_snap is reachable from S_start by some sequence of events that actually occurred, and
  • S_end is reachable from S_snap by the remaining events.

In other words, the recorded state lies on some valid execution consistent with what happened, obtained by permuting concurrent events. That is exactly the guarantee you need for detecting a stable property, one that stays true once it becomes true: deadlock, termination, or an object having become garbage. If a stable property holds in S_snap, it held in S_end as well, because it cannot become false. If it does not hold in S_snap, it did not hold in S_start.

Worth holding on to: the algorithm does not tell you what was true at a moment. It tells you what was true in a possible history, which is sufficient for exactly the class of questions whose answers do not change back.

What it assumes, and what breaks

AssumptionIf violatedPractical response
FIFO channelsA post-snapshot message can overtake the marker and be recorded as in-flight, corrupting the cutSequence numbers per channel, or piggyback the snapshot identifier on every message
No process failuresThe snapshot never completes, because some marker never arrivesTime out and abandon the attempt; snapshotting is usually retried, not repaired
Channels do not lose messagesAn in-flight message is neither in a process nor in a recorded channelReliable delivery beneath the algorithm
Processes cooperateNo snapshot at allThe algorithm assumes a friendly system; Byzantine snapshotting is a different problem

The most-used descendant today is asynchronous barrier snapshotting in Apache Flink, described by Carbone and colleagues in 2015. Markers become barriers injected into the data streams, operators snapshot their state when barriers arrive on all inputs, and the resulting checkpoint is what a job restores from after a failure. The lineage is direct, and the reason for the design is the same: you cannot stop the world to take a picture, so you take a picture that could have been true.

Common misconceptions

  • "A snapshot captures the system at an instant." There is no shared instant. It captures a globally consistent state that could have occurred, which is a different and achievable thing.
  • "Recording each process's state is enough." Messages in flight belong to neither endpoint. Omitting channel state is what makes money disappear in the opening example.
  • "A consistent cut is one where all processes are cut at the same time." It is one where no receive appears without its send. Times are irrelevant, and the cut line may be as jagged as it likes.
  • "The snapshot can be used to check any property." It is sound for stable properties. A transient property, such as a queue being momentarily empty, may appear in a snapshot without ever having held, or hold without appearing.
  • "Markers must be sent by a designated coordinator." Any process may initiate, and several can initiate concurrently as long as snapshots are tagged, which is one reason the algorithm scales.
  • "The algorithm needs synchronized clocks." It never reads a clock. Its ordering comes entirely from FIFO channels and the marker discipline.

The takeaway

  • Polling processes independently produces impossible global states, because in-flight messages are counted twice or not at all.
  • A cut is consistent when every receive it contains has its matching send inside as well; consistency is about happens-before, not about time.
  • Chandy and Lamport's algorithm records local state, sends markers before any further application message on FIFO channels, and records arrivals on each incoming channel until that channel's marker arrives.
  • The channel states are what capture messages in flight, and omitting them is the classic bug.
  • The recorded state need never have occurred, but it is reachable from the state at initiation and reaches the state at completion, which makes it sound for detecting stable properties.
  • The assumptions are FIFO, reliable channels and no failures; Flink's barrier snapshotting is the direct modern descendant.

That closes the causality module. What follows takes the same partial order and asks a harder question: when several replicas hold the same data, which orders of operations are clients allowed to observe?

Sources

  1. Chandy, K. M., & Lamport, L. (1985). Distributed snapshots: determining global states of distributed systems. ACM Transactions on Computer Systems, 3(1), 63-75. lamport.azurewebsites.net
  2. Lamport, L. (n.d.). My writings: notes on the distributed snapshots paper and its origin. lamport.azurewebsites.net
  3. Carbone, P., Fora, G., Ewen, S., Haridi, S., & Tzoumas, K. (2015). Lightweight asynchronous snapshots for distributed dataflow. arXiv. arxiv.org
  4. Tanenbaum, A. S., & van Steen, M. (2017). Global state and distributed snapshots. In Distributed systems: principles and paradigms (3rd ed.). Pearson.
  5. Wikipedia contributors. (n.d.). Snapshot algorithm. en.wikipedia.org
  6. Wikipedia contributors. (n.d.). Chandy-Lamport algorithm. en.wikipedia.org
Key terms
Cut
A set of events closed under the same-process prefix, dividing an execution into a recorded past and an unrecorded future.
Consistent cut
A cut that also contains the send of every receive it contains, so it corresponds to a state the system could have been in.
Channel state
The messages in flight across a cut, which belong to neither endpoint and must be recorded separately.
Marker
A control message sent before any further application message on a FIFO channel, dividing pre-snapshot from post-snapshot traffic.
Stable property
One that stays true once it becomes true, such as deadlock or termination; snapshots are sound for exactly these.
Reachability guarantee
The recorded state is reachable from the state at initiation and reaches the state at completion, even if it never actually occurred.
Barrier snapshotting
The stream-processing descendant of the marker algorithm, in which barriers injected into data streams trigger operator checkpoints.

Module 3: Replication and Consistency

Where copies of the data come from, what clients are allowed to observe, and what the CAP theorem actually says.

Replication Strategies and the Anomalies They Produce

  • Compare single-leader, multi-leader and leaderless replication and say which failure each is designed for.
  • Name the three replication-lag anomalies precisely and give the session guarantee that fixes each.
  • Explain the failover trade-off between lost writes and unavailability, and what a fencing token protects.

On 21 October 2018 a scheduled maintenance job took a network link between two GitHub data centres offline for 43 seconds. The automated failover system did what it was built to do and promoted a database in the other region. Both regions had, for those 43 seconds, been accepting writes. The site was degraded for more than 24 hours afterwards, not because anything was broken, but because reconciling the two divergent write histories was slow, careful work that could not be automated. Their published analysis is worth reading in full, because everything in this lesson is visible in it.

Three reasons to replicate, and they conflict

Copies of your data exist for three different reasons, and the design that serves one poorly serves another.

  • Availability. Keep serving when a machine dies. Requires that another copy can take over without human intervention.
  • Latency. Put data near users. Requires copies far apart, which makes keeping them in step expensive.
  • Throughput. Spread reads across copies. Requires that reading a copy is acceptable, which is a consistency decision.

The tension is immediate: the cheapest way to serve a read from Sydney is to read a Sydney replica, and the cheapest way to guarantee that read is current is to ask the leader in Virginia.

Single-leader replication

One replica accepts writes; the others follow, applying the leader's changes in order. This is what PostgreSQL, MySQL, most managed relational databases, MongoDB replica sets and Kafka partitions all do, and it is the right default because it makes write conflicts impossible: all writes pass through one place, so they have an order.

The interesting parameter is when the leader acknowledges a write.

ModeLeader acknowledgesOn leader failureCost
AsynchronousImmediately after its own writeRecently acknowledged writes can be lostFast; durability is a guess
Semi-synchronousAfter at least one follower confirmsSafe if that follower survivesOne extra round trip
Synchronous to allAfter every follower confirmsNo lossAny single slow follower stalls all writes

The third row is why nobody uses it: replicating synchronously to n followers makes your availability worse than a single machine's, since any one of them failing blocks writes. Semi-synchronous replication to a quorum is the usual compromise, and it is exactly what consensus protocols formalize in Module 4.

The point: asynchronous replication does not merely risk losing writes; it means an acknowledgement is a promise the system may not be able to keep.

Failover, and the two ways it hurts

When the leader stops responding, something must decide whether it is dead and promote a follower. Every part of that sentence is a hazard.

  • Deciding it is dead. There is no way to know, only a timeout. Too long and you are down; too short and you promote during a hiccup, which is what a 43-second partition triggers.
  • Choosing a successor. The most up-to-date follower is the obvious choice, and determining which one that is requires agreement, which is itself the problem you were trying to avoid.
  • Discarding writes. If the old leader had acknowledged writes the new leader never received, they are gone. When the old leader returns and its unreplicated writes conflict with new ones, the usual production practice is to discard them, and that is a silent data loss decision made by a script.
  • Split brain. If the old leader does not know it was demoted, two nodes accept writes. This is the GitHub scenario and the one that costs days rather than minutes.

The defence against the last one is a fencing token: every leadership term carries a monotonically increasing number, and the storage layer rejects any write carrying a number below the highest it has seen. A demoted leader that wakes up from a garbage collection pause, still convinced it is in charge, is then refused by the storage itself rather than by its own good manners. Fencing is the mechanism that makes lease-based leadership safe, and a system without it is relying on the assumption that a paused process notices it was paused.

Multi-leader and leaderless, in one paragraph each

Multi-leader. Several nodes accept writes and replicate to each other. This is what you need when clients are geographically split and each region must accept writes locally, when clients work offline, or when many users edit the same document. The price is unavoidable: two leaders can accept conflicting writes to the same item, so a conflict resolution policy is mandatory. Last write wins loses data, as Module 1 showed; application-level merge is honest; conflict-free replicated data types make merging automatic for certain data shapes.

Leaderless. Clients write to several replicas directly and read from several, using overlapping quorums to obtain recency. Amazon's Dynamo popularized this and Cassandra and Riak followed. It removes failover entirely, since there is nothing to promote, and it moves the entire consistency question into the choice of how many replicas to contact. A later lesson works through the quorum arithmetic.

The three anomalies of replication lag

With any asynchronous scheme, a follower is behind. Three specific user-visible failures follow, and each has a name and a standard fix. These are worth memorizing as a checklist, because they are the bugs users report as "the site is broken" and engineers cannot reproduce.

AnomalyWhat the user seesGuarantee that fixes itImplementation
Reading your own stale writePosts a comment, refreshes, the comment is goneRead-your-writesRead from the leader for data the user may have modified, or route the user to one replica, or carry a write timestamp and wait for the replica to reach it
Moving backwards in timeSees a comment, refreshes, it disappears, because two reads hit replicas at different lagsMonotonic readsPin each user to one replica, chosen by a hash of the user identifier
Effect before causeSees an answer to a question that has not appeared yetConsistent prefix readsEnsure causally related writes are applied in order, either by keeping them in one partition or by tracking causal dependencies

Notice that all three are satisfied automatically by reading from the leader, and that all three are cheap to violate the moment somebody adds a read replica for performance. That is precisely how they arrive in production: as a consequence of an optimization made by someone who was not thinking about consistency.

In short: replication lag is not a performance issue that occasionally shows through. It is a change in the set of histories your users can observe, and the three anomalies are the vocabulary for saying which ones you have permitted.

What actually travels between replicas

The replication log is a design decision with consequences that surface years later.

MethodWhat is shippedProblem
Statement-basedThe SQL statementNondeterminism: a statement using the current time, a random value, or an auto-increment behaves differently on each replica
Write-ahead log shippingThe storage engine's physical log recordsCouples replicas to the exact storage format, so upgrades usually require downtime
Logical (row-based)The rows that changed, described independently of the storage engineLarger; but decoupled, which allows rolling upgrades and external consumers
Trigger-basedApplication-level records written by database triggersSlower and more fragile, but the most flexible

The third row is why change data capture became an architecture rather than a feature: once the log is logical, other systems can consume it, and the database's replication stream becomes the event stream that feeds search indexes, caches and analytics.

Common misconceptions

  • "Adding read replicas is a pure performance win." It changes what clients can observe. Every one of the three anomalies above becomes possible the moment a read can be served by a lagging copy.
  • "Synchronous replication to all followers is the safe choice." It makes availability worse than one machine, since any follower's failure blocks writes. Quorum-based acknowledgement is the safe choice.
  • "Failover is automatic, so leader failure is handled." Failover is a decision made under uncertainty, and it can lose acknowledged writes or produce two leaders. Automating it correctly is the subject of Module 4.
  • "A demoted leader will stop writing once it notices." It may be paused and notice nothing for a minute. Only a fencing token enforced by the storage layer prevents it from acting on stale authority.
  • "Multi-leader replication avoids the single point of failure." It does, and it buys a conflict resolution problem that has no general solution. Choose it when you need local writes or offline operation, not to dodge failover.
  • "Statement-based replication is simplest and therefore safest." It is simplest and least safe, because any nondeterministic function makes the replicas silently diverge.

Summing up

  • Replication serves availability, latency and read throughput, and those three goals pull the design in different directions.
  • Single-leader replication makes write conflicts impossible by construction; the real choice is when the leader acknowledges, and asynchronous acknowledgement means a confirmed write can still be lost.
  • Failover requires deciding a leader is dead without being able to know, choosing the most current follower, possibly discarding writes, and preventing two leaders.
  • Fencing tokens, enforced by the storage layer, are what stop a paused former leader from acting on authority it no longer has.
  • The three lag anomalies are stale reads of your own writes, non-monotonic reads, and effects preceding causes, fixed by read-your-writes, monotonic reads, and consistent prefix guarantees.
  • Logical replication logs decouple replicas from the storage format and turn the replication stream into a reusable change data capture feed.

Each of those anomalies is really a statement about which orderings a client may observe. The next lesson makes that idea precise, and gives the strongest and weakest versions of it their real definitions.

Sources

  1. GitHub. (2018). October 21 post-incident analysis. The GitHub Blog. github.blog
  2. Kleppmann, M. (2017). Replication. In Designing data-intensive applications (ch. 5). O'Reilly Media.
  3. Tanenbaum, A. S., & van Steen, M. (2017). Consistency and replication. In Distributed systems: principles and paradigms (3rd ed.). Pearson.
  4. Oracle. (n.d.). Semisynchronous replication. MySQL 8.0 Reference Manual. dev.mysql.com
  5. PostgreSQL Global Development Group. (n.d.). Log-shipping standby servers. PostgreSQL documentation. postgresql.org
  6. Wikipedia contributors. (n.d.). Split-brain (computing). en.wikipedia.org
  7. Wikipedia contributors. (n.d.). Change data capture. en.wikipedia.org
Key terms
Single-leader replication
One replica accepts all writes and others apply its log in order, which makes write conflicts impossible by construction.
Semi-synchronous replication
Acknowledging a write once at least one follower has confirmed it, trading one round trip for durability.
Split brain
Two nodes simultaneously believing they are the leader, both accepting writes, producing divergent histories.
Fencing token
A monotonically increasing leadership number checked by the storage layer, which rejects writes from a superseded leader.
Read-your-writes
A session guarantee that a client always sees its own prior writes, though not necessarily other clients' writes.
Monotonic reads
A session guarantee that successive reads never move backwards in time, usually implemented by pinning a client to one replica.
Consistent prefix reads
A guarantee that causally ordered writes are observed in that order, so an effect never appears before its cause.
Change data capture
Consuming a database's logical replication stream as an event feed for other systems.

The Consistency Spectrum, Defined Precisely

  • State the definition of linearizability and test a history against it.
  • Distinguish linearizability, sequential consistency, causal consistency and eventual consistency, and say which are composable.
  • Explain why linearizability costs latency even without failures, and what causal consistency buys instead.

Alice and Bob are watching the same football match on two phones. Alice refreshes, sees the final score, and says it out loud. Bob refreshes a second later and his phone still shows the game in progress. Nothing is broken in any obvious sense: both phones queried a real server and got a real answer. What has been violated is a property with a precise definition, published by Maurice Herlihy and Jeannette Wing in 1990, and most of the confusion in this field comes from using the word "strong" where that definition belongs.

Linearizability

Definition. A history of operations is linearizable if each operation appears to take effect atomically at a single instant between its invocation and its response, and the resulting total order of those instants is consistent with real time: if operation A completes before operation B begins, A's instant comes first.

Read that twice, because two separate things are being asserted. The first is atomicity: no operation is observed half-done. The second is a real-time constraint: a completed operation is visible to everything that starts afterwards. That second clause is what rules out Bob's phone.

Test a history. Time runs left to right; a bar spans the interval between a client's request and its response.

A:  |--- write(x, 1) --------|
B:      |- read(x) -> 0 -|                     |- read(x) -> 0 -|
C:              |------ read(x) -> 1 ------|

Is this linearizable?  NO.
C's read returned 1, so the write's instant is before C's response.
B's second read starts AFTER C's read returned, so by the real-time
clause it must observe the write. Returning 0 is impossible.

Remove B's second read and the history IS linearizable: place the
write's instant between B's first read and C's read.

Note what is not required. B's first read may return 0 even though it overlaps the write, because concurrent operations may be ordered either way. Linearizability constrains non-overlapping operations absolutely and overlapping ones not at all.

What matters here: linearizability is a guarantee about single objects and about recency. It says the system behaves as though there were one copy of the data and operations happened one at a time.

Composability, and why it is the killer feature

Herlihy and Wing proved that linearizability is local, meaning composable: if every object in a system is individually linearizable, then the system as a whole is linearizable. Nothing else on this list has that property.

The practical consequence is large. You can build a linearizable register, a linearizable queue and a linearizable counter independently, and reason about a program using all three without re-deriving anything. Sequential consistency, by contrast, is not composable: two sequentially consistent objects can be combined into a system that is not sequentially consistent, which means local reasoning fails and you must analyze the whole program.

Sequential consistency

Definition, from Lamport in 1979: there is some total order on all operations such that every process's own operations appear in that order in the order the process issued them, and every process observes that same total order.

The difference from linearizability is exactly one clause: real time is not mentioned. Everyone agrees on an order, but that order need not match the order in which things actually happened.

real time:   A writes x=1 at 10:00:00
             B reads x    at 10:00:05, gets 0

Linearizable?   No: the read began after the write completed.
Sequentially consistent?  Yes: put B's read before A's write in the
                          agreed total order. Nobody's program order
                          is violated; the order simply does not
                          correspond to the clock.

That is why sequential consistency is the model for a multiprocessor's memory and rarely the one a distributed database advertises: a user with a phone in each hand is an external observer with a clock, and only linearizability accounts for them.

Causal consistency

Definition. Operations related by happens-before are observed by every process in that order; concurrent operations may be observed in different orders by different processes.

This is the direct application of Module 2. It forbids seeing a reply before its message, and it permits two unrelated posts to appear in either order to different readers, which is almost always fine.

What makes it important is a theoretical result rather than a practical one: Mahajan, Alvisi and Dahlin showed in 2011 that causal consistency is the strongest model that a system can provide while remaining available during a partition. Anything stronger requires coordination, and coordination requires a reachable quorum. So the line between causal and sequential consistency is not a matter of engineering effort. It is the boundary of what is achievable while still answering every request.

Eventual consistency, and what it does not say

Definition. If writes stop, all replicas eventually converge to the same value.

Now read it critically, because the definition is much weaker than the phrase suggests. It places no bound on when. It says nothing about what a client may read before convergence, so returning a value written a week ago, or any value at all, is permitted. And the antecedent, writes stopping, never happens in a live system. Eventual consistency is a liveness property with no safety content, which is why practitioners layer session guarantees on top of it, the read-your-writes and monotonic reads of the previous lesson.

Strong eventual consistency is the useful strengthening: any two replicas that have received the same set of updates are in the same state, with no conflict resolution required. That is what conflict-free replicated data types provide, by making the merge operation commutative, associative and idempotent, so that the order of delivery cannot matter.

Remember: eventual consistency without session guarantees does not promise anything a client can rely on within any particular request.

The map

ModelGuaranteesComposableAvailable during a partitionTypical cost per operation
LinearizabilityOne-copy behaviour with a real-time orderYesNoA round trip to a quorum
Sequential consistencyOne agreed order respecting each process's program orderNoNoCoordination, but reads or writes can be local
Causal consistencyCausally related operations ordered; concurrent ones freeNoYesMetadata to track dependencies
Strong eventual consistencySame updates delivered implies same stateDepends on the data typeYesMerge logic in the data type
Eventual consistencyConvergence if writes stopNoYesAlmost nothing

What linearizability costs when nothing is failing

It is tempting to think the price of linearizability appears only during partitions. Attiya and Welch proved otherwise in 1994. In a network where message delay is uncertain within a window u, any implementation of a linearizable shared register forces read operations to take at least about u/4 and writes at least about u/2. Sequential consistency admits implementations where reads are local, or where writes are local, but not both.

So the cost is a floor set by network uncertainty, paid on every operation, in the healthy case. This is the observation that later became the second half of PACELC: even when there is no partition, there is a trade-off between latency and consistency, and it is the one you pay for continuously.

Two axes people conflate

Linearizability and serializability are different properties on different axes, and mixing them up produces confident wrong statements.

LinearizabilitySerializability
ScopeSingle object, single operationMultiple objects, whole transactions
SaysOperations appear in a real-time-respecting total orderTransactions appear in some serial order
Real timeRequiredNot required
CombinedStrict serializability: transactions in a serial order that also respects real time. This is what Spanner provides and what "external consistency" names.

A database can be serializable and still return stale data, because nothing forces the serial order to match the clock. A key-value store can be linearizable and offer no transactions at all. Saying a system is "strongly consistent" distinguishes none of this, which is why the word is best avoided in a design document.

Common misconceptions

  • "Strong consistency means linearizable." The phrase is used for anything from linearizability to snapshot isolation. Name the model instead.
  • "Linearizability and serializability are the same thing." One is about recency of single operations, the other about the atomicity of multi-object transactions. Only their conjunction, strict serializability, gives both.
  • "Eventual consistency means the data is eventually correct." It means replicas converge if writes stop. It permits a client to read an arbitrarily old value at any moment before then.
  • "Consistency is only expensive during failures." Attiya and Welch's bound is a latency floor proportional to network uncertainty, paid on every operation in the healthy case.
  • "If each component is consistent, the system is." True for linearizability, which is composable, and false for every weaker model on the list.
  • "Causal consistency is a weak compromise." It is provably the strongest model compatible with remaining available during a partition, which makes it the natural target for systems that must answer every request.

Looking back

  • Linearizability requires each operation to take effect at an instant inside its interval, with the resulting order respecting real time, so a completed operation is visible to everything that begins later.
  • It is composable, which is the property that permits local reasoning and which no weaker model on the list has.
  • Sequential consistency drops only the real-time clause: everyone agrees on an order, but it need not match the clock.
  • Causal consistency orders only causally related operations and is the strongest model that remains available during a partition.
  • Eventual consistency promises convergence if writes stop and nothing else; strong eventual consistency, as provided by conflict-free replicated data types, adds that equal update sets imply equal state.
  • Linearizability costs latency proportional to network uncertainty even with no failures, and it is a different axis from serializability, whose conjunction with real time is strict serializability.

With the definitions in place, the next lesson can state the theorem that everyone cites about them, in the form in which it was actually proved.

Sources

  1. Herlihy, M. P., & Wing, J. M. (1990). Linearizability: a correctness condition for concurrent objects. ACM Transactions on Programming Languages and Systems, 12(3), 463-492. cs.brown.edu
  2. Lamport, L. (1979). How to make a multiprocessor computer that correctly executes multiprocess programs. IEEE Transactions on Computers, C-28(9), 690-691.
  3. Attiya, H., & Welch, J. L. (1994). Sequential consistency versus linearizability. ACM Transactions on Computer Systems, 12(2), 91-122.
  4. Mahajan, P., Alvisi, L., & Dahlin, M. (2011). Consistency, availability, and convergence (Technical Report TR-11-22). University of Texas at Austin.
  5. Bailis, P. (2014). Linearizability versus serializability. bailis.org
  6. Kingsbury, K. (n.d.). Consistency models. Jepsen. jepsen.io
  7. Vogels, W. (2008). Eventually consistent. All Things Distributed. allthingsdistributed.com
  8. Wikipedia contributors. (n.d.). Linearizability. en.wikipedia.org
Key terms
Linearizability
Each operation takes effect atomically at an instant within its interval, and the resulting order respects real time.
Linearization point
The instant, somewhere between invocation and response, at which an operation is deemed to have taken effect.
Composability
The property that a system of individually correct objects is itself correct; linearizability has it and weaker models do not.
Sequential consistency
A single agreed total order respecting each process's program order, with no real-time requirement.
Causal consistency
Ordering only causally related operations; provably the strongest model available during a partition.
Eventual consistency
Convergence of replicas if writes stop, with no bound on time and no constraint on intermediate reads.
Strong eventual consistency
Replicas that have delivered the same updates are in the same state, achieved by commutative merge functions.
Strict serializability
Serializable transactions in an order that also respects real time; the conjunction of the two axes.

CAP, As It Was Actually Proved

  • State the CAP theorem with the exact definitions Gilbert and Lynch used, and reproduce the proof.
  • Correct the standard misreadings, including pick-two, the meaning of availability, and the meaning of consistency.
  • Use PACELC to describe the trade-off that applies when the network is healthy.

On 19 July 2000 Eric Brewer gave a keynote at the Principles of Distributed Computing symposium in Portland and offered a conjecture about what a distributed data store could guarantee. Two years later Seth Gilbert and Nancy Lynch published a proof of a formal version of it in SIGACT News. In the twenty-odd years since, the conjecture's three-letter summary has been quoted more often than the paper has been read, and the summary is wrong in three specific ways. This lesson states the theorem, proves it in five lines, and then repairs the summary.

Start from the version that is wrong

The popular formulation is: consistency, availability, partition tolerance: pick two. Draw the triangle, choose an edge, done.

Here is why that sentence cannot be right. Partition tolerance is not a feature you select. A partition is an event the network performs regardless of your preferences: a switch fails, a cable is cut, a firewall rule is deployed, or a garbage collection pause makes a node unreachable for a minute. Declining partition tolerance is not a design choice, it is a prediction, and it is false. So the menu never had three items on it.

The definitions, which are where the content is

Gilbert and Lynch's proof concerns a read/write data object replicated over a network, with three properties defined precisely:

TermThe paper's definitionWhat people assume it means
ConsistencyLinearizability of a single read/write register: the execution is equivalent to one in which operations happen atomically in a real-time-respecting orderSomething like the C in ACID, or a vague notion of correctness
AvailabilityEvery request to a non-failing node eventually receives a non-error response. Note the word eventually: there is no time bound at allHigh uptime, low latency, a service-level objective
Partition toleranceThe network is permitted to lose arbitrarily many messages between nodesSurviving a data centre failure

Two of those rows are load-bearing. Because availability requires only an eventual response, a system that takes an hour to answer is available in this sense, and a system that returns an error immediately is not. And because consistency means linearizability of a single register, the theorem says nothing directly about transactions across objects, which is a much larger part of most designs.

The proof

Suppose an algorithm provides all three. Place two nodes, G1 and G2, and partition the network so that no message passes between them.

1. A client writes v1 to G1.
   G1 is non-failing, so by AVAILABILITY it must eventually respond.
   It cannot consult G2, so it responds having applied the write locally.

2. Afterwards, a client reads from G2.
   G2 is non-failing, so by AVAILABILITY it must eventually respond
   with a value, not an error.

3. By PARTITION TOLERANCE, no message from step 1 reached G2,
   so G2 has no knowledge of v1 and returns the old value.

4. The read began after the write completed, and returned a stale value.
   That violates LINEARIZABILITY.

Contradiction. No such algorithm exists.

That is the entire proof, and it is worth noticing how little machinery it needs. It is essentially the same argument as the previous lesson's history test, run on a partitioned pair.

The point: the theorem does not say a partition forces you to choose. It says that during a partition, a system that answers every request cannot also guarantee that answers are current. Outside a partition, the theorem says nothing at all.

The three misreadings

Misreading one: pick two of three. P is not selectable. The real statement is conditional: when a partition occurs, you must give up either availability or linearizability for the affected data. Labelling a system CP or AP describes its behaviour in that situation, not a permanent identity.

Misreading two: CA systems exist. A single-node database is trivially CA, because it has no network to partition. A distributed "CA" system is one whose designers assumed partitions do not happen, which means its behaviour during one is undefined rather than chosen. The honest reading of a CA label is: this system will do something unspecified when the network fails.

Misreading three: availability means uptime. The paper's availability is a total-function requirement with unbounded latency, which is neither what an operations team means by availability nor what a user experiences. A CP system is not "unavailable"; it is unavailable for some operations on some data during a partition, which for a well-designed system may be a tiny fraction of its traffic.

Brewer himself wrote a retrospective in 2012 saying the two-of-three formulation had been misleading, and proposing a better frame: detect that a partition has begun, enter an explicit partition mode with restricted operations, and run a compensation process when it heals. That is what mature systems actually do, and it is a design discipline rather than a label.

What CAP does not cover

Martin Kleppmann's 2015 critique makes the sharpest version of this point: the theorem is a precise statement about a narrow model, and its definitions do not match what practitioners mean by any of the three words. Linearizability is one point on the consistency spectrum, availability with no latency bound is not useful availability, and a single register is not a database.

The most consequential gap is the healthy case. Partitions are rare; the trade-off you pay for continuously is different, and Daniel Abadi named it PACELC:

if (Partition)  then choose Availability or Consistency
Else            choose Latency or Consistency

The second line is the one that governs everyday design, and the previous lesson proved it is not negotiable: Attiya and Welch's bound makes linearizable operations cost a latency proportional to network uncertainty even when nothing is failing. Classifying systems on both axes is far more informative than a CP or AP label:

SystemDuring a partitionNormallyReading
Dynamo-style stores such as Cassandra, with low quorum settingsAvailabilityLatencyPA/EL: answers fast, may answer with stale or conflicting data
SpannerConsistencyConsistencyPC/EC: pays commit-wait latency always, refuses rather than diverges
ZooKeeper, etcd, ConsulConsistencyConsistencyPC/EC: a minority partition stops serving writes entirely
A store configured with quorum reads and writesConsistency for the minority sideConsistencyThe configuration, not the product, determines the classification

The last row is the practical lesson. A single product can sit in several boxes depending on its settings, and often on a per-request basis: many stores let each individual read or write choose its consistency level. So the question "is this database CP or AP" usually has no answer, while "what does this specific operation do when a quorum is unreachable" always does.

So what?: replace the label with the operational question. For each critical operation, ask what happens when the node handling it cannot reach a majority, and make sure the answer is a decision somebody made.

Designing with the partition in mind

Brewer's partition-mode framing turns the theorem into a checklist:

  1. Detect. Decide what evidence counts as a partition, usually a failure to reach a quorum within a timeout.
  2. Restrict. Choose which operations remain available in that mode. Adding an item to a cart may continue; charging a card may not. This is a per-operation decision, and it belongs to the product, not the database.
  3. Record. Keep enough information about what was accepted during the partition to reconcile afterwards, which usually means versioned writes rather than overwrites.
  4. Compensate. When the partition heals, merge and, where the merge is not automatic, take a business action: notify a user, issue a refund, flag for review.

Written that way, CAP stops being a taxonomy and becomes what it always was: a proof that step 2 cannot be skipped.

Common misconceptions

  • "Pick two of three." Partition tolerance is not a choice. The theorem is conditional on a partition occurring, and it forces a choice only then.
  • "Our system is CA." Unless it runs on one node, that means its partition behaviour was never designed, only assumed away.
  • "The C in CAP is the C in ACID." ACID's consistency means the database respects declared invariants. CAP's consistency is linearizability. They are unrelated.
  • "CP systems are unavailable." They are unavailable for particular operations on particular data while a quorum is unreachable, which may be a negligible part of the workload.
  • "CAP explains why our database is slow." CAP concerns partitions. Everyday latency-versus-consistency is the second half of PACELC, and its cost is a proved floor rather than an implementation weakness.
  • "A database is either CP or AP." Most modern stores expose per-operation consistency levels, so the classification belongs to the request, not the product.

What to remember

  • Gilbert and Lynch proved that no replicated register can be simultaneously linearizable and available in a network that may lose arbitrarily many messages.
  • Their availability means an eventual non-error response with no latency bound, and their consistency means linearizability of a single register; both differ from ordinary usage.
  • The proof partitions two nodes, writes to one, reads from the other, and observes that the read must answer without knowledge of the write.
  • Partition tolerance is not selectable, so the real statement is a conditional choice made during a partition, and it can be made per operation and per data item.
  • Brewer's later framing is to detect the partition, restrict operations explicitly, record what was accepted, and compensate on recovery.
  • PACELC adds the trade-off that applies when nothing is broken: consistency costs latency continuously, which is the term that governs most designs.

CAP describes what is impossible when the network fails. The next module asks the harder question: given that failure, what agreement is still achievable, and at what cost?

Sources

  1. Gilbert, S., & Lynch, N. (2002). Brewer's conjecture and the feasibility of consistent, available, partition-tolerant web services. ACM SIGACT News, 33(2), 51-59. groups.csail.mit.edu
  2. Brewer, E. (2012). CAP twelve years later: How the rules have changed. IEEE Computer, 45(2), 23-29. infoq.com
  3. Kleppmann, M. (2015). A critique of the CAP theorem. arXiv. arxiv.org
  4. Kleppmann, M. (2015). Please stop calling databases CP or AP. martin.kleppmann.com
  5. Abadi, D. (2012). Consistency tradeoffs in modern distributed database system design: CAP is only part of the story. IEEE Computer, 45(2), 37-42.
  6. Wikipedia contributors. (n.d.). CAP theorem. en.wikipedia.org
  7. Wikipedia contributors. (n.d.). PACELC design principle. en.wikipedia.org
Key terms
CAP theorem
No replicated register can guarantee linearizability and availability simultaneously in a network that may lose arbitrarily many messages.
CAP availability
Every request to a non-failing node eventually receives a non-error response, with no bound on how long that takes.
CAP consistency
Linearizability of a single read/write register, which is unrelated to the C in ACID.
Partition mode
Brewer's framing in which a system detects a partition, restricts operations explicitly, records what it accepted, and compensates on recovery.
PACELC
If partitioned, choose availability or consistency; else, choose latency or consistency, which is the trade-off paid continuously.
Per-operation consistency level
The common facility letting each request choose its own guarantee, which is why a product cannot be labelled CP or AP as a whole.
Compensation
The business-level action taken after reconciling writes accepted on both sides of a partition, such as a refund or a notification.

Module 4: Agreement

The impossibility result that bounds the field, the two protocols that work around it, and how a group decides who is in charge.

FLP: What Impossibility Actually Forbids

  • State the consensus problem with its three required properties and the FLP hypotheses.
  • Follow the bivalence argument that produces a non-terminating execution.
  • List the four ways real systems escape the result, and identify which property each one relaxes.

In April 1985 the Journal of the ACM published nine pages by Michael Fischer, Nancy Lynch and Michael Paterson proving that a problem the field had been attacking for years has no solution. Not a hard solution: no solution. The paper won the first Dijkstra Prize sixteen years later, and it is now cited in roughly equal measure by people who understand it and people using it to justify not trying. The distinction between those two groups is entirely a matter of reading the hypotheses.

The problem, stated exactly

Each process starts with an input value and must decide on an output. An algorithm solves consensus if every execution satisfies:

  • Agreement. No two correct processes decide different values.
  • Validity. A decided value was proposed by some process. (This rules out the trivial algorithm that always decides 0.)
  • Termination. Every correct process eventually decides.

The first two are safety properties; the third is liveness. Keep that split in mind, because the theorem lands entirely on the third.

The theorem, with its hypotheses on the table

FLP. In an asynchronous message-passing system in which at most one process may crash, there is no deterministic algorithm that solves consensus in every execution.

Every emphasized word is a hypothesis, and each one is an escape hatch. Note especially how weak the failure assumption is: one crash, not one Byzantine liar, and not a partition. The result also holds even when messages are never lost and are eventually delivered. This is not a statement about hostile conditions; it is a statement about the total absence of timing information.

What matters here: FLP forbids guaranteed termination. It does not forbid agreement, it does not forbid validity, and it does not say a real system will fail to decide.

The intuition, before the proof

An asynchronous system cannot distinguish a crashed process from a slow one. So consider an algorithm that has heard from n - 1 of n processes and is waiting for the last.

  • If it waits, and that process has crashed, it waits forever. Termination fails.
  • If it proceeds without the last process, then the message might arrive a moment later and carry information that would have changed the decision. To be safe, the algorithm must be prepared to be interrupted at exactly the wrong instant.

FLP formalizes the second horn: for any algorithm, an adversarial scheduler can always find a moment at which delivering or delaying one message keeps the system undecided, and it can do so forever.

The proof, in its two lemmas

A configuration is the state of every process plus the set of messages in flight. Call a configuration 0-valent if every reachable decision from it is 0, 1-valent if every reachable decision is 1, and bivalent if both outcomes are still reachable. A bivalent configuration is one where the outcome is genuinely undetermined.

Lemma 1: some initial configuration is bivalent. Line up the initial configurations from all-zeros to all-ones, changing one process's input at a time. By validity, the first is 0-valent and the last is 1-valent, so somewhere along the line two adjacent configurations differ in decision. They differ in exactly one process's input, say process p. Now let p crash immediately in both. The remaining processes see identical executions, so they must decide the same value, which contradicts the two configurations being 0-valent and 1-valent respectively. Therefore at least one of the adjacent pair was bivalent all along.

Lemma 2: from any bivalent configuration, some step leads to another bivalent configuration. Take a bivalent C and any message m that is in flight. Suppose, for contradiction, that delivering m after any schedule from C always yields a univalent configuration. Then there exist two schedules from C, one reaching a 0-valent configuration and one reaching a 1-valent configuration, and by taking the earliest point at which they diverge you obtain a critical step by some process p at which the outcome is fixed. Now let p crash immediately after that step. The other processes must still terminate, so they decide some value, and by replaying the two schedules with p absent you find the same execution leading to two different decisions. Contradiction.

Conclusion. Start at a bivalent initial configuration, and by Lemma 2 the scheduler can always take a step that leaves the system bivalent. Doing this forever, while still delivering every message eventually, produces an infinite execution in which no process ever decides. Termination fails.

Two features of that argument are worth naming. The adversary is only allowed to delay messages, not to lose or corrupt them. And the execution it constructs is fair in the sense that every message is eventually delivered. The impossibility does not require anything to go wrong.

The four escapes, and what each gives up

EscapeHypothesis relaxedWhat you getUsed by
RandomizationDeterminismTermination with probability 1, in expected constant or logarithmic roundsBen-Or's protocol and its descendants; Byzantine agreement in some blockchains
Partial synchronyFull asynchronyAlways safe; terminates once the network is stablePaxos, Raft, Zab, Viewstamped Replication
Failure detectorsFull asynchrony, packaged as an oracleConsensus with an eventually accurate detector; Omega is the weakest sufficient oneThe theory behind every heartbeat and timeout in production
Weakening the problemTermination, or agreementEventual consistency, gossip protocols, conflict-free data typesDynamo-style stores

The second row is the one that matters most in practice, and it is worth stating carefully what it delivers. Paxos and Raft are always safe: they never allow two processes to decide differently, no matter how badly the network behaves, and this is unconditional. They are conditionally live: they are guaranteed to reach a decision only during periods when messages are delivered within some bound. During a partition or a leader election storm they may make no progress at all, and that is not a bug. It is FLP, showing up exactly where the theory says it must.

Remember: every timeout in a consensus implementation is the place where the algorithm declines to be purely asynchronous. That is not a hack around the theory; it is the theory, applied.

What the result is often misused to claim

Three arguments you will meet, each wrong in a specific way.

  1. "Consensus is impossible, so we use eventual consistency." Consensus is achievable in the model real networks live in. The honest version of this sentence is that consensus costs a round trip to a quorum and stops during partitions, and for this workload that price is not worth paying, which is a legitimate engineering decision and a different sentence.
  2. "FLP means distributed databases can lose data." FLP concerns liveness. A protocol that stops rather than deciding wrongly has not lost anything.
  3. "Since termination is impossible, we cannot rely on Raft." Raft terminates in every execution where the network eventually behaves, which is every execution anyone has observed for more than a few seconds. The theorem rules out a guarantee, not the behaviour.

Consensus is the same problem in several costumes

A last observation that makes the impossibility feel broader than it looks. Several problems are equivalent to consensus, in the sense that a solution to one gives a solution to the others:

  • Atomic broadcast, delivering messages to all processes in the same order.
  • Leader election that guarantees a unique leader.
  • Atomic commit across participants, the subject of a later lesson.
  • Linearizable read-write registers with compare-and-swap, by Herlihy's consensus hierarchy.

So FLP is not a statement about one algorithm's problem. Any time a system needs all participants to agree on a single unambiguous fact, the same impossibility applies and the same four escapes are the only ones available.

Common misconceptions

  • "FLP says consensus can never be achieved." It says no deterministic algorithm guarantees termination in a fully asynchronous system with one possible crash. Real systems change one of those hypotheses.
  • "FLP is about network partitions." The proof needs no lost messages at all. Delay alone is enough.
  • "The result assumes many failures." One crash. Adding more does not make it worse, and removing all of them makes consensus easy.
  • "Raft violates FLP." Raft assumes partial synchrony. Under full asynchrony its randomized election timeouts still leave it unable to guarantee a leader is ever elected, exactly as required.
  • "Randomization is cheating." It changes the guarantee from certain termination to termination with probability 1, which is a genuinely different and perfectly respectable statement.
  • "FLP and CAP say the same thing." They have different hypotheses and different conclusions. FLP concerns termination under asynchrony with no message loss; CAP concerns linearizability and availability under message loss.

The short version

  • Consensus requires agreement, validity and termination; the first two are safety and the third is liveness.
  • FLP proves that no deterministic algorithm guarantees all three in an asynchronous system where a single process may crash, even with reliable message delivery.
  • The proof shows an initial configuration is bivalent and that a scheduler can always keep it bivalent, producing an infinite undecided execution.
  • The four escapes are randomization, partial synchrony, failure detectors, and weakening the problem itself.
  • Paxos and Raft are unconditionally safe and conditionally live, so their stalls during network trouble are the theorem appearing exactly where it must.
  • Atomic broadcast, unique leader election and atomic commit are all equivalent to consensus, so the same limits and the same escapes apply to each.

The next lesson takes the second escape and follows it all the way to a working protocol.

Sources

  1. Fischer, M. J., Lynch, N. A., & Paterson, M. S. (1985). Impossibility of distributed consensus with one faulty process. Journal of the ACM, 32(2), 374-382. groups.csail.mit.edu
  2. Ben-Or, M. (1983). Another advantage of free choice: completely asynchronous agreement protocols. In Proceedings of the 2nd ACM Symposium on Principles of Distributed Computing, 27-30.
  3. Dwork, C., Lynch, N., & Stockmeyer, L. (1988). Consensus in the presence of partial synchrony. Journal of the ACM, 35(2), 288-323.
  4. Chandra, T. D., Hadzilacos, V., & Toueg, S. (1996). The weakest failure detector for solving consensus. Journal of the ACM, 43(4), 685-722.
  5. Herlihy, M. (1991). Wait-free synchronization. ACM Transactions on Programming Languages and Systems, 13(1), 124-149.
  6. Wikipedia contributors. (n.d.). Consensus (computer science), including the FLP result and its workarounds. en.wikipedia.org
  7. Wikipedia contributors. (n.d.). Atomic broadcast and its equivalence with consensus. en.wikipedia.org
Key terms
Consensus
Deciding one value satisfying agreement, validity and termination, where the first two are safety and the third is liveness.
Configuration
The combined state of all processes plus the messages currently in flight.
Bivalent configuration
One from which both decision values are still reachable, so the outcome is genuinely undetermined.
Critical step
The event at which a bivalent configuration becomes univalent; FLP shows the scheduler can always postpone it.
Conditionally live
Guaranteed to terminate only during periods when the network satisfies a timing bound; the property Paxos and Raft actually have.
Randomized consensus
Protocols that escape FLP by using coin flips, achieving termination with probability 1 rather than certainty.
Consensus equivalence
Atomic broadcast, unique leader election and atomic commit are all reducible to consensus and to each other.

Paxos, Developed Carefully

  • Derive the two Paxos phases from the safety requirement rather than memorizing them.
  • Trace a run with competing proposers and show why the value already chosen is carried forward.
  • Explain why Paxos can livelock, and how multi-Paxos reduces the steady state to one round trip.

Leslie Lamport submitted a paper describing the Paxos algorithm in 1990. He wrote it as the report of an archaeological dig on a Greek island, complete with fictional scholars and a parliament that legislated while its members wandered in and out of the chamber. Almost nobody read it, and it was not published until 1998, eight years later. In 2001 he wrote a five-page version with no archaeology in it, and that is the one everyone learns from. This lesson follows the second one, deriving the protocol from what safety requires rather than presenting it as a set of rules to memorize.

The problem, reduced to its smallest form

Forget replicated logs for now. Solve exactly one thing: a group of processes must choose a single value, and once chosen, every process that learns a value must learn the same one. Three roles, which in practice are usually the same machines wearing different hats:

  • Proposers suggest values.
  • Acceptors vote. A value is chosen when a majority of acceptors have accepted it.
  • Learners find out what was chosen.

Assume the crash-recovery model, asynchronous messages, and 2f + 1 acceptors so that f may fail. Majorities are the whole trick: any two majorities of the same set share at least one member, so information accepted by one majority cannot be invisible to the next.

Deriving the algorithm

First requirement. If a single proposer proposes a single value and nothing fails, some value must be chosen. So an acceptor must accept the first proposal it receives.

But if acceptors always accept the first thing they see, three acceptors receiving three different proposals choose nothing, and worse, different values could each collect a majority in different runs. So proposals need to be ordered, and acceptors need permission to accept more than one. Give every proposal a unique ballot number, and let acceptors accept several proposals over time. Uniqueness is usually arranged by making the number a pair of a counter and a server identifier.

Second requirement, the safety property. If a proposal with value v has been chosen, then every higher-numbered proposal that is ever issued must also have value v.

That is stronger than it needs to be, and deliberately so: it is easier to enforce a rule about what proposers issue than about what gets chosen. So a proposer about to issue ballot n must first find out whether any value might already have been chosen by a lower ballot. It cannot ask every acceptor, since some may be down. It asks a majority, and the intersection property does the rest: if a value was chosen by some majority, any majority the proposer contacts contains an acceptor that accepted it.

The core of it: Paxos is one idea. Before proposing, ask a majority what they have already accepted, and if any of them has accepted anything, propose that instead of your own value.

The protocol

PHASE 1  (prepare)
  proposer: choose a ballot number n, higher than any it has used.
            send PREPARE(n) to at least a majority of acceptors.

  acceptor: on PREPARE(n):
      if n > the highest ballot I have already promised:
          record the promise DURABLY, then
          reply PROMISE(n, highestAccepted)  where highestAccepted is the
              (ballot, value) pair of the highest-numbered proposal
              I have accepted, or nothing if I have accepted none.
          I will now refuse any ACCEPT with a ballot below n.
      else: ignore, or reply with a rejection carrying my highest promise.

PHASE 2  (accept)
  proposer: if PROMISEs arrive from a majority:
      if any of them reported an accepted proposal:
          v = the value from the HIGHEST-numbered reported proposal
      else:
          v = my own proposed value
      send ACCEPT(n, v) to at least a majority.

  acceptor: on ACCEPT(n, v):
      if I have not promised to a ballot higher than n:
          record (n, v) DURABLY, then reply ACCEPTED(n, v).

LEARNING
  a value is CHOSEN as soon as a majority has accepted the same (n, v).
  learners find out by listening to acceptors, or through the proposer.

The word durably appears twice, and it is not decoration. In the crash-recovery model an acceptor that promises, crashes, restarts having forgotten the promise, and then accepts an older ballot has broken the algorithm. Every real Paxos implementation writes and flushes before replying, and that flush is a significant part of its latency.

Tracing a run

Three acceptors A1, A2, A3. Proposer P1 wants "red"; proposer P2 wants "blue".

1.  P1 sends PREPARE(1) to A1, A2, A3.
    All three have promised nothing, so all reply PROMISE(1, none).

2.  P1 has a majority and no reported values, so it uses its own:
    P1 sends ACCEPT(1, red).
    A1 and A2 accept and record (1, red). The message to A3 is delayed.

    "red" is now CHOSEN: a majority has accepted it. Nobody knows this yet.

3.  P2 starts, unaware of any of the above, and sends PREPARE(2)
    to A2 and A3.
    A2 replies PROMISE(2, (1, red))  -- it has accepted (1, red).
    A3 replies PROMISE(2, none).

4.  P2 has a majority of promises, and one of them reported a value.
    The rule forces it: v = red, the value of the highest-numbered
    reported proposal. P2 must abandon "blue".
    P2 sends ACCEPT(2, red).

5.  All acceptors accept (2, red). Same value, higher ballot. Safe.

Step 3 is the load-bearing one. P2 chose to contact A2 and A3, but it could have contacted any two of the three, and every such pair includes A1 or A2, both of which accepted "red". There is no majority that misses the previous majority entirely. That is the argument, and it is the whole reason Paxos is correct.

Now change one detail. Suppose in step 2 only A1 had accepted before P1 crashed. Then "red" was not chosen, since one acceptor is not a majority. If P2 then contacts A2 and A3, neither reports a value, so P2 is free to propose "blue", and "blue" may be chosen. This is correct behaviour: no value had been chosen, so no promise was broken. Paxos guarantees that once a value is chosen it never changes; it guarantees nothing about which value that will be.

Where liveness goes

Two proposers can prevent each other from making progress indefinitely.

P1: PREPARE(1)   promised
P2: PREPARE(2)   promised, invalidating P1's ballot
P1: ACCEPT(1)    refused; P1 retries with PREPARE(3)
P2: ACCEPT(2)    refused; P2 retries with PREPARE(4)
P1: ACCEPT(3)    refused; P1 retries with PREPARE(5)
...forever, with no value ever chosen

Nothing is unsafe here; nothing is decided either. This is FLP arriving on schedule: the algorithm is always safe and cannot guarantee termination. The standard remedy is to elect a distinguished proposer and have everyone else stand down, with randomized backoff so that two candidates do not keep colliding. Note the circularity, which is the honest situation: electing a unique leader is itself consensus, so the leader election need only be a good heuristic. If it elects two leaders, safety survives and progress stalls; if it elects one, progress resumes.

Key idea: the leader in Paxos is a performance optimization and a liveness heuristic, never a correctness requirement. Two leaders make it slow, not wrong.

From one value to a log: multi-Paxos

A replicated service needs a sequence of decisions, not one. Run an independent instance of the algorithm for each log position, and you have multi-Paxos. As written, every entry costs two round trips, which is unacceptable.

The optimization is to notice that Phase 1 does not mention the value. A stable leader can run Phase 1 once for all future instances, claiming ballot n for every position from here on, and thereafter commit each new entry with a single Phase 2 round trip.

Basic Paxos, per decisionMulti-Paxos, steady state
Round tripsTwo (prepare, then accept)One (accept only)
Durable writesTwo per acceptorOne per acceptor
When Phase 1 runsEvery decisionOnly when leadership changes

A new leader's first act is a Phase 1 covering all positions, whose replies tell it which entries may already have been chosen; it must re-propose those before proposing anything new. That recovery step is where most of multi-Paxos's real complexity lives, and it is the part the original papers leave to the reader.

Why Paxos has a reputation

Three separate difficulties get conflated. The single-decree algorithm above is genuinely simple. Turning it into a replicated log requires decisions the papers do not make for you: how to handle gaps in the log, when to garbage collect, how to change the membership of the acceptor set, how to serve reads without a full round. And the presentations are unusually indirect, the 1998 paper by design and the 2001 rewrite by compression.

The engineering gap is documented from the inside: Google's team building Chubby wrote a paper in 2007 describing what it took to get from the published algorithm to a production system, and their answer involved substantial machinery the algorithm never mentions. That gap is exactly what motivated the protocol in the next lesson.

Common misconceptions

  • "Paxos requires a leader." It requires nothing of the sort for safety. A leader is how it makes progress in practice.
  • "A proposer proposes its own value." Only when no acceptor in its promise majority reports a previously accepted proposal. Otherwise it must adopt the highest-numbered reported value.
  • "Ballot numbers are timestamps." They must be unique per proposer and increasing, which is why implementations use a counter paired with a server identifier. Wall-clock timestamps collide and go backwards.
  • "Once a majority has accepted, everyone knows." A value can be chosen without any process being aware of it. Learning is a separate step, and this is why a new leader must run Phase 1 before assuming a position is empty.
  • "An acceptor can reply first and persist later." That breaks the algorithm in the crash-recovery model, because a restarted acceptor may contradict a promise it already gave.
  • "Paxos guarantees the first proposed value wins." It guarantees only that a chosen value never changes. Which value is chosen depends on timing.

Where this leaves us

  • Paxos chooses one value among proposers, with a value counting as chosen once a majority of acceptors have accepted it.
  • Its correctness rests entirely on majority intersection: any majority a new proposer contacts includes an acceptor from any earlier accepting majority.
  • Phase 1 asks a majority to promise and to report what they have accepted; Phase 2 proposes the highest-numbered reported value, or the proposer's own if none was reported.
  • Promises and accepts must be written durably before replying, because a restarted acceptor with amnesia can violate safety.
  • Competing proposers can livelock indefinitely without violating safety, which is FLP appearing exactly where it must; a distinguished proposer with randomized backoff is the standard remedy.
  • Multi-Paxos runs Phase 1 once per leadership term and one Phase 2 per log entry, reducing the steady state to a single round trip.

Everything above is correct and, by common consent, hard to implement from the papers. The next lesson covers a protocol designed with that complaint as its explicit starting point.

Sources

  1. Lamport, L. (1998). The part-time parliament. ACM Transactions on Computer Systems, 16(2), 133-169. lamport.azurewebsites.net
  2. Lamport, L. (2001). Paxos made simple. ACM SIGACT News, 32(4), 51-58. lamport.azurewebsites.net
  3. Chandra, T. D., Griesemer, R., & Redstone, J. (2007). Paxos made live: an engineering perspective. In Proceedings of the 26th ACM Symposium on Principles of Distributed Computing, 398-407.
  4. van Renesse, R., & Altinbuken, D. (2015). Paxos made moderately complex. ACM Computing Surveys, 47(3), article 42.
  5. Microsoft Research. (n.d.). Paxos made simple: publication record. microsoft.com
  6. Wikipedia contributors. (n.d.). Paxos (computer science). en.wikipedia.org
Key terms
Acceptor
A process that votes on proposals; a value is chosen when a majority of acceptors have accepted it.
Ballot number
A unique, increasing proposal identifier, usually a counter paired with a server identifier so that two proposers never collide.
Majority intersection
The property that any two majorities of the same set share a member, which is what carries an accepted value forward.
Prepare phase
Asking a majority to promise to reject lower ballots and to report the highest proposal each has accepted.
Accept phase
Proposing the highest-numbered reported value, or the proposer's own if none was reported, to a majority.
Chosen value
One accepted by a majority; it may be chosen without any process yet knowing it, which is why learning is a separate step.
Distinguished proposer
A leader elected to avoid duelling proposals; a liveness heuristic that is never required for safety.
Multi-Paxos
Running one instance per log position, with a stable leader executing the prepare phase once per term.

Raft: The Same Idea, Made Teachable

  • Trace a Raft leader election, including the log-completeness restriction on voting.
  • Explain the log matching property and how the consistency check enforces it.
  • State why a leader may not commit an entry from a previous term by replica count alone.

Diego Ongaro and John Ousterhout did something unusual for a distributed systems paper: they ran a controlled experiment on their readers. A few dozen graduate students at two universities were taught Paxos and Raft in equivalent video lectures and then quizzed on both. The students scored substantially higher on Raft. The paper's claim is not that Raft is more capable than Paxos, because it is not; the claim is that a protocol people can implement correctly from the paper is worth more than an equally powerful one they cannot.

The design decision that changes everything

Raft imposes a strong leader. In Paxos any proposer may propose at any time, and a new leader must reconstruct what earlier proposers may have done. In Raft:

  • All client requests go to the leader.
  • Log entries flow only from leader to followers, never the reverse.
  • A follower whose log disagrees with the leader's is overwritten to match.

That third rule is the one that buys the simplicity. There is never a merge, never a reconciliation of two divergent logs. The leader's log is authoritative, and the entire safety argument becomes: make sure that whoever gets elected already has every committed entry.

Raft splits the problem into three pieces that can be understood separately, and that decomposition is the paper's real contribution: leader election, log replication, and the safety restriction that ties them together.

Terms, which are ballot numbers with a better name

Time is divided into terms, numbered consecutively. Each term begins with an election and has at most one leader; some terms have none, if the election fails. Every message carries its sender's term, and two rules govern them:

if a server sees a term GREATER than its own:
    it updates its term and reverts to follower
if a server receives a message with a term LESS than its own:
    it rejects the message

Those two lines are Raft's entire mechanism for handling a stale leader that wakes from a pause. It sends an AppendEntries, learns from the rejection that the term has moved on, and steps down. That is the same job a fencing token does in the replication lesson, built into the protocol.

Leader election

Every follower runs an election timeout, chosen randomly from a range such as 150 to 300 milliseconds, and resets it whenever it hears from a current leader. When the timeout fires:

follower becomes CANDIDATE:
    increment my term
    vote for myself
    send RequestVote(term, lastLogIndex, lastLogTerm) to all servers

a server grants its vote if ALL of:
    the candidate's term is at least my term
    I have not already voted in this term
    the candidate's log is AT LEAST AS UP TO DATE as mine

candidate becomes LEADER on receiving votes from a majority
candidate reverts to FOLLOWER if it learns of a current leader
             or of a higher term
if the timeout fires again with no result: start a new term and retry

The randomization is doing real work. With fixed timeouts, servers would time out together, split the vote, and repeat indefinitely, which is Paxos's duelling-proposers livelock in another form. Randomized timeouts mean one server almost always wakes first and collects a majority before the others start. This is the same escape from FLP as before, appearing as a line in a configuration file.

The third voting condition is the safety restriction, and it deserves its own definition. A log is at least as up to date as another if its last entry has a higher term, or the same term and an index at least as large. Since a candidate needs a majority of votes, and any committed entry is on a majority, the two majorities intersect: at least one voter holds every committed entry, and it will refuse a candidate missing them. Therefore any elected leader already contains every committed entry, and no reconstruction phase is needed.

Why this matters: Paxos handles a new leader by having it discover and re-propose what may already have been chosen. Raft handles it by refusing to elect a leader that would need to.

Log replication and the consistency check

Each entry holds a command, its index, and the term in which the leader created it. The leader sends AppendEntries containing the new entries plus the index and term of the entry immediately preceding them.

AppendEntries(term, prevLogIndex, prevLogTerm, entries[], leaderCommit)

follower:
    reject if term < my term
    reject if my log has no entry at prevLogIndex with term prevLogTerm
    otherwise: delete any conflicting suffix, append the new entries,
               and advance my commit index up to leaderCommit

The rejection path is what repairs divergent logs. The leader keeps a nextIndex per follower, decrements it on rejection, and retries, walking backwards until it finds the last point of agreement, then ships everything after it. That loop is simple to write and it is the only reconciliation mechanism in the protocol.

From the check follows the Log Matching Property: if two logs contain an entry with the same index and term, then those logs are identical in every entry up to that index. The proof is an induction: the leader creates at most one entry per index per term, so index and term identify an entry uniquely; and the consistency check means a follower only appends after the previous entry matched, which by induction means everything before it matched too. One field in one message gives an invariant over entire logs.

An entry is committed once the leader has replicated it to a majority. Committed entries are applied to the state machine in index order, and the leader tells followers the commit index in subsequent messages.

The rule that catches everyone

Here is the subtle part, and it is the one place where Raft's simplicity has a sharp edge. A new leader may find entries from previous terms sitting on a majority of servers. It is tempting to conclude they are committed and apply them. That is wrong, and the paper's Figure 8 is the counterexample.

The dangerous sequence, in outline:

term 2: leader S1 replicates entry E to S1 and S2, then crashes.
        E is on 2 of 5 servers. Not committed.
term 3: S5 is elected with votes from S3, S4, S5, and appends its own
        entry at the same index, then crashes before replicating.
term 4: S1 is elected again, and copies E to S3, so E is now on a
        MAJORITY: S1, S2, S3.
        If S1 commits E on that basis and then crashes...
term 5: S5 can still be elected (its log is up to date by the last-term
        rule with respect to S2, S3, S4), and it will overwrite E
        on every follower.

A committed entry was overwritten. Safety broken.

The fix is a single restriction: a leader may only commit an entry from its own term by counting replicas. Entries from earlier terms become committed indirectly, when a later entry from the current term is committed above them, which by the log matching property drags everything before it along. Many implementations arrange this by having a new leader immediately append a no-op entry in its own term.

Worth holding on to: the trap is not that Raft is subtle everywhere. It is that consensus has exactly a few subtle places, and moving them around does not remove them. Raft concentrated them into this rule and the voting restriction.

Raft against Paxos

Multi-PaxosRaft
LeadershipOptional; an optimization for livenessMandatory; all entries flow from the leader
New leader's first jobRun phase 1 across positions and re-propose whatever may have been chosenNothing; the election restriction guarantees it already has everything committed
Log divergenceHandled by re-proposal per positionFollowers are overwritten to match the leader
Entries out of orderPositions can be decided independently, allowing gapsStrictly sequential, which is simpler and slightly less parallel
Ballot identityBallot number, unique per proposerTerm number, at most one leader per term
Steady-state costOne round trip to a majorityOne round trip to a majority

The two protocols solve the same problem with the same message complexity. The difference is where the complexity is spent, and Raft's decision to forbid gaps in the log is a real, if usually acceptable, loss of pipelining.

Changing the membership

Adding or removing a server cannot be done by editing configuration files, because during the change two different majorities can exist simultaneously and elect two leaders. Raft offers two answers. Joint consensus introduces a transitional configuration in which decisions require majorities of both the old and the new sets, so no split is possible; once that transition is committed, the cluster moves to the new configuration alone. The simpler single-server change adds or removes one member at a time, which guarantees the old and new majorities overlap without a transitional phase. Both are part of the protocol, not an operational procedure, and getting them wrong is one of the more common sources of real incidents.

Common misconceptions

  • "Raft is stronger than Paxos." They solve the same problem with the same fault tolerance and the same steady-state cost. Raft's contribution is comprehensibility and a complete specification.
  • "Any server with an up-to-date log can be elected." A candidate needs a majority of votes, and each voter compares logs, so a candidate missing committed entries is refused by the intersecting voter.
  • "An entry on a majority is committed." Only if it was created in the current leader's term. Entries from earlier terms are committed indirectly, and treating them otherwise is the Figure 8 bug.
  • "Randomized election timeouts are a tuning detail." They are how Raft escapes the split-vote livelock, and setting the range too narrow reintroduces it.
  • "A leader can serve reads from its own state without coordination." Not safely: it may have been deposed without knowing. Correct implementations confirm leadership with a heartbeat round or hold a time-based lease, which is why linearizable reads cost something.
  • "Membership changes are an operations task." They are part of the consensus protocol, because a naive change can create two disjoint majorities and therefore two leaders.

Putting it together

  • Raft is a strong-leader protocol: entries flow only from leader to follower, and disagreeing followers are overwritten.
  • Terms are ballot numbers with an ordering rule that makes a stale leader step down automatically on the first rejection.
  • Randomized election timeouts avoid split votes, which is the same escape from FLP that a distinguished proposer provides in Paxos.
  • The election restriction, that a voter refuses a candidate whose log is less up to date than its own, guarantees an elected leader holds every committed entry, so no recovery phase is needed.
  • The AppendEntries consistency check yields the log matching property, and its rejection path is the sole mechanism for repairing divergent logs.
  • A leader may commit by replica count only for entries from its own term; earlier entries commit indirectly, and ignoring this permits a committed entry to be overwritten.

Both protocols assume they know who the members are, and both need someone to notice when a member is gone. That is the subject of the next lesson.

Sources

  1. Ongaro, D., & Ousterhout, J. (2014). In search of an understandable consensus algorithm. USENIX Annual Technical Conference. raft.github.io
  2. Ongaro, D. (2014). Consensus: bridging theory and practice (PhD dissertation). Stanford University. web.stanford.edu
  3. USENIX. (2014). In search of an understandable consensus algorithm: presentation record. usenix.org
  4. The Raft Consensus Algorithm. (n.d.). Raft: paper, visualization, and implementations. raft.github.io
  5. etcd. (n.d.). API guarantees: linearizability and serializable reads. etcd documentation. etcd.io
  6. Wikipedia contributors. (n.d.). Raft (algorithm). en.wikipedia.org
Key terms
Strong leader
Raft's rule that log entries flow only from leader to followers, so divergent follower logs are overwritten rather than merged.
Term
A numbered period with at most one leader; a message carrying a higher term forces the recipient to step down.
Election timeout
A randomized interval after which a follower becomes a candidate; randomization prevents perpetual split votes.
Election restriction
A voter refuses a candidate whose log is less up to date, guaranteeing that an elected leader has every committed entry.
Log matching property
If two logs share an entry with the same index and term, they are identical in all preceding entries.
Commit index
The highest log index known to be replicated on a majority and therefore safe to apply to the state machine.
Current-term commit rule
A leader may commit by counting replicas only for entries created in its own term; earlier entries commit indirectly.
Joint consensus
A transitional configuration requiring majorities of both the old and new membership, preventing two disjoint majorities during a change.

Leader Election, Failure Detection, and Membership

  • Design a heartbeat-based failure detector and reason about its timeout trade-off.
  • Explain why a lease alone is unsafe and what a fencing token adds.
  • Compare consensus-based membership with gossip-based membership and say when each is appropriate.

Google's Chubby lock service, described by Mike Burrows in a 2006 paper, does something a local mutex never does: when a client acquires a lock, Chubby hands it a sequencer, a short string containing the lock's name and a generation number. The client is expected to pass that string along with every request it makes to whatever resource the lock protects, and the resource is expected to check it and reject anything stale. The reason for this apparatus is one sentence long: a client holding a lock may be paused, may lose its lease without noticing, and may then issue a request that arrives after somebody else has taken over.

Detecting failure with heartbeats

Every failure detector in production is a heartbeat and a timeout. Process A expects a message from B every T milliseconds and suspects B when none arrives within some window. The only design question is the window, and it is a genuine trade-off with no correct answer.

TimeoutDetection timeFalse positivesConsequence of getting it wrong
Short, for example 500 msFastFrequent, on any garbage collection pause or network hiccupUnnecessary failovers, leader churn, and in the worst case a cluster that spends its time electing rather than serving
Long, for example 30 sSlowRareThirty seconds of requests going to a dead node

Note the asymmetry hiding in the first row. A false positive does not merely cause a spurious failover; under load it can cause a correlated one, because whatever slowed the first node is probably slowing the others. Aggressive timeouts are a classic ingredient of metastable failures, where a system that would have recovered on its own instead thrashes indefinitely.

A refinement worth knowing is the accrual failure detector of Hayashibara and colleagues, used in Cassandra and Akka. Rather than a boolean, it maintains a sliding window of observed heartbeat inter-arrival times, fits a distribution, and outputs a continuously varying suspicion value phi, the negative logarithm of the probability that a heartbeat this late would still arrive. A link that is consistently slow raises phi slowly; a link that has genuinely gone silent raises it fast. The application then chooses its own threshold, which decouples the measurement from the policy.

In short: you cannot detect failure, so the engineering question is what error rate you can tolerate in each direction, and the answer differs for a leader election and for a load balancer's health check.

Leases, and why a lease alone is not enough

A lease is a lock with an expiry. The holder may act as leader until time T, and must renew before then. If the holder dies, the lease simply expires and someone else may take it, with no need to detect anything.

Leases rest on a clock assumption, and the clock lesson already showed how that goes wrong. Worse, there is a failure mode leases do not address at all:

t=0    node A acquires the lease, valid until t=10
t=3    A begins a 20-second garbage collection pause
t=10   the lease expires. A does not notice; it is paused.
t=11   node B acquires the lease and begins acting as leader
t=23   A wakes, still believing it holds the lease, and sends a write

Two leaders have now written. Nothing detected the overlap.

The repair is not a shorter lease, which only shrinks the window. It is a fencing token: every lease grant carries a monotonically increasing number, the holder attaches it to every request, and the resource being protected refuses any request with a number below the highest it has seen. A's write at t=23 carries token 7 while the storage has already seen token 8 from B, so it is rejected by the storage rather than by A's good judgement.

lease granted to A  -> token 7
lease granted to B  -> token 8
storage sees write(token=8) from B: accepts, remembers 8
storage sees write(token=7) from A: 7 < 8, REJECT

Chubby's sequencer is exactly this, ZooKeeper's zxid and epoch serve the same purpose, and a Raft term number does the job inside the protocol. A lock without fencing is a suggestion.

Electing a leader, in practice

Since unique leader election is equivalent to consensus, systems do not invent one. They either run a consensus protocol themselves, or delegate to a coordination service that already does.

The standard recipe in ZooKeeper or etcd: every candidate creates an ephemeral, sequentially numbered node under a well-known path. Ephemeral means the node vanishes when the candidate's session ends, including when it crashes. The candidate with the lowest sequence number is the leader.

create /election/n_  as EPHEMERAL SEQUENTIAL   -> /election/n_0000000042
list /election
if my node has the smallest number:  I am the leader
else: WATCH the node immediately below mine, and wait

The last line matters more than it looks. Watching the whole directory instead would wake every candidate on every change, a herd of clients that all re-read and all but one go back to sleep. Watching only your immediate predecessor means exactly one client wakes per departure. That is a design pattern rather than a protocol detail, and it recurs whenever a coordination service is used for anything at scale.

Remember: a coordination service does not make leader election free. It moves the consensus round trip into a service designed to do it correctly, which is almost always the right decision.

Membership: who is in the group?

Both Paxos and Raft assume the participant set is known. Determining it is a separate problem, and there are two families of answer with genuinely different properties.

Consensus-based membershipGossip-based membership
MechanismMembership changes are entries in the replicated logNodes exchange partial state with random peers
GuaranteeAll members agree on the current viewViews converge eventually and may disagree meanwhile
ScaleTens of nodesThousands
During a partitionThe minority side cannot change membershipBoth sides continue with divergent views
Used byetcd, ZooKeeper, Raft clusters generallyCassandra, Consul's data plane, Serf, Akka Cluster

The gossip approach has a well-specified representative: SWIM, from Das, Gupta and Motivala in 2002, whose two ideas are worth stating precisely. First, failure detection is separated from dissemination. Each node periodically probes one random peer; if the probe times out, it asks k other randomly chosen nodes to probe on its behalf, and only if all the indirect probes also fail does it declare the peer suspect. That indirect step removes most false positives caused by one bad link rather than a dead node. Second, membership updates ride along on the probe messages themselves rather than being broadcast, so the network cost per node is constant regardless of cluster size.

A refinement most implementations add is a suspicion state between alive and dead, with a timeout, during which the suspected node can refute the claim about itself. That single state converts many false positives into brief noise instead of a removal.

Choosing between them

The choice is not a matter of taste, and a large system usually needs both.

  • Use consensus-based membership when the membership itself determines correctness: the set of Raft voters, the replicas that must acknowledge a write, the holder of a lock. Disagreement here produces two majorities and therefore two leaders.
  • Use gossip when membership is an optimization: which backends to route to, which nodes to gossip with next, which peers hold which shards approximately. Disagreement here produces a request sent to a node that is going away, which retries handle.

Cassandra is the illustration: it gossips membership across a large ring, and when it needs an actual agreement, such as a lightweight transaction, it runs Paxos over a quorum of replicas. Two mechanisms, two purposes, and it would be a mistake to use either for the other's job.

Bottom line: ask what happens if two nodes disagree about the membership for thirty seconds. If the answer is a retry, gossip is fine. If the answer is two leaders, you need consensus.

Common misconceptions

  • "A shorter timeout gives faster failure detection at no cost." It raises the false positive rate, and false positives under load are correlated across nodes, which is how a slow cluster becomes an unavailable one.
  • "A lease prevents two leaders." It prevents two leaders only if every process notices its lease expiring in time, which a paused process does not do. Fencing tokens are what make the guarantee real.
  • "A distributed lock behaves like a mutex." A mutex is enforced by the operating system, which cannot be bypassed. A distributed lock is enforced by convention unless the protected resource checks a token.
  • "Gossip protocols are eventually consistent, so membership is eventually correct." It converges only if the churn rate is lower than the convergence rate. A cluster with continuous membership change may never settle.
  • "Using ZooKeeper means you do not need to think about consensus." It means the hard part is implemented correctly elsewhere. You still choose the timeouts, handle session expiry, and decide what happens to in-flight work when leadership moves.
  • "Watching the whole election directory is simpler and equivalent." It is simpler and produces a thundering herd on every change; watching only your predecessor wakes exactly one client.

What you now know

  • Every failure detector is a heartbeat and a timeout, and the timeout trades detection speed against a false positive rate whose errors are correlated under load.
  • Accrual detectors output a continuous suspicion level from observed inter-arrival times, separating measurement from policy.
  • A lease bounds how long a dead leader's authority persists, but a paused process can outlive its own lease without noticing.
  • Fencing tokens, checked by the protected resource, are what make leases safe; Chubby's sequencers, ZooKeeper's zxid and Raft's term all serve this role.
  • The standard election recipe uses ephemeral sequential nodes with each candidate watching only its immediate predecessor, avoiding a thundering herd.
  • Consensus-based membership is required wherever disagreement would produce two majorities; gossip-based membership scales to thousands and is right wherever disagreement merely causes a retry.

With agreement covered, the remaining question is what happens when several machines must commit or abort the same transaction together, which turns out to be consensus wearing yet another costume.

Sources

  1. Burrows, M. (2006). The Chubby lock service for loosely-coupled distributed systems. 7th USENIX Symposium on Operating Systems Design and Implementation. research.google.com
  2. Hunt, P., Konar, M., Junqueira, F. P., & Reed, B. (2010). ZooKeeper: wait-free coordination for internet-scale systems. USENIX Annual Technical Conference. usenix.org
  3. Das, A., Gupta, I., & Motivala, A. (2002). SWIM: Scalable weakly-consistent infection-style process group membership protocol. International Conference on Dependable Systems and Networks. cs.cornell.edu
  4. Hayashibara, N., Defago, X., Yared, R., & Katayama, T. (2004). The phi accrual failure detector. In Proceedings of the 23rd IEEE International Symposium on Reliable Distributed Systems, 66-78.
  5. Apache Software Foundation. (n.d.). ZooKeeper internals: sessions, zxids and the atomic broadcast protocol. zookeeper.apache.org
  6. Google. (2017). Managing critical state: distributed consensus for reliability. In Site reliability engineering. sre.google
  7. Wikipedia contributors. (n.d.). Gossip protocol. en.wikipedia.org
Key terms
Heartbeat
A periodic liveness message whose absence within a chosen window causes suspicion; the basis of every practical failure detector.
Accrual failure detector
A detector reporting a continuous suspicion level derived from observed heartbeat inter-arrival times, leaving the threshold to the application.
Lease
A time-bounded grant of authority that expires without needing anyone to detect a failure.
Fencing token
A monotonically increasing number attached to every request and checked by the protected resource, which rejects stale holders.
Sequencer
Chubby's name for a fencing token, passed by a lock holder to the resource it is protecting.
Ephemeral sequential node
A coordination-service object that disappears when its creator's session ends and carries an ordering number, the basis of the standard election recipe.
Thundering herd
The wasteful wake-up of every waiter on a shared change, avoided by watching only one's immediate predecessor.
SWIM
A gossip membership protocol using randomized direct and indirect probing, with updates piggybacked on probe messages.

Module 5: Transactions Across Machines

Committing together, quorum systems that never commit at all, and what an isolation level actually promises.

Atomic Commit and the Blocking Problem

  • State the atomic commit problem and explain why a prepared participant has surrendered its right to abort.
  • Show exactly how two-phase commit blocks, and why a participant cannot decide on its own.
  • Explain why three-phase commit is not the fix and what replacing the coordinator with a consensus group achieves.

A participant in a distributed transaction has voted yes. It has written the transaction's changes to its log, it is holding every lock the transaction touched, and it is waiting to be told whether to commit. The coordinator's machine has just failed. The participant may not commit, because some other participant might have voted no. It may not abort, because the coordinator might already have told someone else to commit. It may not release the locks, because it may yet have to do either. So it waits, and every transaction that needs those rows waits behind it.

The problem, which is not quite consensus

Atomic commit. Each of n participants votes yes or no. All participants must reach the same decision, and:

  • If any participant votes no, the decision must be abort.
  • If all vote yes and no failures occur, the decision must be commit.

Compare that with consensus, where validity requires only that the decided value was proposed by someone. Atomic commit has a much stronger validity condition: a single no is a veto. That difference is what makes the problem harder in practice, because a protocol cannot simply pick whichever answer is convenient.

Two-phase commit

PHASE 1: PREPARE
  coordinator -> all participants:  PREPARE
  each participant:
      decide whether it CAN commit (constraints, locks, disk space)
      if yes: write the transaction and a PREPARED record DURABLY,
              keep all locks, reply YES
      if no:  reply NO and abort locally

PHASE 2: DECIDE
  coordinator:
      if all replies are YES: write COMMIT to its own log, then
                              send COMMIT to all
      otherwise:              write ABORT, then send ABORT to all
  each participant: apply the decision, release locks, acknowledge

The interesting moment is the middle. When a participant replies YES it is making a promise it cannot retract: whatever happens, including its own crash and restart, it will be able to commit if told to. That is why the prepared state must be durable and why the locks must be held. The participant has surrendered its autonomy, and the only entity that can release it is the coordinator's decision.

The point: the prepare phase is not a query. It is the participant giving up its right to change its mind, which is exactly what makes atomicity possible and exactly what makes the failure case so painful.

How it blocks

Consider the failure at the worst possible instant.

coordinator sends PREPARE to A and B
A replies YES (durably prepared, holding locks)
B replies YES (durably prepared, holding locks)
coordinator writes COMMIT to its log
coordinator sends COMMIT to A       <-- A commits and releases locks
coordinator CRASHES before sending anything to B

B is now in doubt. It knows only that it voted yes.
  - It may not abort: A may already have committed.
  - It may not commit: it does not know the coordinator's decision.
  - It may not release locks: either outcome may still require them.

B waits for the coordinator to come back. Not for a timeout. Forever.

Participants can help each other a little through a cooperative termination protocol: an in-doubt participant asks the others what they know, and if any of them has heard the decision, or has not yet voted, the outcome can be resolved. But if every surviving participant is prepared and none has heard the decision, they are all equally ignorant, and the protocol is genuinely stuck.

This is why two-phase commit is described as a blocking protocol. The blocking is not a bug in a particular implementation; Dale Skeen proved in 1981 that no protocol with only these two phases can avoid it, because there is a state, prepared and undecided, from which the correct action depends on information a participant does not have.

The operational shape of this

What a database administrator sees is a transaction stuck in a prepared state, holding locks, blocking unrelated work. Most systems offer a manual override, a heuristic decision, letting an operator force a commit or an abort. It resolves the immediate outage and it can silently violate atomicity: if the operator forces abort while the coordinator had decided commit, one participant has committed and another has not, and the databases are now inconsistent with no error anywhere. The XA specification has a name for this outcome, a heuristic exception, which tells you how routine it became.

That experience is most of the reason distributed transactions acquired their reputation in the 2000s. The protocol was correct. Its failure mode required a human, at an unpredictable hour, making a decision with insufficient information.

Three-phase commit, and why it did not save anyone

The obvious repair is to insert a phase. Three-phase commit adds a pre-commit round between voting and committing: once all participants have voted yes, the coordinator tells everyone "prepare to commit", and only after that acknowledgement does it send the commit. The point is that a participant in the pre-commit state knows everyone voted yes, so if the coordinator dies, the survivors can safely commit on their own after a timeout.

It works, under assumptions nobody has:

  • It requires a synchronous network with a known bound on message delay, because the recovery step is a timeout that must be provably long enough. In a partially synchronous network the timeout can fire early.
  • It fails under network partitions. Two groups can time out on opposite sides and reach opposite decisions, because each side reasons about the other's state from its own incomplete view.

So three-phase commit trades a blocking protocol that is always safe for a non-blocking protocol that can be unsafe. Given the choice between stalling and corrupting, production systems take the stall, which is why 3PC appears in textbooks and almost nowhere else.

The fix that works: make the coordinator not die

The real problem in the blocking scenario is not the number of phases. It is that the coordinator is a single point of failure holding a decision nobody else has. Replace it with a replicated state machine: run the coordinator's log through Paxos or Raft, so that the decision to commit is itself a consensus decision recorded on a majority.

ordinary 2PC:      decision lives on one coordinator, which may vanish
replicated 2PC:    the decision is a committed entry in a consensus log,
                   so a new coordinator can be elected and will find it

Now the failure of a coordinator is a leader election, not an indefinite stall. The transaction is delayed by an election timeout rather than by a human. This is what Google's Spanner does: each shard is a Paxos group, and a transaction spanning shards runs two-phase commit between the groups, with the coordinator role held by one group's leader and the decision durable in that group's replicated log. Two-phase commit did not go away; it was placed on top of consensus, which is exactly the layering the theory recommends since atomic commit and consensus are equivalent problems.

Why this matters: the useful question about a distributed transaction system is not whether it uses two-phase commit. It is whether the coordinator's decision is replicated, because that is what determines how long a failure blocks you.

Or avoid the problem entirely

Many systems decide that atomicity across services is not worth its price and restructure the work instead.

ApproachMechanismWhat you give up
SagaA sequence of local transactions, each with a compensating action to undo it if a later step failsIsolation: intermediate states are visible, and compensation is a business action rather than a rollback
Transactional outboxWrite the business change and an outbox row in one local transaction; a relay publishes the message afterwardsImmediacy: the message is at-least-once and slightly delayed, so consumers must deduplicate
Idempotent receiverMake every downstream effect safe to repeat, keyed by a request identifierSimplicity of the data model, which now carries deduplication state
Single service ownershipRedraw boundaries so the operation touches one datastoreFlexibility of the service decomposition

The last row is the one most often overlooked and most often correct. A distributed transaction is frequently the symptom of a boundary drawn in the wrong place, and moving the boundary is cheaper than making the transaction work.

Common misconceptions

  • "Two-phase commit is unsafe." It is safe: it never produces a partial commit. It is blocking, which is a liveness failure, and the two are frequently confused.
  • "A participant can time out and abort." Before voting, yes. After voting yes, no: another participant may already have committed on the coordinator's instruction.
  • "Three-phase commit solves the problem." It solves blocking under crash failures in a synchronous network, and it is unsafe under partitions, so it trades a liveness problem for a safety problem.
  • "A heuristic decision resolves a stuck transaction." It unblocks the locks and may leave participants permanently inconsistent, which is why the specification treats it as an exception requiring reconciliation.
  • "Sagas are two-phase commit without the coordinator." Sagas abandon isolation entirely: intermediate states are visible to everyone, and undoing means running a compensating business action, not a rollback.
  • "Distributed transactions are always the wrong choice." With a replicated coordinator they are practical, and Spanner runs them at scale. The historical objection was to an unreplicated coordinator.

Recap

  • Atomic commit requires unanimity, with a single no acting as a veto, which is a stronger validity condition than consensus.
  • Voting yes in the prepare phase is an irrevocable promise, which is why the prepared state is durable and the locks are held.
  • If the coordinator fails after some participants have voted yes, those participants are in doubt and may neither commit nor abort nor release locks.
  • Cooperative termination helps only if some surviving participant knows the decision or has not yet voted.
  • Three-phase commit removes blocking under crash failures but is unsafe under partitions, so production systems keep the blocking protocol.
  • The practical fix is to replicate the coordinator's log through consensus, as Spanner does with Paxos groups, turning an indefinite stall into an election timeout; alternatively, restructure the work as sagas or move the transaction inside one service.

The next lesson takes the opposite path: systems that never coordinate at all, and the arithmetic that tells you what their reads are worth.

Sources

  1. Gray, J. (1978). Notes on data base operating systems. In Operating Systems: An Advanced Course (Lecture Notes in Computer Science, vol. 60, pp. 393-481). Springer.
  2. Skeen, D. (1981). Nonblocking commit protocols. In Proceedings of the ACM SIGMOD International Conference on Management of Data, 133-142.
  3. Bernstein, P. A., Hadzilacos, V., & Goodman, N. (1987). Concurrency control and recovery in database systems. Addison-Wesley.
  4. Corbett, J. C., et al. (2012). Spanner: Google's globally-distributed database, sections on two-phase commit over Paxos groups. USENIX OSDI. research.google.com
  5. Richardson, C. (n.d.). Pattern: Saga. Microservices.io. microservices.io
  6. Wikipedia contributors. (n.d.). Two-phase commit protocol. en.wikipedia.org
  7. Wikipedia contributors. (n.d.). Three-phase commit protocol. en.wikipedia.org
Key terms
Atomic commit
Reaching a unanimous commit-or-abort decision among participants, where any single no forces abort.
Prepared state
A durable participant state in which it has promised to be able to commit and has surrendered the right to abort unilaterally.
In-doubt transaction
One whose participant has voted yes and has not learned the decision, holding locks until it does.
Blocking protocol
One in which a failure can leave correct participants unable to proceed; two-phase commit is safe but blocking.
Cooperative termination
In-doubt participants querying each other, which resolves the transaction only if someone knows the decision or has not yet voted.
Heuristic decision
An operator forcing a stuck transaction to commit or abort, which unblocks locks and can silently break atomicity.
Replicated coordinator
Running the coordinator's decision log through consensus so a failure becomes an election rather than an indefinite stall.
Saga
A sequence of local transactions with compensating actions, trading isolation for the removal of distributed atomicity.

Quorums, and What Dynamo-Style Systems Actually Promise

  • Compute quorum configurations and state what R plus W greater than N does and does not guarantee.
  • Explain sloppy quorums, hinted handoff, read repair and anti-entropy, and the guarantee each one weakens or restores.
  • Choose a consistency level for a given operation and name the anomaly it permits.

Amazon's internal service agreements, as described in the 2007 Dynamo paper, are written at the 99.9th percentile rather than the mean. A service that answers within 300 milliseconds for 999 requests out of every thousand has met its obligation; one with an excellent average and an ugly tail has not. Designing to that number, for a shopping cart that must accept a write even while parts of the infrastructure are on fire, produced a storage system with no leader, no locks, and no ability to refuse a write. It also produced a set of guarantees that are widely misread, and the arithmetic in this lesson is how you read them correctly.

The quorum arithmetic

Each key is stored on N replicas. A write is acknowledged once W replicas confirm it; a read consults R replicas and returns the most recent version among them. Two inequalities matter:

R + W > N     the read set and the write set must overlap in
              at least one replica, so a read sees the latest
              successful write

W > N / 2     two write quorums must overlap, so two concurrent
              writes cannot both succeed without any replica
              seeing both

Some configurations, with N = 3 and then N = 5:

NWROverlap?ToleratesCharacter
322Yes, 4 > 31 replica down for both reads and writesThe usual balanced choice
331Yes, 4 > 3No replica down for writesFast reads, fragile writes
313Yes, 4 > 3No replica down for readsFast writes, fragile reads
311No, 2 < 32 replicas downFastest, no recency guarantee at all
533Yes, 6 > 52 replicas downBalanced with more headroom

Two further observations that the table does not show. Response time is governed by the Wth fastest replica, not the slowest, so raising N while holding W fixed improves tail latency by giving the request more chances to avoid a slow node. And availability improves in the same way: with N = 5 and W = 3, two replicas may be down and writes still succeed.

What the inequality does not buy you

Here is the part that gets misquoted. Satisfying R + W greater than N does not make the store linearizable, and the reasons are concrete rather than theoretical.

  • Sloppy quorums. If the designated replicas for a key are unreachable, a Dynamo-style system will write to other nodes instead, so that the write can succeed. Then the W nodes that acknowledged are not from the N that a reader will consult, and the overlap guarantee simply does not apply.
  • Partially successful writes. A write that reaches two replicas but not the required three is reported as failed, and it is not rolled back. A later read may or may not see it, depending on which replicas it consults, and the client that received an error has no way to know.
  • Concurrent writes. Two writes with no causal relationship must be resolved somehow. Last-write-wins by timestamp silently discards one of them, as the clock lesson showed; keeping siblings pushes the decision to the application.
  • Reads concurrent with writes. A read overlapping a write may return either value, and a subsequent read may return the older one, violating monotonic reads even though the arithmetic is satisfied.
  • Restored replicas. A node restored from a backup, or one that lost its data and rejoined empty, can reduce the number of replicas actually holding a value below W without anyone noticing.

What matters here: R plus W greater than N is a statement about set intersection under ideal conditions. It is not a consistency model, and treating it as one is the most common mistake made with these systems.

The machinery around the quorum

Consistent hashing. Keys and nodes are hashed onto a ring; a key belongs to the first node clockwise from it, and its N replicas are the next N distinct nodes. Adding or removing a node moves only the keys between it and its neighbour, rather than reshuffling everything as a modulo-based scheme would. Real systems give each physical machine many virtual nodes scattered around the ring, which evens out the load and spreads the recovery work for a failed machine across many peers rather than one.

Hinted handoff. When a write goes to a substitute node because the intended replica is unreachable, the substitute records a hint saying who the data really belongs to, and delivers it when that node returns. This is what makes a sloppy quorum eventually converge, and it is also what makes the sloppy quorum acceptable: the data is not lost, only temporarily misfiled.

Read repair. When a read consults R replicas and finds one of them stale, the coordinator writes the newer value back to it. Repair is therefore driven by traffic, which means hot keys stay consistent and cold keys can be stale for a very long time.

Anti-entropy with Merkle trees. Cold keys are handled by a background process comparing replicas. Comparing every key would be prohibitive, so each replica maintains a hash tree over its key range: two replicas compare root hashes, and if they differ, descend only into the subtrees that differ. Two replicas holding a million identical keys confirm agreement in one comparison, and two differing in one key locate it in about twenty.

Conflict resolution, and what the application must decide

Concurrent writes are detected with version vectors, as in Module 2, and then something must be done. Two families:

PolicyBehaviourWhen it is acceptable
Last write winsKeep the value with the larger timestamp, discard the otherWhen losing a concurrent write is genuinely harmless: caches, immutable-by-convention data, telemetry
Siblings returned to the clientBoth versions are stored and handed back on the next readWhen the application has a merge rule: cart union, set union, largest value, user prompt
Conflict-free replicated data typeThe data type's merge is commutative, associative and idempotent, so replicas converge automaticallyCounters, sets, registers with a defined merge, collaborative text

The Dynamo paper's own example is the shopping cart, merged by union, which is why an item removed during a partition can reappear. That is a deliberate choice: for a cart, a resurrected item costs a customer one click and a lost purchase costs a sale.

Bottom line: a leaderless store cannot resolve a conflict correctly, because correctness is a property of the application's semantics. What it can do is detect the conflict reliably and refuse to hide it.

Tunable consistency, per request

Cassandra exposes the arithmetic directly as a per-query consistency level, which is the practical form all of this takes.

LevelMeaningTypical use
ONEOne replica respondsHigh-volume writes where loss of a recent value is tolerable
QUORUMA majority of all replicasPaired with QUORUM reads, gives the overlap property
LOCAL_QUORUMA majority within the local data centreAvoiding cross-region latency, at the cost of cross-region recency
ALLEvery replicaRarely: any single replica failure blocks the operation

Two implications. First, the classification of the system as AP or CP belongs to the request, not to the database, exactly as the CAP lesson argued. Second, LOCAL_QUORUM deserves scrutiny: it gives quorum semantics within a region and no recency guarantee across regions, which is usually what is wanted and is frequently described as though it were full quorum consistency.

When genuine agreement is required, these systems reach for consensus rather than quorum arithmetic. Cassandra's lightweight transactions run Paxos across the replicas of a partition, at roughly four round trips, precisely because compare-and-set cannot be built out of overlapping read and write sets.

Common misconceptions

  • "R plus W greater than N gives strong consistency." It gives set overlap under ideal conditions. Sloppy quorums, partially failed writes and concurrent writes each break it in a different way.
  • "A failed write leaves no trace." A write that reached fewer than W replicas is reported as an error and is not undone, so its value may appear in a later read.
  • "Read repair keeps replicas in sync." It repairs what is read. Data nobody reads is repaired only by anti-entropy, which is why that background process is not optional.
  • "Higher N always means better durability." It means more copies and more repair traffic; durability also depends on W, on failure correlation, and on whether the copies are in the same rack.
  • "Consistent hashing exists to spread load evenly." Its purpose is to minimize the keys that move when membership changes. Even load comes from virtual nodes layered on top.
  • "These systems cannot do compare-and-set." They can, by running consensus for that operation, which is why it costs several round trips and is offered as a separate, slower facility.

What to carry forward

  • R plus W greater than N makes read and write sets overlap, and W greater than half of N makes two write sets overlap; both are necessary and neither is sufficient for linearizability.
  • Sloppy quorums, unrolled-back partial writes, concurrent writes and restored replicas each defeat the overlap argument in practice.
  • Consistent hashing with virtual nodes bounds the data movement on membership change and spreads recovery work across many peers.
  • Hinted handoff makes sloppy quorums converge, read repair fixes what is read, and Merkle-tree anti-entropy fixes what is not.
  • Conflicts can be detected reliably by version vectors but resolved correctly only by application semantics, whether by merge rules or conflict-free data types.
  • Per-request consistency levels put the CAP classification on the operation rather than the product, and genuine compare-and-set requires consensus rather than quorums.

Everything so far has concerned single values. The next lesson puts transactions back on top and asks what an isolation level is actually promising.

Sources

  1. DeCandia, G., Hastorun, D., Jampani, M., Kakulapati, G., Lakshman, A., Pilchin, A., Sivasubramanian, S., Vosshall, P., & Vogels, W. (2007). Dynamo: Amazon's highly available key-value store. 21st ACM Symposium on Operating Systems Principles. allthingsdistributed.com
  2. Vogels, W. (2007). Amazon's Dynamo. All Things Distributed. allthingsdistributed.com
  3. Kleppmann, M. (2017). Leaderless replication and the limitations of quorum consistency. In Designing data-intensive applications (ch. 5). O'Reilly Media.
  4. Apache Software Foundation. (n.d.). Dynamo: consistent hashing, replication and repair in Cassandra. Apache Cassandra documentation. cassandra.apache.org
  5. Wikipedia contributors. (n.d.). Quorum (distributed computing). en.wikipedia.org
  6. Wikipedia contributors. (n.d.). Consistent hashing. en.wikipedia.org
  7. Wikipedia contributors. (n.d.). Merkle tree. en.wikipedia.org
Key terms
Quorum condition
R plus W greater than N, which forces the read set and the write set to share at least one replica.
Write quorum overlap
W greater than half of N, which prevents two concurrent writes from both succeeding with no replica observing both.
Sloppy quorum
Accepting a write on substitute nodes when the designated replicas are unreachable, which voids the overlap guarantee.
Hinted handoff
A substitute node recording who the data belongs to and delivering it when that node returns.
Read repair
Writing a fresher value back to a stale replica discovered during a read; it repairs only what is read.
Anti-entropy
A background comparison of replicas, usually via Merkle trees, that repairs data nobody is reading.
Consistent hashing
Mapping keys and nodes onto a ring so that a membership change moves only the keys near the changed node.
Virtual nodes
Many ring positions per physical machine, evening out load and spreading recovery work across peers.

Distributed Transactions and Isolation Levels

  • Name the anomalies each isolation level permits, and identify write skew in a concrete example.
  • Explain why snapshot isolation is not serializability, and how serializable snapshot isolation closes the gap.
  • Distinguish serializability from strict serializability and say what the extra real-time requirement costs.

In 1995 five authors including Jim Gray published a paper at SIGMOD with an uncomfortable finding. The ANSI SQL standard defined its isolation levels by listing phenomena that must not occur, and the definitions were ambiguous enough that a database could satisfy the letter of "serializable" while permitting behaviour no serial execution could produce. Several commercial systems were shipping something called serializable that was in fact a weaker level with a different failure mode. Thirty years later the naming confusion is still with us, and it is the reason this lesson defines everything by the anomaly rather than by the label.

The anomalies, in order of severity

AnomalyWhat happensConcrete damage
Dirty writeA transaction overwrites data written by another uncommitted transactionTwo half-applied updates interleave; nobody permits this
Dirty readReading data written by an uncommitted transactionActing on a value that is about to be rolled back
Non-repeatable readReading the same row twice in one transaction and getting different valuesA report whose totals do not add up because rows changed mid-scan
PhantomA query's result set gains rows when repeated, because another transaction inserted matching rowsA uniqueness check that passes and is then violated
Lost updateTwo read-modify-write cycles interleave and one is silently overwrittenTwo concurrent increments produce one increment
Write skewTwo transactions each read a set, each check a condition over it, and each write to a different row, jointly breaking the conditionThe subtle one; worked below

The levels, defined by what they permit

LevelDirty readNon-repeatable readPhantomLost updateWrite skew
Read uncommittedAllowedAllowedAllowedAllowedAllowed
Read committedPreventedAllowedAllowedAllowedAllowed
Snapshot isolationPreventedPreventedPreventedPreventedAllowed
SerializablePreventedPreventedPreventedPreventedPrevented

The one bold cell is the whole subject of this lesson.

Snapshot isolation, and the anomaly it cannot see

Snapshot isolation gives each transaction a consistent view of the database as of the moment it began, implemented with multiversion concurrency control: writes create new versions, and each transaction reads the versions that were committed at its start timestamp. Readers never block writers and writers never block readers, which is why it is the default in PostgreSQL's repeatable read, in Oracle, and in most distributed databases. Conflicting writes to the same row are caught, usually by a first-committer-wins rule.

Now the failure. Two doctors are on call for the same shift, and the rule is that at least one must remain on call.

Initial state: Alice and Bob are both on call. Count = 2.

Transaction A (Alice)                Transaction B (Bob)
------------------------------------------------------------------
BEGIN                                BEGIN
SELECT count(*) WHERE on_call        SELECT count(*) WHERE on_call
   -> 2, so it is safe to leave         -> 2, so it is safe to leave
UPDATE doctors SET on_call = false   UPDATE doctors SET on_call = false
   WHERE name = 'Alice'                 WHERE name = 'Bob'
COMMIT                               COMMIT

Result: nobody is on call. The invariant is broken.

Snapshot isolation permits this and is not malfunctioning. The two transactions wrote to different rows, so there is no write-write conflict to detect. Each read a snapshot in which the condition held. Yet no serial execution could produce this outcome: run A then B and B's count is 1, so B does not leave.

The pattern generalizes to any check-then-act over a set: booking the last seat, enforcing a maximum number of concurrent bookings, allocating a unique username, keeping an account balance non-negative across two accounts, claiming a job from a queue. Every one of these is write skew waiting to happen, and it appears under load rather than in testing.

Remember: if a transaction reads a set, makes a decision from what it read, and writes something that would change the answer, snapshot isolation is not enough.

Three routes to serializability

TechniqueHowCost
Two-phase lockingAcquire shared locks on reads and exclusive locks on writes, release none until commit; predicate or index-range locks catch phantomsReaders block writers and writers block readers; deadlocks require detection and abort; latency is unpredictable under contention
Actual serial executionRun transactions one at a time on a single thread, with all logic in stored procedures and the data in memoryThroughput limited to one core per partition; every transaction must be short and non-interactive
Serializable snapshot isolationRun snapshot isolation optimistically, track read-write dependencies, and abort a transaction when a pattern that could not arise serially is detectedAborts under contention, which the application must retry; excellent when conflicts are rare

The third deserves a sentence of mechanism, because it is what PostgreSQL has implemented since version 9.1 and what most distributed systems reach for. Cahill, Rohm and Fekete showed that every anomaly permitted by snapshot isolation contains a specific structure: a transaction with an incoming and an outgoing read-write dependency, two such edges in a row. The database tracks those edges at runtime and, when the dangerous structure appears, aborts one of the participants. Most transactions never notice; the ones that would have produced write skew get a serialization failure and are retried.

The consequence for application code is worth stating plainly: under serializable snapshot isolation, every transaction must be prepared to be aborted and retried. Code that treats a commit as guaranteed will fail intermittently under load, in a way that looks like a database problem and is not.

Distributed serializability, and the extra requirement

Serializability says the outcome equals some serial order. It does not say which one, and in particular it does not require that order to match real time. A distributed database can be perfectly serializable and still return you a snapshot from ten seconds ago, because ordering your read before a transaction that has already committed is a legal serial order.

For a user this is indistinguishable from a bug. Adding the real-time requirement gives strict serializability, also called external consistency, which is serializability plus the linearizability constraint from Module 3: if transaction A commits before transaction B begins, A precedes B in the serial order.

PropertyMulti-object atomicityReal-time orderExample
Snapshot isolationYes, over a snapshotNoPostgreSQL repeatable read
SerializabilityYesNoPostgreSQL serializable, on one node
LinearizabilityNo, single objectYesetcd, a Raft-backed register
Strict serializabilityYesYesSpanner

The bottom row is what commit-wait buys in the clock lesson. Spanner pays a few milliseconds of deliberate delay on every write so that timestamp order equals real-time order, and that delay is the visible price of the strongest row in the table.

The point: in a distributed system, "serializable" and "you will see the write that just committed" are two separate promises, and most systems make only the first.

How distributed transactions are actually built

  • Two-phase locking plus two-phase commit. The traditional route, and the one whose blocking behaviour the previous lesson described. Correct, and slow under contention.
  • Multiversion timestamp ordering with a replicated coordinator. Spanner assigns commit timestamps from TrueTime, uses two-phase commit across Paxos groups, and serves read-only transactions from a snapshot at a timestamp with no locks and no coordination at all. That last point is a large part of why it performs acceptably: read-only work does not participate in the expensive path.
  • Client-driven snapshot isolation. Google's Percolator built cross-row transactions on top of Bigtable entirely in the client, using extra columns as locks and a timestamp oracle for ordering, which showed that a single-row store can be extended without changing the storage system.
  • Deterministic execution. Calvin and its descendants agree on the transaction order first, via consensus, and then execute deterministically on every replica, which removes the need for an agreement protocol at commit time. The constraint is that the read and write sets must be known before execution begins.

The naming problem

Because the standard defined levels by phenomena rather than by mechanism, vendors implemented different things under the same names.

You ask forYou may get
PostgreSQL repeatable readSnapshot isolation, which is stronger than the standard requires and still permits write skew
Oracle serializableSnapshot isolation, not serializability
MySQL InnoDB repeatable readA snapshot for reads, with locking reads seeing the latest committed data, which is a different combination again
PostgreSQL serializableGenuine serializability via serializable snapshot isolation, with serialization failures

The operational rule that follows is simple. Do not ask what isolation level a system offers; ask, for the specific invariant you care about, whether the system will preserve it. Then write the two-transaction interleaving that would break it and check the documentation, or the behaviour, against that.

Common misconceptions

  • "Snapshot isolation is serializable." It permits write skew, which no serial execution can produce, and write skew is exactly the shape of most business invariants.
  • "Repeatable read means what the standard says." It means whatever the vendor implemented, and the implementations genuinely differ.
  • "Serializable means my read sees the latest commit." It does not. That is the additional real-time requirement of strict serializability.
  • "Serializable snapshot isolation is free." It is cheap when conflicts are rare, and it converts contention into aborts, so every transaction needs a retry loop.
  • "Isolation and consistency mean the same thing." Isolation levels concern concurrent transactions on one logical database; the consistency models of Module 3 concern what different clients observe across replicas. A system has a position on both axes.
  • "A unique constraint prevents write skew." It prevents one instance of it, where the conflict is a duplicate key. Nothing enforces an invariant over a set that no single row violates.

The takeaway

  • Isolation levels are best understood by the anomalies they permit, because the standard's names are ambiguous and vendors differ.
  • Snapshot isolation prevents dirty reads, non-repeatable reads, phantoms and lost updates, and permits write skew.
  • Write skew arises whenever transactions read a set, decide from it, and write different rows, which covers most check-then-act business rules.
  • Serializability is reached by two-phase locking, by actual serial execution, or by serializable snapshot isolation, which detects the dangerous read-write dependency structure and aborts.
  • Serializability alone permits a stale but legal serial order; strict serializability adds the real-time constraint, and Spanner pays commit-wait latency for it.
  • Ask whether your specific invariant survives a specific interleaving, rather than which level a system claims to implement.

One question remains: how are all these pieces assembled into a service that stays up? The last module answers it with the replication abstraction the whole course has been building toward, and then with the systems that shipped it.

Sources

  1. Berenson, H., Bernstein, P., Gray, J., Melton, J., O'Neil, E., & O'Neil, P. (1995). A critique of ANSI SQL isolation levels. ACM SIGMOD. arxiv.org
  2. Adya, A. (1999). Weak consistency: a generalized theory and optimistic implementations for distributed transactions (PhD thesis). Massachusetts Institute of Technology.
  3. Cahill, M. J., Rohm, U., & Fekete, A. D. (2009). Serializable isolation for snapshot databases. ACM Transactions on Database Systems, 34(4), article 20.
  4. Bailis, P., Davidson, A., Fekete, A., Ghodsi, A., Hellerstein, J. M., & Stoica, I. (2014). Highly available transactions: virtues and limitations. VLDB. bailis.org
  5. PostgreSQL Global Development Group. (n.d.). Transaction isolation. PostgreSQL documentation. postgresql.org
  6. Kleppmann, M. (2017). Transactions, and weak isolation levels. In Designing data-intensive applications (ch. 7). O'Reilly Media.
  7. Wikipedia contributors. (n.d.). Snapshot isolation. en.wikipedia.org
  8. Wikipedia contributors. (n.d.). Isolation (database systems). en.wikipedia.org
Key terms
Write skew
Two transactions each reading a set, checking a condition, and writing different rows, jointly violating the condition.
Snapshot isolation
Each transaction reads a consistent view as of its start, with multiversion storage; it prevents most anomalies and permits write skew.
Phantom
A row appearing in a repeated query because another transaction inserted a matching row.
Two-phase locking
Holding shared and exclusive locks until commit, which achieves serializability at the cost of blocking and deadlocks.
Serializable snapshot isolation
Optimistic serializability that tracks read-write dependencies and aborts transactions forming a structure no serial order could produce.
Serialization failure
The abort a database returns when optimistic checking detects a conflict; every transaction needs a retry loop.
Strict serializability
Serializability plus the real-time constraint that a committed transaction precedes any transaction beginning afterwards.
Deterministic execution
Agreeing on transaction order first and executing identically on every replica, removing the need for agreement at commit time.

Module 6: Systems in the Wild

The replication abstraction the whole course has been building toward, and the published systems that show what each design gave up.

State Machine Replication in Practice

  • State the three requirements of state machine replication and identify sources of nondeterminism that break the third.
  • Explain log compaction, snapshot transfer, and the three ways to serve a linearizable read.
  • Compare a replicated log with chain replication and say when each is appropriate.

Every object in a Kubernetes cluster, every pod specification, every secret, every config map, lives in etcd. etcd is not a database with replication attached. It is a replicated log with a key-value store computed from it: if you replayed that log from the first entry on an empty machine, you would obtain the same cluster state, entry for entry. Fred Schneider wrote down the requirements for this in a 1990 tutorial, and the whole of Module 4 exists to satisfy one of them.

The three requirements

A service is replicated correctly by the state machine approach if:

  1. The service is a deterministic state machine: the same command applied to the same state always produces the same new state and the same output.
  2. All replicas start in the same state.
  3. All replicas apply the same commands in the same order.

Given those, every replica holds identical state at every point in the command sequence, and any of them can answer a query. Requirement 3 is exactly the consensus problem, which is why Paxos and Raft are the engines of every such system. Requirement 1 is where working implementations actually go wrong.

Key idea: consensus is the hard part in theory and determinism is the hard part in practice.

The determinism traps

Anything that makes the same command produce different results on different machines corrupts a replica silently, and the corruption is discovered later, in a place unrelated to its cause.

SourceSymptomRepair
Reading the clock inside a commandReplicas store different timestamps for the same writeThe leader reads the clock once and puts the value in the log entry
Random numbers, including UUID generationDifferent identifiers on different replicasGenerate at the client or the leader and include the value in the command
Iteration order over a hash mapDivergence only when the command's effect depends on orderSort explicitly, or use an ordered container
Floating-point differences across compilers or CPUsRare, tiny, and eventually a mismatched checksumUse integers or fixed-point for anything stored
Reading external state: a file, an environment variable, another serviceReplicas disagree whenever the external thing differsFetch it before proposing and embed the result in the command
Wall-clock expiry evaluated at apply timeA key expires on one replica and not anotherStore an absolute expiry decided by the leader, and evaluate against the log's own ordering where possible

The pattern in the repair column is one rule: nondeterministic inputs must be resolved before the command enters the log, not after. The log is the source of truth, so anything the state machine needs must be in it.

The log cannot grow forever

Two problems arrive together. Storage fills up, and a replica that has been down for a week would need to replay a week of commands to catch up.

The answer is a snapshot: serialize the state machine's current state, record the log index it corresponds to, and discard log entries below that index. The interesting parts are operational rather than theoretical.

  • Snapshotting must not stop the service. Implementations use copy-on-write, a fork, or an immutable data structure so that the snapshot is taken from a consistent view while writes continue.
  • A lagging follower may need the snapshot itself. If a follower's next required entry has already been compacted away, the leader cannot send log entries and must transfer the whole snapshot instead. Raft has an explicit InstallSnapshot message for this, and it is the only case in the protocol where state, rather than log, moves between servers.
  • Snapshot frequency is a real trade-off. Too often and you burn IO on redundant work; too rarely and recovery is slow and disks fill.

Serving reads

Writes must go through the log. Reads need not, and how a system handles them determines both its performance and its actual consistency guarantee.

MethodMechanismGuaranteeCost
Read through the logAppend a no-op read entry and answer when it commitsLinearizableA full consensus round, plus a durable write
ReadIndexRecord the current commit index, exchange heartbeats with a quorum to confirm still-leader, wait for the state machine to reach that index, then read locallyLinearizableOne round trip of messages, no disk write
Lease readRely on a time-based leadership lease and read locally with no messages at allLinearizable only if clock drift stays within the assumed boundNothing, plus a clock assumption
Follower readRead any replica's state directlyStale, not linearizable; etcd calls this serializableNothing, and it scales with replicas

The second row is what most systems use, and it explains something that puzzles newcomers: why a read from a distributed key-value store costs a network round trip when the leader obviously has the answer in memory. The leader has the value; what it does not have, without asking, is the knowledge that it is still the leader.

The fourth row is the honest name for reading a follower, and it should be chosen deliberately. It is exactly right for a configuration value that changes rarely, and exactly wrong for a lock.

In short: a linearizable read costs a round trip because leadership is not a fact a node can verify alone.

Two shapes of replicated service

The replicated log. Every replica holds the whole log and the whole state; a leader orders writes; a majority must acknowledge. Throughput is bounded by one leader and by the slowest member of each quorum, and adding replicas does not increase write capacity. This is etcd, ZooKeeper, and a single Raft group.

Chain replication. Van Renesse and Schneider proposed arranging replicas in a chain: writes enter at the head and propagate along it, reads are served by the tail, and a write is committed once it reaches the tail. Because the tail has by definition seen every committed write, its reads are linearizable with no coordination at all, and the write path is a pipeline rather than a fan-out, which uses network capacity more evenly. The cost is latency proportional to the chain length and a more complex reconfiguration procedure when a link fails, usually delegated to a separate coordination service. The CRAQ variant lets every node serve reads by tracking clean and dirty versions, which greatly improves read throughput for workloads dominated by reads.

Consensus logChain replication
Write latencyOne round trip to a majorityOne traversal of the chain
ReadsLeader, with a leadership checkTail, with no check needed
Failure handlingBuilt in, via electionRequires an external configuration service
Tolerates f failures with2f + 1 nodesf + 1 nodes, given a reliable reconfiguration service

Scaling past one log

A single replicated log is a single ordering point, and therefore a single machine's worth of throughput. Every large system that uses consensus does the same thing about it: partition the data and run one consensus group per partition. CockroachDB calls them ranges, TiKV calls them regions, Spanner calls them Paxos groups, and each is an independent Raft or Paxos instance with its own leader and its own log.

That buys linear scaling for operations confined to one partition, and it hands you the cross-partition problem, which is the two-phase commit of Lesson 13 layered on top of the per-partition consensus. The architecture of every modern distributed SQL database is exactly that sentence.

Kafka is the same idea with the log itself as the product. A partition is an ordered log replicated to a set of followers; the leader tracks the in-sync replica set, and a producer using acknowledgement from all in-sync replicas, together with a minimum in-sync replica count, obtains a durability guarantee equivalent to a quorum write. Setting the minimum to one while requiring all acknowledgements is the classic misconfiguration, since the set can shrink to the leader alone and the guarantee quietly becomes nothing.

Bottom line: consensus does not scale, and it does not need to. Partition until each group's write rate fits one leader, then pay for the cross-partition operations you actually need.

Common misconceptions

  • "Adding replicas increases write throughput." It increases fault tolerance and read capacity in some designs, and it makes writes slower, since a larger quorum must acknowledge.
  • "The leader can answer reads from memory." It can answer, but not linearizably, because it may have been deposed without knowing. ReadIndex or a lease is what makes the answer safe.
  • "Determinism means avoiding threads." It means the same command on the same state gives the same result. Clocks, random numbers, hash iteration order and external lookups are the usual offenders, and concurrency is only one of many.
  • "Snapshots are an optimization." Without them the log grows without bound and recovery time grows with it, so they are a requirement for any long-lived system.
  • "Chain replication is strictly better because it needs fewer nodes." It needs an external service to reconfigure the chain, and that service is itself a consensus system, so the total node count is not the whole comparison.
  • "A Kafka producer with acknowledgement from all replicas cannot lose data." Not unless the minimum in-sync replica count is set above one; otherwise the in-sync set can shrink to the leader and all means one.

Summing up

  • State machine replication needs a deterministic state machine, identical initial state, and identical command order, and the third requirement is precisely consensus.
  • Nondeterminism from clocks, randomness, iteration order, floating point and external lookups must be resolved before a command enters the log.
  • Snapshots bound log growth and recovery time, and a follower whose next entry has been compacted must receive the snapshot itself.
  • A linearizable read costs at least a round trip, because a leader cannot confirm it is still the leader without asking; lease reads trade that for a clock assumption, and follower reads are stale by design.
  • Chain replication puts writes at the head and reads at the tail, giving coordination-free linearizable reads, at the cost of chain-length latency and external reconfiguration.
  • Consensus scales by partitioning into many independent groups, which turns cross-partition work into two-phase commit over those groups.

The final lesson reads the systems themselves, and asks of each one the question this course has been building toward: what did it give up, and what did that buy?

Sources

  1. Schneider, F. B. (1990). Implementing fault-tolerant services using the state machine approach: a tutorial. ACM Computing Surveys, 22(4), 299-319.
  2. van Renesse, R., & Schneider, F. B. (2004). Chain replication for supporting high throughput and availability. USENIX OSDI. usenix.org
  3. Ongaro, D. (2014). Consensus: bridging theory and practice, chapters on log compaction and client interaction. Stanford University. web.stanford.edu
  4. etcd. (n.d.). Why etcd and API guarantees: linearizable versus serializable reads. etcd documentation. etcd.io
  5. Apache Software Foundation. (n.d.). Replication: in-sync replicas and acknowledgement settings. Apache Kafka documentation. kafka.apache.org
  6. Wikipedia contributors. (n.d.). State machine replication. en.wikipedia.org
Key terms
State machine replication
Keeping replicas identical by applying the same deterministic commands in the same order from a shared starting state.
Determinism requirement
The rule that a command must produce the same result on every replica, which forbids clocks, randomness and external lookups at apply time.
Log compaction
Replacing a log prefix with a snapshot of the state it produced, bounding storage and recovery time.
InstallSnapshot
The transfer of a full state snapshot to a follower whose next required log entry has already been compacted away.
ReadIndex
Confirming leadership with a quorum heartbeat and then reading locally at the recorded commit index, giving a linearizable read without a log write.
Lease read
Reading locally on the strength of a time-based leadership lease, which is linearizable only under a clock drift assumption.
Chain replication
Replicas in a chain with writes at the head and reads at the tail, giving coordination-free linearizable reads.
In-sync replica set
Kafka's set of followers currently caught up with the leader; the durability of an acknowledged write depends on its minimum size.

Case Readings: What Each System Gave Up

  • Read a systems paper by identifying its workload assumptions, its failure model, and the guarantee it deliberately abandoned.
  • Compare GFS, Bigtable, Dynamo, Spanner, ZooKeeper and Kafka on what each traded and what it bought.
  • Apply the same analysis to a system you use, and check its claims rather than assuming them.

The Google File System paper opens with an assumption that reads like a complaint. With thousands of commodity machines, it says, component failures are the norm rather than the exception, so monitoring, error detection and automatic recovery must be built into the system rather than added around it. Everything else in that paper follows from taking that sentence seriously, including several decisions that look like mistakes until you know the workload. Reading these papers well means finding the sentence like that one, and then asking what it forced the designers to give up.

How to read a systems paper

Five questions, in order. Applied to any paper, they extract more than a careful summary would.

  1. What is the workload? Read sizes, write patterns, read-to-write ratio, access locality. Almost every surprising decision is explained here.
  2. What is the failure model? Crash-stop or crash-recovery, how many simultaneous failures, whether partitions are considered, whether corruption is considered.
  3. What invariant is preserved, exactly? Not the marketing word. The precise statement, in the vocabulary of Module 3.
  4. What was given up? Every design trades something. If you cannot name it, you have not finished reading.
  5. What would break it? Which change in the workload or the environment would invalidate the design.

Why this matters: a paper's contribution is rarely the mechanism. It is the discovery that a particular guarantee was not needed, which made a cheaper mechanism sufficient.

GFS, 2003: give up file system semantics

The workload was enormous files, mostly appended to rather than overwritten, read mostly in large sequential streams, by applications the same organization controlled. That last clause is the licence for everything else.

A single master holds all metadata in memory and is not on the data path: clients ask it which chunkserver holds a chunk, then talk to the chunkserver directly. Chunks are 64 megabytes, which is enormous by 2003 standards and is precisely what lets one master's memory index a petabyte-scale file system. The consistency model is deliberately weak: a concurrent record append is guaranteed to appear at least once, at an offset the system chooses, and regions of a file may be consistent but undefined, meaning all replicas agree yet the contents are a mixture of concurrent writes.

Given up: POSIX semantics, efficient small files, and any useful guarantee about concurrent overwrites. Bought: a design simple enough to reason about, since a single master means no distributed metadata consensus, and throughput limited by disks rather than coordination. What would break it: many small files, or applications that cannot tolerate duplicate records.

Bigtable, 2006: give up the transaction

A sparse, sorted, multidimensional map, partitioned into tablets, each stored as immutable sorted files with an in-memory write buffer, with a master assigning tablets and Chubby holding the coordination state.

The decision that defines it: transactions are single-row only. No joins, no multi-row atomicity, no secondary indexes. That restriction is what allows a tablet to be moved between servers freely, since no operation spans two tablets, and it is why Bigtable scaled when contemporaneous relational systems did not.

Given up: multi-row atomicity and the relational model. Bought: horizontal scalability with predictable operations, and a data model that survives being resharded. What would break it: an application that genuinely needs cross-row invariants, which is what pushed Google to build Spanner.

Dynamo, 2007: give up consistency, on purpose

Covered in Lesson 14, and worth restating as a trade. The requirement was that the shopping cart must accept a write even during a data centre event, with service objectives written at the 99.9th percentile. The design follows: no leader, no locks, sloppy quorums with hinted handoff, version vectors, siblings returned to the application, gossip membership.

Given up: linearizability, and the convenience of the database resolving conflicts. Bought: a store that cannot be made unavailable by any single failure, and predictable tail latency. What would break it: data whose conflicts have no sensible merge, which is why the paper's example is a shopping cart and not an account balance.

Spanner, 2012: buy consistency back, with hardware

Spanner is interesting precisely because it reverses Dynamo's trade. Data is sharded, each shard is a Paxos group, cross-shard transactions use two-phase commit between group leaders, and TrueTime provides bounded clock uncertainty from GPS receivers and atomic clocks. Commit-wait makes timestamp order agree with real-time order, which gives strict serializability at global scale.

The paper is candid about the price. Every commit waits out the clock uncertainty, so write latency has a floor set by the quality of the timekeeping hardware, and the whole design assumes a private network and infrastructure a single organization controls. Read-only transactions escape the cost by executing at a timestamp with no locks and no coordination, which is what keeps the system usable.

Given up: commit latency and the ability to run on ordinary infrastructure. Bought: globally distributed transactions with the strongest guarantee on the map. What would break it: clock uncertainty growing large, which turns commit-wait from milliseconds into something unusable.

ZooKeeper and Kafka: two ways to sell a log

ZooKeeper replicates a small tree of nodes through an atomic broadcast protocol. Its guarantees are worth stating precisely, because they are often misquoted: writes are linearizable, requests from a single client are applied in the order that client issued them, and reads are served locally by whichever server the client is connected to and may be stale. A client needing a current read must call sync first. That default is a deliberate trade of read recency for read scalability, and it is the source of many surprises for people who assumed the whole system was linearizable.

Kafka makes the log itself the product. A topic is partitioned; each partition is an ordered, replicated, durable log; ordering is guaranteed within a partition and not across them. A leader maintains the set of in-sync replicas, and durability comes from requiring acknowledgement from all in-sync replicas together with a minimum size for that set.

Given up, in both cases: in ZooKeeper, fresh reads by default and any large data volume; in Kafka, global ordering and random access. Bought: read scalability and a simple coordination primitive; and throughput plus replayability, which turns a message queue into a durable record of what happened.

The comparison

SystemCore guaranteeDeliberately abandonedCoordination cost
GFSRecord append occurs at least once; replicas agreePOSIX semantics, defined content under concurrent writesMaster consulted for metadata only
BigtableSingle-row atomicityMulti-row transactions, joinsChubby for tablet assignment
DynamoAlways writable; eventual convergenceLinearizability, automatic conflict resolutionNone on the write path
SpannerStrict serializabilityCommit latency, commodity infrastructurePaxos per shard, two-phase commit across shards, commit-wait
ZooKeeperLinearizable writes, per-client FIFO orderFresh reads by default, data volumeConsensus on writes; reads are local
KafkaPer-partition order and durabilityGlobal ordering, random accessLeader plus in-sync replica acknowledgement

Worth holding on to: read the table's third column first. The abandoned guarantee is what made each system possible, and choosing a system means choosing which absence you can live with.

Do not take the claims on faith

A published guarantee is a claim about an implementation, and implementations have bugs. Kyle Kingsbury's Jepsen project tests distributed systems by generating concurrent operations, injecting partitions and clock skew, and checking the resulting histories against the claimed consistency model, using the history-checking machinery of Module 3. The reports have repeatedly found real violations of documented guarantees in widely deployed systems, and, just as often, ambiguity in what a system claimed in the first place.

Two habits follow. When evaluating a system, look for an independent analysis before believing the documentation. And when writing documentation, state the guarantee in the vocabulary of Module 3, because a claim that cannot be tested cannot be trusted.

Common misconceptions

  • "These papers describe best practices." They describe designs fitted to specific workloads at specific companies. GFS's chunk size is wrong for your photo storage; Dynamo's sibling model is wrong for your ledger.
  • "Newer systems supersede older ones." Spanner did not replace Dynamo; they occupy opposite corners of the same trade-off, and Amazon still runs Dynamo-derived systems for the workloads that need them.
  • "ZooKeeper is linearizable." Its writes are. Its reads are served locally and may be stale unless the client synchronizes first.
  • "Kafka guarantees ordering." Within a partition. Across partitions there is no order at all, which is why the partitioning key is a correctness decision and not a performance one.
  • "Spanner proves the CAP theorem is wrong." Spanner chooses consistency during a partition and is unavailable for affected shards. Its achievement is making partitions rare enough, on a private network, that the choice seldom shows.
  • "A system's documented guarantee is what it provides." Frequently, and not always. Independent testing has found the gap often enough to make checking the default.

Pulling it together

  • Read a systems paper by extracting the workload, the failure model, the exact invariant, the abandoned guarantee, and the change that would break it.
  • GFS gave up file system semantics and defined concurrent-write content in exchange for a single-master design and enormous streaming throughput.
  • Bigtable gave up multi-row transactions, which is exactly what made tablets movable and the system scalable.
  • Dynamo gave up linearizability and automatic conflict resolution to obtain a store that is always writable and predictable in the tail.
  • Spanner bought strict serializability back with bounded-uncertainty clocks and commit-wait, paying latency and requiring specialized infrastructure.
  • ZooKeeper trades fresh reads for read scalability, Kafka trades global ordering for throughput, and independent testing exists because documented guarantees and implemented guarantees sometimes differ.

That closes the course. Seventeen lessons have circled one fact: you cannot tell a slow machine from a dead one, and every guarantee in this field is a statement about what remains true despite that. When you next evaluate a system, ask the five questions from this lesson, name the guarantee in the vocabulary of Module 3, and find out what it costs. Those three habits will outlast every protocol in this file.

Sources

  1. Ghemawat, S., Gobioff, H., & Leung, S.-T. (2003). The Google File System. 19th ACM Symposium on Operating Systems Principles. research.google.com
  2. Chang, F., et al. (2006). Bigtable: a distributed storage system for structured data. USENIX OSDI. research.google.com
  3. DeCandia, G., et al. (2007). Dynamo: Amazon's highly available key-value store. ACM SOSP. allthingsdistributed.com
  4. Corbett, J. C., et al. (2012). Spanner: Google's globally-distributed database. USENIX OSDI. research.google.com
  5. Hunt, P., Konar, M., Junqueira, F. P., & Reed, B. (2010). ZooKeeper: wait-free coordination for internet-scale systems. USENIX ATC. usenix.org
  6. Kingsbury, K. (n.d.). Jepsen analyses: testing distributed systems against their claimed guarantees. jepsen.io
  7. Massachusetts Institute of Technology. (n.d.). 6.824 Distributed Systems: paper reading list and lecture notes. MIT PDOS. pdos.csail.mit.edu
Key terms
Record append
GFS's append operation, guaranteed to place a record at least once at a system-chosen offset, permitting duplicates.
Consistent but undefined
A GFS file region on which all replicas agree while the contents are an interleaving of concurrent writes.
Tablet
Bigtable's unit of partitioning and movement, made freely relocatable by the single-row transaction restriction.
Per-client FIFO order
ZooKeeper's guarantee that a single client's requests are applied in the order it issued them, distinct from linearizable reads.
sync
The ZooKeeper call a client makes before a read when it needs the connected server to have caught up.
Partition key
The field determining a Kafka partition, and therefore a correctness decision, since ordering holds only within a partition.
Jepsen testing
Generating concurrent operations under injected faults and checking the recorded history against a claimed consistency model.

Open the interactive version with quizzes and progress →