// OSSeva Blog
OperationsThe Raft Consensus Algorithm Explained: How It Works and Where It Runs
The short answer
Raft is a consensus algorithm: a way for a small cluster of servers to agree on one ordered log of commands, so that every server applies the same commands in the same order and ends up in the same state. Diego Ongaro and John Ousterhout published it at Stanford in 2014, in the paper "In Search of an Understandable Consensus Algorithm". A shorter version won a best paper award at that year's USENIX Annual Technical Conference.
The paper says the Raft consensus algorithm produces a result equivalent to multi-Paxos and is as efficient. Its goal was to be easy to understand and to implement correctly, which it gets from a strong leader, randomised election timers and a short list of rules. It now sits under etcd (and so Kubernetes), Consul, Kafka's KRaft mode, ClickHouse Keeper and Apache Ratis.
What a consensus algorithm solves
Distributed systems need a few facts every node agrees on: who the leader is, which configuration is current, where each partition lives. One server holding them is a single point of failure. Several servers holding copies need a rule for agreeing on changes when some crash or the network splits, without ending up with two leaders in a split-brain. That rule is the consensus algorithm. Raft's answer is a replicated log: once a majority of servers has stored an entry, it is committed and can never be lost or reordered. Five servers keep working with two down; three tolerate one.
How the Raft algorithm works
Roles and terms
Every server is in one of three states: leader, follower or candidate. In normal operation there is one leader and the rest are followers. Followers are passive. They answer the leader and candidates, and any client that reaches a follower is redirected to the leader.
Raft divides time into terms, numbered with consecutive integers. Each term starts with an election. If a candidate wins, it leads for the rest of the term. If the vote splits, the term ends with no leader and a new term begins. Terms act as a logical clock: a server that sees a higher term than its own updates its term at once, and a leader that sees a higher term steps down.
Leader election
The leader sends a periodic heartbeat, an AppendEntries RPC with no log entries. A follower that hears nothing for an election timeout assumes the leader is gone. It increments its term, votes for itself and asks every other server for a vote with the RequestVote RPC. Each server votes for at most one candidate per term, first come first served. A candidate that gets votes from a majority of the full cluster becomes leader and starts sending heartbeats.
Split votes are handled by randomness. Each server picks its election timeout at random from a fixed range; the paper's example is 150 to 300 ms. Usually one server times out first, wins, and sends heartbeats before anyone else starts an election.
There is one more rule, and it carries most of the safety. A server refuses its vote to a candidate whose log is less up to date than its own. So only a candidate holding every committed entry can win.
Log replication
The leader appends each client command to its log and sends it to the followers in AppendEntries RPCs. An entry is committed once the leader that created it has replicated it on a majority of servers. That also commits every earlier entry. The leader includes its commit index in later RPCs, and each follower then applies committed entries to its state machine in log order.
Each AppendEntries call carries the index and term of the entry just before the new ones. A follower whose log does not match rejects the call, and the leader walks back until the logs agree, then overwrites the follower's conflicting entries. Log entries only ever flow from the leader to the followers.
Safety properties
- Election safety: at most one leader per term.
- Leader append-only: a leader never overwrites or deletes its own entries.
- Log matching: if two logs hold an entry with the same index and term, they are identical up to that entry.
- Leader completeness: a committed entry is present in the log of every later leader.
- State machine safety: no two servers apply different entries at the same index.
The paper also covers membership changes, using a joint consensus phase in which the old and new configurations must both agree, and log compaction through snapshots.
Raft vs ZooKeeper's ZAB
ZooKeeper does not use Raft. It uses its own protocol, ZooKeeper Atomic Broadcast (ZAB), which is older and was designed for ZooKeeper alone. Both are leader-based and majority-based, so they look alike from a distance.
| Raft | ZAB (ZooKeeper) | |
|---|---|---|
| Leader generation | Term number | Epoch, the high 32 bits of the zxid; the low 32 bits count proposals |
| Leader election | Part of the algorithm; randomised timeouts; candidate needs an up-to-date log | Leader must have seen the highest zxid, then becomes active only after a quorum has synced with it |
| Commit | Majority has stored the entry; commit index piggybacked on AppendEntries | Leader sends a COMMIT once a quorum has acknowledged the proposal |
| Log direction | Leader to followers only | The published description moves entries to and from the leader during recovery |
| Reads | Linearizable by default in etcd; serializable reads optional | Served locally by each server; a client calls sync() first if it needs the latest value |
The last row matters most to operators. A ZooKeeper client can read a slightly stale value from its own server, while etcd makes reads linearizable by default and lets clients opt into faster serializable ones.
Where Raft runs
- etcd uses the etcd-io raft library, a Go implementation of the paper. Kubernetes stores all API server data in etcd, so every Kubernetes cluster runs Raft underneath. Kubernetes itself does not.
- Consul servers form a Raft peer set through the hashicorp/raft library. Client agents are not in the peer set, which is why three or five servers can serve thousands of clients.
- Kafka KRaft replaced ZooKeeper with a Raft-based metadata quorum. KIP-595 calls it a Raft dialect: pull-based where Raft pushes, and using Kafka's offset and epoch in place of index and term. Kafka 4.0 removed ZooKeeper mode entirely.
- ClickHouse Keeper is written in C++ on eBay's NuRaft library and speaks the ZooKeeper client protocol. See ClickHouse Keeper vs ZooKeeper.
- Apache Ratis is a Java Raft library meant to be embedded in applications, not a standalone server like ZooKeeper or Consul.
Why most products still sit on ZooKeeper
New coordination layers default to Raft. The installed base is another matter. Solr 10, released in March 2026, still requires ZooKeeper. HBase servers still need it, Hadoop's automatic NameNode failover runs through a ZooKeeper quorum, Storm 3.1 needs it, Druid still defaults to it, and Pinot depends on it through Apache Helix. NiFi 2.x defaults to ZooKeeper-based leader election, with Kubernetes as an option.
The coordination layer is product code, not configuration. These systems are written against ZooKeeper's znodes, ephemeral nodes and watches, usually through Apache Curator, so only the product's own developers can replace it, as Kafka did over several releases. Most of the Hadoop-era stack has not. So the ZooKeeper ensemble stays, often as an end-of-life version that nobody has looked at for years. Our ZooKeeper dependency map lists which products need it and which version each one bundles.
Moving your coordination layer to Raft
Where the product supports it, the move is a planned migration: ZooKeeper to KRaft for Kafka, ZooKeeper to ClickHouse Keeper, or ZooKeeper to etcd for Patroni. OSSeva runs those migrations through its ZooKeeper to Raft migration service. Where the product still needs ZooKeeper, OSSeva ships patched, signed ZooKeeper builds for 3.4, 3.5, 3.6 and 3.7 today, supports 3.8 and 3.9, and provides CVE attestation for auditors; see ZooKeeper extended support. For a side-by-side of the coordination services, read ZooKeeper vs etcd, Consul and KRaft.
Frequently asked questions
Is Raft a gossip protocol?
No. Gossip protocols spread state between peers with eventual consistency and no leader. Raft is a leader-based consensus protocol: nothing is committed until a majority has stored it.
Does Kubernetes use Raft?
Indirectly. Kubernetes keeps its cluster state in etcd, and etcd uses Raft. The Kubernetes components themselves do not run a Raft quorum.
How many nodes does a Raft cluster need?
An odd number, normally three or five. Three tolerates one failure and five tolerates two. More servers add write latency, because each commit needs a larger majority.
What happens when the Raft leader fails?
Followers stop receiving heartbeats, one of them times out, starts a new term and wins an election if its log is up to date. Committed entries survive because every eligible candidate already holds them.
Is Raft better than Paxos?
Not in fault tolerance or performance, where the paper treats them as equivalent. The difference is structure: in Paxos, leader election sits outside the core protocol and the multi-decree form real systems need was never fully specified, while Raft builds election in and concentrates work in the leader. In the paper's user study, 33 of 43 students answered questions about Raft better than about Paxos.
Tags
Related articles
ZooKeeper Vulnerabilities by Version: CVEs in 3.4 to 3.9
September 29, 2026MigrationZooKeeper Alternatives: ZooKeeper vs etcd, Consul, KRaft and ClickHouse Keeper
September 29, 2026ComplianceWhy Your Scanner Flags the ZooKeeper Inside a Product You Bought, and How VEX Attestation Answers It
September 29, 2026Ready to get your open source under control?
Talk to an OSSeva engineer about CVE coverage, compliance, and migration support for your stack.