Skip to main content
Create your own
Lesson illustration

Understanding Raft: Distributed Consensus Explained

Welcome to the first lesson of our "Advanced Distributed Concepts" module. In the previous module, we completed our tour of the three pillars of observability by building a centralized logging pipeline. You learned how to see what happened, where, and when across your services. Now, we shift our focus from observing failures to building systems that are fundamentally resilient to them.

Today, your learning outcome is to explain the need for distributed consensus and the basic mechanics of the Raft algorithm. In the systems you've built, a primary database failure might have triggered a pager alert for you or your team to perform a manual failover. At scale, this process must be automatic, fast, and, most importantly, correct. An incorrect failover can lead to catastrophic data corruption. We will explore why this problem of "agreement" is so critical and how Raft, a popular and understandable consensus algorithm, provides an elegant solution.

The Need for Automated, Correct Agreement

Imagine a replicated database cluster (like PostgreSQL with streaming replication) where you have one primary node handling writes and several replicas for fault tolerance. If the primary node fails, one of the replicas must be promoted to become the new primary. But which one? And how do all the other nodes agree on this new leader?

This is the core problem of distributed consensus. If the nodes can't agree, you risk a "split-brain" scenario, where two different nodes both believe they are the leader. This leads to them accepting different writes, causing the database state to diverge and become corrupt.

To understand why automating this leader transition is crucial and how it relates to the broader problem of consensus, let's start with a foundational video from Martin Kleppmann.

Distributed Systems 6.1: Consensus

This excerpt from Martin Kleppmann's "Distributed Systems" series explains why manual failover is insufficient for modern systems and introduces the challenge that consensus algorithms aim to solve.

Please watch the section from the beginning until the introduction of consensus. This will frame the problem by connecting the concept of state machine replication (which you've used with databases) to the need for automated leader election.

As the video explains, relying on a human to intervene is too slow and error-prone for highly available systems. The system itself must be able to detect a leader's failure and elect a new one. This automatic election process is what consensus algorithms provide. They are the bedrock of reliability for many critical systems you'll encounter, including:

  • etcd: The configuration store for Kubernetes.
  • CockroachDB & TiDB: Distributed SQL databases.
  • Consul: A service discovery and configuration tool.

Knowing how consensus works is not just theoretical; it's a practical necessity for understanding how these modern, scalable tools achieve their fault tolerance guarantees.

Understanding Raft

This article clearly articulates why consensus is fundamental to building fault-tolerant distributed systems using replicated state machines.

Read the first three sections, from the introduction down to the end of the "Why do we need consensus?" section. This will solidify the connection between the practical problem (keeping a service running) and the abstract solution (a consensus algorithm).

Introducing Raft: Consensus Designed for Understandability

While Paxos was the original groundbreaking consensus algorithm, it is notoriously difficult to understand and implement correctly. Raft was created with a primary goal of being more understandable while providing the same safety guarantees. It achieves this by breaking the problem down into three more-or-less independent parts:

  1. Leader Election: Electing one node from the cluster to be in charge.
  2. Log Replication: Ensuring all nodes have an identical, ordered log of operations, managed by the leader.
  3. Safety: A set of rules that guarantee the system's correctness, even in the face of failures and network partitions.

Let's watch a brief, high-level overview of how these pieces fit together.

Understand RAFT without breaking your brain

This video, "Understand RAFT without breaking your brain," provides an excellent and intuitive introduction to the core ideas of Raft.

Watch the first two sections, covering the problem of inconsistency and introducing the log as the source of truth. This will give you a simple mental model before we dive into the mechanics.

Raft's Core Concepts: States and Terms

To orchestrate its work, Raft defines a few key concepts. At any given moment, every server in the cluster is in one of three states:

  • Follower: The default state. Followers are passive; they simply respond to requests from the leader and candidates. They don't initiate any communication.
  • Candidate: A temporary state used during an election. A follower becomes a candidate when it believes the leader has failed.
  • Leader: There is at most one leader at any time. The leader handles all client requests and manages the replication of the log to followers.
This diagram illustrates the transitions between the Follower, Candidate, and Leader states. A follower times out and becomes a candidate to start an election. A candidate can win the election to become leader, or it can discover another leader and revert to being a follower.

Raft also divides time into Terms. A term is a monotonically increasing number that acts as a logical clock. Each term begins with an election. If an election succeeds, that term has a single leader. If it fails (due to a split vote), the term ends with no leader, and a new election (and a new term) begins.

Terms are crucial for maintaining consistency. If a server receives a message with a term number higher than its own, it knows its information is outdated. It immediately updates its term and, if it's a leader or candidate, steps down to become a follower. This simple rule is how Raft deals with stale leaders that might have been isolated by a network partition.

Mechanic 1: Leader Election

The leader election process is triggered by a follower timing out. Here's how it works:

  1. Timeout: While it's a follower, a server listens for periodic "heartbeat" messages from the current leader. If it doesn't receive a heartbeat within its election timeout, it assumes the leader has crashed.
  2. Become a Candidate: The follower increments its current term number, transitions to the Candidate state, and votes for itself.
  3. Request Votes: The new candidate sends a RequestVote RPC (Remote Procedure Call) to all other servers in the cluster, asking them to vote for it.
  4. Voting: Other followers will grant their vote to the candidate only if they haven't already voted in the current term. They also perform a safety check to ensure the candidate's log is at least as up-to-date as their own.
  5. Winning the Election: A candidate becomes the new leader if it receives votes from a majority of the servers in the cluster (e.g., 3 out of 5). Once it becomes leader, it immediately starts sending heartbeats to all other servers to establish its authority and prevent new elections.

What happens if multiple followers time out and become candidates at roughly the same time? They could split the vote, with no single candidate achieving a majority. To prevent this from happening indefinitely, Raft uses randomized election timeouts. Each follower waits for a random period (e.g., between 150-300ms) before starting an election. This staggering makes it highly probable that one follower will time out, become a candidate, and win the election before others even begin.

System Design Interview: Raft Consensus Algorithm Deep Dive – Tech Interview Dot Org

This article is framed for system design interviews and provides a concise, clear explanation of Raft's mechanics.

Please read the sections Server States and Terms and Leader Election. Pay close attention to the role of randomized timeouts and the majority vote requirement. This directly addresses common interview questions about Raft's resilience.

Mechanic 2: Log Replication

Once a leader is elected, it's responsible for managing the replicated log. The goal is to ensure that every server eventually has an identical copy of the log, which contains the sequence of commands to be applied to the state machine (e.g., your database).

This diagram shows the leader (Term 5) receiving a client request, appending it to its log as entry 4:5, and replicating it to Followers A and B. Follower C is partitioned and does not receive the entry. The entry is committed once a majority (the leader plus Follower A) have acknowledged it.

The log replication process follows these steps:

  1. Client Request: A client sends a command (e.g., SET x = 10) to the leader.
  2. Local Append: The leader appends the command to its own log as a new entry. Each entry is tagged with the current term number. At this point, the entry is uncommitted.
  3. Parallel Replication: The leader sends AppendEntries RPCs to all followers, asking them to add the new entry to their logs.
  4. Majority Acknowledgment: The leader waits for responses. Once a majority of followers confirm they have appended the entry, the leader considers the entry committed. This is the point of no return. A committed entry is guaranteed to be durable and will eventually be executed by all available servers.
  5. Apply and Respond: The leader applies the command to its own state machine and sends a success response back to the client.
  6. Notify Followers: In subsequent AppendEntries RPCs (including heartbeats), the leader informs the followers of the latest committed entry index. Followers then apply all committed entries to their own state machines in order.

This mechanism ensures that even if the leader crashes right after an entry is committed (but before it has notified all followers), the new leader elected will have that committed entry in its log (due to the voting safety rules), guaranteeing it is never lost.

Understanding Raft

This section of the "Understanding Raft" article details the log replication process, including how inconsistencies are handled.

Please read the section on Log Replication. Focus on the distinction between appending an entry and committing it, and the role of the majority in making an entry durable.

Conclusion

In this lesson, we've unpacked one of the most fundamental algorithms in distributed systems. You've learned that building reliable, automated systems requires a formal way for servers to agree, a problem known as distributed consensus. You've also explored the basic mechanics of Raft, a practical and widely-used algorithm for achieving it.

Key Takeaways:

  • The Need for Consensus: To build fault-tolerant replicated systems (like databases or configuration stores), we need an automated way to handle leader failure and prevent "split-brain" scenarios that lead to data corruption.
  • Raft's Approach: Raft provides consensus through a combination of leader election and log replication. It is designed to be more understandable than its predecessor, Paxos.
  • Leader Election: A follower becomes a candidate if its election timeout fires. It wins by getting a majority vote in a given term. Randomized timeouts are a key mechanism to prevent election deadlocks.
  • Log Replication: The leader replicates log entries to followers. An entry is committed only after it has been stored on a majority of servers, making it durable against failures.
  • Safety: The combination of majority rules and term numbers ensures that there is at most one leader per term and that committed log entries are never lost.

For your interview preparation, the "Frequently Asked Questions" section of the "Raft Consensus Algorithm Deep Dive" article is an invaluable resource. It directly addresses questions like "How does Raft prevent split-brain?" and "What happens to writes when a leader fails?"

System Design Interview: Raft Consensus Algorithm Deep Dive – Tech Interview Dot Org

This is an excellent resource for self-study to prepare for interviews.

I recommend reading through the FAQ section to solidify your understanding of these critical scenarios.

Having mastered how to achieve consistency within a single distributed cluster, we are now ready to zoom out. In our next lesson, we will explore techniques for managing systems that span multiple geographic regions, focusing on DNS load balancing and GeoDNS.

Can't find a good explanation? Sign up and we'll make it for you

Sign up