TT Lab
Get started
Learn Learning paths Courses

There Were Two Leaders

There Were Two Leaders

Continue in TT Lab

One-line summary

When the network splits, there can be two leaders at the same moment. What Raft prevents is not that state but having two leaders in the same term, and the leader on the minority side cannot finalize anything and only later learns that its records were discarded.

Why this was needed

Incidents in distributed systems usually come not from the failure itself but from incomplete information. A leader that has been cut off does not know it has been cut off. All it knows is "responses haven't been coming lately," which is indistinguishable from the other side being dead. And a Raft leader has no mechanism for stepping down on its own. It keeps believing it is the leader until it sees a larger term.

So it keeps acting like a leader: it accepts client values, writes them to its own log, and returns success. At the same moment, on the other side, a majority gathers and elects a new leader with a higher term. At this point there really are two leaders. The starting point of this algorithm is that there is no way to prevent that; what it prevents instead is two leaders finalizing different values at the same position.

How it works

The leader on the minority side cannot get confirmation from a majority. It can write entries to its own log, but the commit index does not rise. So those entries remain "recorded but not finalized."

분할 중                        임기  로그                  커밋
  옛 리더 (혼자)                 1   alpha, ghost           1
  새 리더 + 팔로워 (둘)          2   alpha, real            2

The key is that the two entries are fighting over the same position number, 2. Only one can remain in a position, and which one will remain is already decided: the one that got confirmation from a majority.

Recovery happens quietly. When the connection returns, the old leader sees a heartbeat of term 2, raises its own term to 2, and steps down to follower. Then the new leader says, "If position 1 of your log is term 1, append real after it," and the old leader's position 2 is cut off because its term diverges. ghost disappears. No warning is raised.

Here it becomes clear why the election restriction is needed. If, right after recovery, the old leader timed out first and became a candidate, and could win votes on log length alone, it could overwrite the already committed real with ghost. The rule of looking at the term of the last entry first blocks that path.

It is also worth seeing why this design is accepted. Stopping the minority side from accepting writes is a choice to give up availability. If both sides of the split accept writes, the service keeps running, but at recovery a person has to merge the two branches of records. Systems that use Raft choose the opposite: the minority side stops, and in exchange recovery finishes automatically. For things like metadata, where "a missing value is better than a wrong one," this is the right choice.

So what the title of this course means is not "a bug was found." The moment of having two leaders is normal operation, and the promise of this algorithm is that even at that moment only one value is finalized. The place where incidents happen is not inside the algorithm but outside it: the point at which a client reads the success response returned by the minority-side leader as final.

What it looks like in practice

In operations, this event comes in the form of "the write said it succeeded but the value is gone." The client got success from the minority-side leader, and a few seconds later, when it reads, the value is not there. There is no error in the logs, because nothing failed.

Real implementations try to narrow this window. A representative approach is for the leader to step down by itself if it does not receive responses from a majority to its heartbeats for a certain time (lease-based). The etcd documentation says that if quorum is lost to a temporary partition, the cluster automatically and safely resumes when the network recovers, but permanent quorum loss is fatal (etcd FAQ). The same document has a table showing that a 3-node cluster has a majority of 2 and tolerates 1 failure, and a 5-node cluster has 3 and 2 respectively.

The advice not to use an even number of nodes comes from here too. The majority of 4 nodes is 3, so it tolerates 1 failure, the same as 3 nodes, while the number of machines that can fail goes up by one. The etcd documentation recommends keeping the cluster at seven nodes or fewer, and notes that Google Chubby suggests five.

Another thing to remember is that a "partition" does not only mean a cable being cut. A single firewall rule, a single security group change, an asymmetric path where only one direction is blocked, a long garbage collection pause: all appear in the same shape. The nodes are alive and the processes are healthy, but messages just do not get through. The case where only one side is blocked is especially nasty: if heartbeats go out but responses do not come back, the leader believes it is talking, while the follower hears nothing and opens an election. That is why this lab does not cut any wire either. Rejecting messages is enough.

What you will do in the next lab

You detach the leader from the other two of the three nodes and then bring it back. Since the Pod has no permission to set up a firewall, instead of cutting the wire you build a switch that rejects each other's messages. You record the moment when two leaders exist at once, write one value to each side, and confirm with your own eyes which side disappears on recovery.