The farm came online on a Thursday morning. By a quarter to three that afternoon it did not exist, and every machine involved was healthy, and the row that said it existed was still sitting on two of the three machines that were meant to be holding it.
Galewatch, which collects a reading every two seconds from each wind turbine on the farms it monitors, spent April learning what eighty seconds of downtime on the small Postgres holding lesson 014's routing table can cost. Lesson 032 is that story: eighty seconds of planned maintenance on that database, two movers that both believed they held the job, and a year of one turbine's history deleted by the one that was being polite about it.
So in May they did what lesson 011 had told them to do and nobody had got round to. The routing table moved onto three machines, a primary and two standbys, which I will call route-1, route-2 and route-3. Lesson 011 published both the recipe and the reason for the third machine: name one standby as synchronous and a commit waits for it, and then that standby dying stops every commit on the primary, which lesson 011 called a remarkable thing to build in the name of staying up. So you require any one of two standbys instead, and own three machines to get what you thought two would buy. In Postgres it is one setting, synchronous_standby_names = 'ANY 1 (route-2, route-3)', and from that afternoon every committed routing change sat on at least two machines before the statement returned.
At 13:55 on the Thursday, route-2 was restarted to finish the minor version upgrade that was going round all three. It came back at 13:58 with its replication pointing at the hostname of the old single routing database, the one they had decommissioned three weeks earlier. Postgres logged a connection failure, waited, and tried again. It would have done that forever. route-2 was up. It answered read-only queries. It was replaying nothing.
Nothing complained, and this is the part worth sitting with. Not one commit slowed by a microsecond, because route-3 answered all of them and ANY 1 means the primary never waits for the second standby. A quorum's entire job is to stop waiting for the member that is late, and what you give up in exchange is any way of noticing that the member has stopped.
At 14:12:04 an engineer inserted one row. Galewatch's twenty-first wind farm, Drumlea Moss, forty turbines, had come online that morning, and that row is the only thing in the world that tells the four ingest boxes the farm exists and which machine its readings belong on. The insert returned as soon as route-3 confirmed it, half a millisecond on lesson 002's ladder. Two of the three machines had it.
Forty turbines at a reading every two seconds is twenty readings a second, and for twenty eight minutes they landed exactly where the row said.
At 14:40 the primary's host failed, in the ordinary way hosts fail. The on-call engineer opened the runbook written when the cluster was built, and the runbook said promote route-2.
At 14:41 route-2 was the routing database, with everything repointed at it, holding the routing table as it had stood at 13:55.
The four ingest boxes reload that table every five seconds, wholesale. That was the other thing the three weeks after April bought: lesson 031's incident was four boxes disagreeing for twelve minutes because the divisor shipped as configuration, one box at a time, and moving it into the table turns twelve minutes into five seconds. A cache refreshed wholesale also inherits a deletion nobody asked for. By 14:41:10 none of the four boxes had ever heard of Drumlea Moss.
What happened next is the only reason this is a story about a dashboard and not a story about data. Lesson 031 ended on making the wrong owner loud, and the same instinct had been applied a layer further out: a box with no row for a farm refuses the reading rather than choosing a default. Forty turbines started collecting errors, and forty turbines did what turbines on weather-beaten links have always done, which is keep the reading and try again later.
So the dashboard went flat.
At 14:48 the engineer at Drumlea Moss, who had a handover call at three and a new dashboard to show on it, asked why. The row went back in at 15:06, the backlogs drained, and the bill was twenty five minutes of a new customer's first afternoon and thirty thousand readings arriving late.
The stale hostname was a plain mistake and the least interesting thing here, because the cluster had been built so that exactly that mistake costs nothing until the afternoon it costs everything. route-2 had been replaying nothing for forty two minutes and no part of the system was arranged to care. Had the host failed at seven in the evening the gap would have been five hours, and the promotion would have looked every bit as routine.
Two of three machines had the row. The runbook named the third.
You built a quorum in lesson 011 and left out half of it
A quorum is the number of machines that have to agree before something counts. A majority is the usual choice, and for three machines a majority is two.
Count what lesson 011 actually bought. Any one of two standbys means every acknowledged write sits on the primary plus one standby: two of three. That is a write quorum, nobody called it one, and the property it buys is exactly the one you want, which is that no acknowledged write is held by a single machine.
Now put a promotion next to it. Suppose the choice of the next leader also had to be made by a majority of the three. The voting set has two members and the write set has two members, and two plus two is four, which is more than three, so the sets cannot avoid each other: whoever votes, at least one voter holds the row. Give that voter the right to refuse a candidate whose log is behind its own, and the machine missing the row cannot win. There is no proof in that paragraph, only counting. Two subsets of a three element set, each with two elements, must share an element.
So what May was missing is not a mechanism.
Nobody asked. The runbook had picked the next leader three weeks before the write happened, and a runbook written at build time cannot know which machine will be holding the last row. A write quorum and a vote quorum have to add up to more than the cluster. When they do, the group can refuse to elect a machine that is behind. When they don't, every promotion is a guess that happens to be right most days.
Generalise it with the same supposed vote on top, because this is the step where people hurt themselves. Five machines, one primary and four standbys, ANY 1, set up by somebody reasoning that five is safer than three. Every write lands on two of the five. A majority of five is three. Two plus three is five, which is not more than five, so the write set and a voting set can miss each other completely:
five machines, write on 2 of them, leader chosen by 3 of them
m1 m2 m3 m4 m5
has the write x x . . .
voted for m5 . . x x x
\_____ _____/
v
no machine appears in both rows, so m5
wins the election without the write
Said out loud: m1 and m2 took the write, m3 and m4 and m5 elected m5, those groups share nobody, and the new leader has never seen the entry. Five machines want ANY 2, which is three written against three voting, six against a cluster of five, and the overlap is back. That inequality is worth deriving rather than memorising, because the version you memorise will be the three machine one.
One honest qualification, since the May story is unusually clean. In a healthy cluster the standby that did not confirm catches up in a millisecond or two, so a blind promotion is a coin flip inside a window a few milliseconds wide, which is why thousands of teams run this arrangement and never get bitten. route-2's window was forty two minutes wide, and the thing that widened it was a quorum doing its job.
Agree on the order, not the value
Consensus gets taught as getting a group of machines to agree on a value, which makes it sound like something you do once. The object you want is a list.
Galewatch's routing table is a sequence of decisions. Split Ardnave Point in two. Change its divisor to three. Add Drumlea Moss on machine 11. Apply those to an empty table in that order, on any machine you like, and you get the same table; apply them in a different order and you may not. So the thing three machines have to agree about is not what the table says. It is the order of the changes, and once they agree on that, every copy that replays it is identical by construction.
The name is state machine replication, and lesson 012 showed you one without naming it: a Postgres replica does not re-run your queries, it replays the write ahead log, which is why lesson 012 could say the primary read a table in page order while the replica read it in log order.
So: an ordered log. Entries numbered from one, appended at one end, each machine holding its own copy and its own idea of how far down that copy it may go.
Now the clock, which is the part that surprised me the first time somebody walked me through it. Raft numbers its periods. A term is a stretch of time with at most one leader in it, terms only count upwards, and every message carries the term its sender believes it is in. A message from an older term gets refused. A message from a newer term makes the receiver step down, whatever it thought it was a moment earlier.
You have met that number. Lesson 032 called it a fencing token: a number that only increases, minted by the only party allowed to mint it, compared by the receiver against the highest it has already obeyed. Raft's term is lesson 032's fencing token used as the system's only clock. Nothing in the protocol asks what time it is. It asks which term you are from, and that is answered by comparing against a number the machine already has, which is the one sort of question a machine can answer with nobody's help.
An election that reads the log
Raft was published in 2014 by Diego Ongaro and John Ousterhout under a title that says what it is for: in search of an understandable consensus algorithm. Paxos had been the correct answer for two decades and almost nobody could build a working system from reading it. etcd, which holds all of a Kubernetes cluster's state, runs Raft, and so does Consul. ZooKeeper is the older cousin and runs a protocol of its own, ZAB, which arrives in much the same place by a different road.
Three roles, and a machine is in exactly one of them. A follower does what it is told. A candidate is standing for election. A leader is the one machine appending to the log.
A follower that hears nothing from a leader for the length of its election timeout assumes there isn't one. It increments the term, votes for itself, becomes a candidate and asks every other machine for a vote. A majority makes it leader. Two rules decide those votes, and they are the whole of the safety argument.
A machine grants at most one vote per term. That is what makes "at most one leader per term" true rather than hoped for, by the counting already done: two majorities of the same set share a machine, and the shared machine has already voted.
A machine refuses a candidate whose log is behind its own, comparing the term of the last entry first and the position second. That is the rule May did not have. Replay the Thursday with it. route-1 is gone, both standbys time out, and route-2 asks route-3 for a vote in the new term. route-3 holds the 14:12 entry and route-2 does not, so route-3 refuses, and with its own vote and nothing else route-2 is stuck at one of three. Nobody wins that term. Timeouts fire again, route-3 asks first, and route-2 has no grounds to refuse a candidate ahead of itself. route-3 leads, Drumlea Moss survives, and nobody at the farm notices anything at all.
Two candidates can wake in the same term and split the vote so neither reaches a majority, and the fix is not clever: each machine draws its timeout at random from a range, so the next round starts at different moments and somebody gets there first. Make the range several times wider than a round trip to your peers, or you will hold elections because a packet was slow.
The log side is the same shape. The leader appends an entry and sends it out, and the message carries the position and term of the entry immediately before it. A follower that cannot match that preceding entry refuses the new one, and the leader walks backwards until they agree and then fills the follower forwards. That refusal is what stops two copies quietly diverging in the middle.
An entry is committed when a majority of the cluster has stored it. Only then does the leader move its commit index, apply the entry to its own copy of the table and answer the client, and followers learn how far the commit index has moved on the next message they get.
Which pays off the promise lesson 032 made about today. The heartbeat that holds the leadership is an append with no entries in it. The term stamps both the election and every entry written during it. You win an election by having the best log, and you demonstrate you are still leader by having a majority keep accepting your entries. One number and one message do both jobs. A leader is not elected and then handed a log; it is elected by its log.
One limit, because it is the limit of the whole family. All of this assumes machines crash, lose their memory and come back with their log intact. A machine that lies, or returns with a log quietly corrupted, is a harder problem, and nothing above defends against it.
What a commit costs
Every write waits for a round trip to a majority before anybody is told it happened. Lesson 011 priced that exact round trip and called it the cost of synchronous replication: half a millisecond inside one data centre, two hundred milliseconds from Mumbai to Virginia. Galewatch's routing table takes a handful of writes a week, so nobody will ever look at the half millisecond twice. That is the happy case and it is the common one, because the things worth agreeing about are usually rare.
An even sized cluster is the next bill and it is a strange one, because you pay it for nothing. Four machines have a majority of three and survive one failure, which is what three machines already did. Then count the ways to lose. Call the chance a given machine is down right now one percent, a round number chosen for the ratio rather than as a claim about anybody's hardware. Three machines lose their quorum when two of the three are down: three pairs, each at one in ten thousand, so about three in ten thousand. Four machines lose it when any two of four are down, and four machines make six pairs, so about six in ten thousand.
| Machines | Majority | Survives | Quorum lost |
|---|---|---|---|
| 3 | 2 | 1 failure | 3 in 10,000 |
| 4 | 3 | 1 failure | 6 in 10,000 |
| 5 | 3 | 2 failures | 1 in 100,000 |
The ratio does not depend on the one percent. Three pairs became six pairs, so the fourth machine doubles your chance of losing the quorum at any failure rate small enough for this arithmetic, and it buys no extra tolerance at all. Five machines need three of five down, which is ten triples rather than three pairs, about one in a hundred thousand, thirty times better than three. The useful sizes are three and five, and the step between them is paid for in commit latency.
Geography is the bill nobody budgets. Put the three machines one per region, leader in Mumbai, followers in Singapore and Virginia. A majority is the leader plus one follower, so every commit waits for the nearer of the two. Singapore is about four thousand kilometres, which on lesson 002's rule of a millisecond per hundred kilometres round trip is forty milliseconds. That is a floor, and lesson 002 says so in the same breath: fibre does not run in straight lines, so expect two to four times the straight line number. Use forty anyway, because it flatters the design and still ruins it. Forty milliseconds a commit taken in series is twenty five commits a second for the whole cluster, forever, and the real wire will not give you twenty five. No hardware moves it, because the limit is the speed of light in glass. Then the sharp edge: the furthest member is free until the day the nearest one dies, at which point commits go from forty milliseconds to two hundred and nothing is down. A quorum's latency is set by its second fastest member, and that member can change without anybody being paged. Lesson 050 owns the rest.
Cut three machines into two and one and the lonely one can neither elect nor commit, which is lesson 024's two answers with no third arriving as a protocol rule rather than an architectural choice. It can still serve reads if you let it, stale by an unbounded amount, which is lesson 023's vocabulary and lesson 023's argument.
Reads are the part most explanations skip. Writes are safe; reads are not free. A leader answering a read out of its own memory is doing precisely what box 2 did in lesson 032, acting on a title it may already have lost, and a leader that has been partitioned away and deposed will answer reads confidently until it hears otherwise. The honest fixes are a round trip to the majority before every read, which Raft calls a read index and which makes a read cost what a write costs, or a lease on the leadership, which is cheap and brings back every clock problem lesson 032 spent a section on. Whichever your library does by default, you want to know which.
Who should run one, and who already does
Galewatch does not need to write any of this, and three things fall out of the May afternoon in increasing order of effort.
An alert when the primary can see fewer than two streaming standbys. One query against the primary's own view of its replicas. It would have gone off at 13:55 for the restart, which is the false alarm you expect and silence for ten minutes, and it would have been back on at 14:05 with nobody having touched route-2, which is not one. Its absence is the whole incident. Galewatch had an alert for the routing database being down, bought in April, and none at all for it being down to one copy. Lesson 026 called Marlow's monitoring a museum of its own outages, and this is how the exhibits get added.
Then the runbook stops naming a machine. Before promoting, read the last received log position off both standbys and promote the higher one:
-- run on each standby; the question the runbook never asked
select pg_last_wal_receive_lsn();
route-2 3A/7C0001F8
route-3 3A/8E0004B0 <- further ahead, so promote this one
Two log positions, and the one that is further ahead wins. That is Raft's vote rule performed by hand, with what lesson 032 called a human majority of one, and for a three machine cluster that fails over twice a year it is the right answer.
Then the real one, which Galewatch has not built. Doing it by hand needs a person awake; automating it means something has to decide, and a thing that decides needs a majority of its own or it is the fifty lines of Python lesson 024 deleted. That is why serious failover controllers keep their own state in a consensus cluster rather than in the database they are failing over. Lesson 031's question about who is in the cluster is now answerable, at least in shape: it is an entry in a log that a majority accepted. Lesson 048 owns getting that answer to the four ingest boxes.
Marlow Books, the four person online bookshop in this course, should run none of it, and today sharpens why. Lesson 032 said the founder has no automatic failover because there is no third machine and no fence. The third machine is not a nice-to-have for when there is budget: with two machines every majority is both of them, so one failure stops the cluster, which is strictly worse than the warm replica and the forty one minute rebuild lesson 025 measured.
Two machine high availability is a story people tell.
If Marlow ever moves that Postgres onto somebody's managed service, the automatic failover on the feature list is one of those controllers doing the promotion for them, and they will be depending on consensus without having read a word about it.
Stagefront, the ticketing service in this course where a stadium show goes on sale at ten in the morning, is the interesting refusal. A seat claim is a decision whose second copy contradicts the first, which is lesson 032's test, so it looks like something consensus is for. Price it. One consensus round per claim is at best one data centre round trip, half a millisecond, serialised through one leader: two thousand claims a second before anybody touches a disk. Lesson 017 put 3,333 requests a second at the door, against the five hundred a second lesson 004 measured the purchase path converting, so two thousand lands awkwardly between the two at the one minute of the year the business exists for. Lesson 015 had already settled it without a network call: claim the seat in one fast transaction on one machine, do the slow unreliable thing with nothing locked, confirm it in a second. And lesson 024 named the trade: keeping the seats on one machine is how Stagefront escapes this whole argument, and capacity is what it pays.
Where consensus would earn its keep there is that one machine being one machine. A three node group holding the seats survives losing one, and every claim then pays a round trip and a durable write at both ends instead of a local commit. At five hundred a second that fits; at 3,333 it does not, and lesson 021's limiter and lesson 017's queue are already why the inner tier never sees 3,333.
Recap
A write quorum and a vote quorum have to add up to more than the cluster. Lesson 011's any one of two standbys puts every acknowledged write on two of three machines, which is a write quorum nobody named. Decide the promotion by two of three as well and the sets must overlap, so a voter holding the newest entry can refuse the candidate missing it. Five machines with any one of four need not overlap at all, and that arrangement looks more careful than the one that works.
Agree on the order, not the value. Machines that apply the same changes in the same order to the same starting table hold the same table by construction. That is state machine replication, and lesson 012's replica replaying the write ahead log in log order is one you have already met.
Raft's term is lesson 032's fencing token used as the system's only clock. A numbered period with at most one leader in it, carried on every message, refused when old and obeyed when new. The protocol never asks what time it is, which is what lets a machine answer on its own.
A leader is not elected and then handed a log; it is elected by its log. A candidate wins by being at least as up to date as the majority it asks, an entry commits when a majority has stored it, and the heartbeat holding the leadership is an append with nothing in it.
The fourth machine doubles your chance of losing the quorum. Three pairs become six pairs and the tolerance stays at one failure. Three and five are the sizes worth having, and five costs latency to buy about thirty times the margin.
A quorum's latency is set by its second fastest member. One machine per region with the leader in Mumbai means every commit waits at least forty milliseconds for Singapore, and twenty five a second is a ceiling the real fibre will not reach. Virginia costs nothing until Singapore dies, and then it is two hundred milliseconds and nothing has gone down.
Check your understanding
A colleague is setting up a five machine Postgres cluster for a configuration database and proposes requiring any one standby, reasoning that five machines are safer than three. Work out which promotions are unsafe, say what you would change, and name the one alert you would add before changing anything.
Somebody argues that Raft would have prevented Galewatch's May incident. Say what would have been different at 14:41, what would have been exactly the same at 13:58, and which of the two you think is more dangerous to leave unfixed.
You inherit a three node cluster with one node in each of three regions and a leader whose commits take forty milliseconds. The nearest follower is taken out for a planned rebuild. Say what happens to your write latency, what your monitoring will show, and what you would have put in place last month.
A service reads from its Raft leader's memory on every request, on the grounds that the leader is authoritative. Name the lesson 032 mechanism this repeats, give the two ways to make the read honest, and say which you would ship for a config store read ten thousand times a second.
Stagefront wants its seat machine replicated by a three node consensus group so that losing it is survivable. Use lesson 017's arrival rate and lesson 004's conversion rate to say at which of the two the design fits, and name what has to sit in front of it for that to be the rate that arrives.
Marlow's founder reads about Raft and proposes running it across box A and box B, the two machines the shop already has. Give the arithmetic that makes this worse than what they have today, then say what you would spend the same money on instead.
Next lesson
034 Clocks, Ordering and Why "Before" Is Hard. Today put an entire protocol on a counter rather than a clock and never once asked what time it was; next lesson asks what the clocks you do have can and cannot tell you about which of two things happened first.