Skip to the lesson
Little Builders system design, explained small

Distributed Systems Core

Imagine a group of friends in different treehouses who keep one shared scorebook. They can only talk through string telephones, and sometimes a string snaps or a friend dozes off. How do they make sure everyone writes down the same score? Who decides when they disagree? How do they make sure only one friend uses the shared paint set at a time? A distributed system is many computers working as one team, and this page is about the rules that help them agree when things go wrong.

CAP theorem

Two treehouses and a cut string

Two treehouses keep copies of the same scoreboard. When a kid changes the score in treehouse A, A tells treehouse B through a string telephone so both stay the same. CAP names three things we would love to have all at once.

Consistency (C): every question gets the latest answer, as if there were only one scoreboard. Grown-ups call this linearizable. Availability (A): every treehouse that is still working gives a real answer, not an error. Partition tolerance (P): the team keeps going even when messages between treehouses are lost, for example when the string is cut.

Now the string gets cut. A kid in treehouse B asks for the score. B knows it might be behind, so it has two choices. It can say sorry, I cannot be sure right now, please try later. That keeps every answer correct but does not answer, so it picks consistency. Or it can answer with the score it has. That always answers but might be old, so it picks availability.

Strings do break in real life: cables get cut, switches fail, and computers freeze. So partition tolerance is not really optional. The honest meaning of CAP is: while the string is cut, you must choose between always correct and always answering. While the string works, you can have both.

Two treehouses share a string telephone. Cut the string and it falls slack: each house must now pick, answer with old news or wait.

Remember

During a network split, a system must choose: refuse to answer, or answer with data that may be old.

PACELC theorem

Even on a calm day, checking takes time

CAP only talks about the bad day when the string is cut. PACELC adds the normal days. Read it like a sentence: if there is a Partition, choose Availability or Consistency. Else, choose Latency or Consistency.

Latency means waiting time. Even when every string works, making sure all treehouses agree takes time, because B has to check with A before it answers. So on a calm day the choice is: answer fast from the nearest copy, which might be a moment behind, or check with the others first, which is always up to date but slower.

Different systems choose differently. Dynamo-style databases such as Cassandra, with their default settings, choose availability during a split and speed on normal days, which is written PA/EL. Google Spanner chooses consistency both times, PC/EC, and pays for it with extra waiting. Many databases let you pick per request how many copies must agree, so you can choose fast for a like counter and careful for a bank balance.

A seesaw: one block on one end answers fast, three on the other check first and answer slower. Press an end down to choose.

Remember

Split or not, every system trades being up to date against being fast or always answering.

Consensus and majorities

Agreeing on one answer, even if some friends nap

Consensus means a group of computers agrees on one value, like who the class leader is or what the next line in the scorebook says, and once it is decided, it never changes. It must work even if a few members crash or their messages are very slow.

The secret is the majority, which grown-ups call a quorum. A choice counts once more than half of the group agrees. With 5 computers, 3 is a majority, so the group keeps working with 2 of them broken. With 3 computers, 2 is a majority, so it survives 1 broken.

Why a majority? Any two majorities of the same group always share at least one member. If one group of 3 picked red, any other group of 3 must include someone from the first group, who will speak up and say we already picked red. So two different answers can never both win.

Consensus expects members to crash or be slow, but not to lie. Handling members that lie or cheat is a different and harder problem called Byzantine fault tolerance. There is also a cost: every decision needs a majority to reply, so it is slower than one computer deciding alone, and if a majority cannot be reached, the group has to stop and wait.

  • 3 computers: need 2 to agree, can lose 1.
  • 5 computers: need 3 to agree, can lose 2.
  • 4 computers: need 3 to agree, can still lose only 1. That is why groups usually have an odd size.
Five posts in a ring. Point to raise the three nearest, a majority. Any two groups of three share a post, so two answers never both win.

Remember

Any two majorities overlap, so the group can never pick two different answers.

Paxos

Collect promises first, then ask for agreement

Paxos is the classic way to reach consensus. There are three jobs. Proposers suggest a value, acceptors vote on it, and learners find out what was chosen. One computer can do several jobs at once.

Phase one is prepare and promise. A proposer picks a ticket number, say 5, bigger than any it has used before, and asks the acceptors to listen. Each acceptor that has not already promised a bigger number promises to ignore anything with a number below 5. In its reply it also reports any value it has already accepted, and that value's number.

Phase two is accept. Once a majority has promised, the proposer asks them to accept a value with ticket 5. It is not free to choose: if any reply reported an earlier accepted value, it must propose the one with the highest number. Only if nobody reported anything may it use its own idea. When a majority accepts the same proposal, the value is chosen for good, and the learners are told.

That must-reuse rule is what keeps Paxos safe: once a value might have won, every later proposal carries it forward. Paxos is proven correct, but it is famously hard to understand and to build. Two proposers can keep interrupting each other with bigger numbers, so real systems pick one main proposer. And real systems need a long list of decisions, not one, so they run a version called Multi-Paxos with many details left to the builder.

Each post is as tall as the biggest ticket it promised. Pick a ticket by height: lower posts promise and lift, taller ones refuse.

Remember

Paxos collects promises from a majority, then gets a majority to accept, always reusing any value that may already have won.

Raft

One class leader keeps the shared notebook

Raft was designed to do the same job as Paxos but to be much easier to understand. Its big idea is to pick one leader and let the leader run everything.

Time is split into terms, numbered 1, 2, 3 and so on, like school years. Each term has at most one leader. Every change, like set the score to 7, goes to the leader. The leader writes it as the next line in its log, which is a numbered notebook of changes, and sends that line to the followers.

Once a majority of the computers, counting the leader, have stored the line, it is committed. Committed means safe: it will not be lost or changed, even if the leader crashes. Every computer applies committed lines in the same order, so all the copies end up the same.

The leader also sends small heartbeats, I am still here messages, often enough that no follower's timer runs out. The weak spots: every write goes through one leader, which can get busy, and when the leader dies there is a short pause with no leader while a new one is chosen.

Remember

Raft has one leader per term writing a numbered notebook, and a line is safe once a majority has it.

Leader election

When the leader goes quiet, someone steps up

What if the leader falls asleep for good? Each follower has a timer, like an egg timer, that restarts every time a heartbeat arrives. Each one picks a random length for its timer, for example somewhere between 150 and 300 milliseconds. If the timer rings with no heartbeat, that follower decides the leader is gone.

It becomes a candidate. It moves to the next term number, votes for itself, and asks everyone else to vote for it. Each computer votes at most once per term, for the first fit candidate that asks. The candidate that gets votes from a majority becomes the leader and starts sending heartbeats right away.

Why random timers? If every timer rang at the same moment, everyone would ask for votes at once, the votes would split, and nobody would win. With random timers, one usually rings first and wins before the others wake up. If a vote does split, they simply try again in a new term with new random timers.

Two more safety rules. A computer refuses to vote for a candidate whose notebook is behind its own, so a new leader always has every committed line. And if an old leader wakes up and sees a bigger term number, it knows it was replaced and quietly becomes a follower.

Five egg timers, each with its own amount of sand. Slide right to let time pass: the first to run out stands up as the new leader.

Remember

Random timers pick a candidate, a majority of votes makes a leader, and a bigger term number always wins.

Distributed locks

One talking stick for a crowd of computers

Sometimes only one computer at a time should do a job, like printing tonight's report or editing one file. In class, only the kid holding the talking stick may speak. A distributed lock is a talking stick for computers: a lock keeper lends it to one asker at a time.

But what if the kid holding the stick falls asleep? Nobody could speak again. So the stick is lent with a time limit. This is called a lease: the lock is yours for, say, 10 seconds, and you must renew it before the time runs out. If you crash, the lease ends on its own and someone else gets a turn.

Leases bring a sneaky danger. A computer can freeze for a while without noticing, for example during a long cleanup pause called garbage collection, or while its network is very slow. It wakes up still believing it holds the stick, even though its lease ended and someone else has it now. Suddenly two computers think they are in charge.

The fix is a fencing token. Each time the lock keeper lends the stick, it also hands out a number that only ever goes up: 33, then 34, and so on. Every write to storage must carry that number, and storage remembers the highest number it has seen. A write with 33 that arrives after a write with 34 is rejected, so the sleepy computer cannot do harm.

  • For locks that protect correctness, use a lock keeper built on consensus, such as ZooKeeper or etcd.
  • If a lock only saves effort, like not printing the same report twice, a simple lock in a single Redis server (a fast, popular in-memory store) is fine, because a rare double only wastes a little work.
  • Fencing only works if the storage checks the number. The lock alone cannot stop a frozen computer.
One key on a hook, its tag dotted with a count. Each lend adds a dot, a bigger number, so an old holder with a smaller one is turned away.

Remember

Lend the lock with a time limit, and stamp every write with a growing number so stale holders are turned away.

Quick recap

  1. Network splits will happen. During one, a system picks consistency (refuse or wait) or availability (answer, maybe with old data).
  2. PACELC adds normal days: even without a split, being up to date costs waiting time.
  3. Consensus lets a group agree on one value using majorities, and any two majorities always overlap.
  4. Paxos gathers promises, then acceptances, and must carry forward any value that may already have won.
  5. Raft uses one leader per term and a shared log, and a line is committed once a majority stores it.
  6. Leader election uses random timers, one vote per computer per term, and a majority to win.
  7. Distributed locks need leases so a crash cannot block forever, and fencing tokens so a frozen holder cannot do damage.

Grown-up words

and what they mean in plain words

Node
One computer in the team.
Partition
A break in the network, so some computers cannot hear others.
Linearizable
Acting as if there were only one copy of the data, always up to date.
Latency
Waiting time before an answer comes back.
Quorum
A big enough group to decide, usually more than half.
Term
A numbered round in Raft. Each term has at most one leader.
Heartbeat
A small I am still here message from the leader.
Log
A numbered list of changes, kept in order.
Lease
A lock that is only yours for a set time unless you renew it.
Fencing token
A growing number handed out with a lock, so storage can turn away old holders.