Search

How Consensus Algorithms Like Raft Keep Servers in Agreement

The short answer

Quick answer: A consensus algorithm lets a group of servers agree on the same sequence of operations, even if some of them crash. Raft does it by electing one leader. Clients send commands to the leader, which appends them to its log and replicates them to the other servers (followers). Once a majority has stored an entry, it is committed and applied. If the leader fails, the followers notice the silence, and one of them wins an election to take over. Because any two majorities overlap, a new leader is guaranteed to have every committed entry. A cluster of five keeps working with two servers down.

What consensus is for

Some pieces of state must be identical everywhere and must survive machine failures:

  • Who is the current primary database?
  • What is the cluster's configuration?
  • Who holds this lock?

Storing that on one server makes it a single point of failure. Storing it on several raises the question of how they stay in agreement.

The standard answer is a replicated state machine. If every server starts in the same state and applies the same commands in the same order, they all end in the same state. So the problem reduces to agreeing on one ordered log of commands. That agreement is consensus.

Raft was designed by Diego Ongaro and John Ousterhout in 2014 with understandability as an explicit goal, because the older Paxos algorithm was notoriously hard to grasp and implement. The Raft website hosts the paper and an interactive visualisation.

The basics

Each server is in one of three states:

  • Follower: passive; responds to the leader and to candidates.
  • Candidate: trying to become leader.
  • Leader: handles all client requests and replication.

Time is divided into numbered terms. Each term begins with an election and has at most one leader. The term number acts as a logical clock: any server that sees a higher term than its own immediately updates and steps down to follower. That is how a stale leader discovers it has been replaced.

Decisions need a majority (a quorum): 2 of 3, 3 of 5, 4 of 7.

Leader election

  1. The leader sends regular heartbeats to all followers.
  2. Each follower has an election timeout, set to a random value (for example, between 150 and 300 ms).
  3. If a follower's timer expires without hearing from a leader, it becomes a candidate: it increments the term, votes for itself, and asks the others for votes.
  4. Each server gives at most one vote per term, to the first valid candidate that asks.
  5. A candidate with votes from a majority becomes leader and starts sending heartbeats.

Randomised timeouts are the simple trick that makes this work. Usually one follower times out before the others, wins, and suppresses further elections. If two candidates split the vote, nobody wins, the timers reset to new random values, and the next round almost always succeeds.

There is one important restriction: a server refuses to vote for a candidate whose log is less up to date than its own. This guarantees that only a server holding all committed entries can win.

Log replication

  1. A client sends a command to the leader.
  2. The leader appends it to its own log as a new entry.
  3. It sends the entry to all followers.
  4. Each follower appends it and acknowledges.
  5. When a majority (including the leader) has the entry, the leader marks it committed, applies it to its state machine, and replies to the client.
  6. Followers learn of the commit in the next message and apply it too.

A slow or crashed follower does not hold anything up; the leader simply keeps retrying until that follower catches up.

Each replication message also carries the index and term of the entry just before the new ones. If the follower's log does not match at that point, it rejects the message, and the leader steps back and tries again until it finds where their logs agree. It then overwrites anything after that point. A follower's uncommitted, conflicting entries are discarded; committed entries are never lost.

Why it is safe

Raft's central guarantee: once an entry is committed, it will be present in the log of every future leader.

The reasoning:

  • A committed entry is stored on a majority of servers.
  • A candidate needs votes from a majority.
  • Any two majorities share at least one server.
  • That server will not vote for a candidate with a less up-to-date log.

So a candidate lacking a committed entry cannot gather enough votes.

What happens in a partition

Suppose a five-node cluster splits into a group of two (containing the old leader) and a group of three.

  • The group of three elects a new leader in a higher term and continues.
  • The old leader still accepts requests but cannot reach a majority, so it can never commit anything. Its clients time out.
  • When the network heals, the old leader sees the higher term, steps down, and its uncommitted entries are replaced.

There is no split brain. The price is that the minority side cannot make progress. Raft chooses consistency over availability; in the terms of the CAP theorem, it is a CP design.

Cluster size

ServersMajorityFailures tolerated
110
321
431
532
743

Clusters use odd sizes because an even number adds cost without adding fault tolerance. Three or five is typical. Larger clusters tolerate more failures but make every write slower, since more acknowledgements are needed.

Practical extras

  • Snapshots. The log cannot grow forever. Servers periodically save their state and discard old entries. A follower far behind receives a snapshot.
  • Membership changes. Adding or removing servers must be done carefully, through the log itself, so that two separate majorities can never form.
  • Durability. Entries and votes are written to disk before acknowledging, using the same idea as a write-ahead log.
  • Consistent reads. A leader that has been partitioned might serve stale data, so reads that must be current are confirmed with a majority or protected by a lease.

Where Raft is used

  • etcd, the store behind Kubernetes. Every object in a cluster lives there; see what Kubernetes does.
  • Consul for service discovery and configuration.
  • CockroachDB, TiKV and YugabyteDB, which run one Raft group per range of data.
  • Kafka's KRaft mode, which replaced ZooKeeper for cluster metadata.

Consensus is used for small, critical data and for coordinating who leads. It is too slow for every bulk write, so large systems use it to elect leaders and then replicate ordinary data by cheaper means. See how database replication works.

Raft and Paxos

Paxos, by Leslie Lamport, came first and underlies systems such as Google's Chubby and Spanner. The two offer equivalent guarantees. Raft differs mainly in presentation and structure: it insists on a strong leader and breaks the problem into election, replication and safety, which makes it easier to teach and to implement correctly.

Both tolerate crashed nodes, not malicious ones. Handling nodes that lie requires Byzantine fault-tolerant protocols, which are more expensive.

Frequently asked questions

What is a quorum?

A majority of the servers in the cluster. Requiring one for every decision ensures any two decisions involve at least one common server.

What happens when the Raft leader fails?

Followers stop receiving heartbeats, one times out and starts an election, and a new leader is chosen, usually within a second or so.

Why do Raft clusters have an odd number of nodes?

A 4-node cluster tolerates one failure, the same as a 3-node cluster, while requiring more acknowledgements. Odd sizes give the best fault tolerance per node.

Can Raft lose data?

Not committed data, as long as a majority of servers survive with their disks intact. Entries a leader accepted but had not yet committed can be discarded when it fails.

Conclusion

Raft turns the hard problem of agreement into three understandable parts: elect a leader with randomised timeouts, replicate a log and commit on a majority, and restrict elections so committed entries are never lost. That is enough to keep a cluster consistent through crashes and partitions, and it is why Raft sits quietly underneath so much modern infrastructure.

Related articles

Sources and further reading

Usama Muneer

Usama Muneer

Coder, Blogger, Tech Speaker & Web Technologies Enthusiast. Passionate about working on open-source Programming languages & Tools while utilizing my Product Development skills.

Your experience on this site will be improved by allowing cookies Cookie Policy