TT Lab
Get started
Learn Learning paths Courses

There Were Two Leaders

Receiving and committing are not the same

Continue in TT Lab

One-line summary

The fact that the leader received a value and wrote it to its own log guarantees nothing. Only when a majority holds the same entry at the same position, and that entry belongs to the leader's current term, is it "committed," and only then can you tell the client it succeeded.

Why this was needed

Once an election ends, there is a single place that accepts writes. But how does that one spread the values it receives to the others? If you simply "send and forget," a value is nowhere to be found if the leader dies right after receiving it. Conversely, if you "wait until everyone has received it," the whole system stalls when just one node is slow. The point that consensus chooses is the majority in between.

The reason for choosing a majority is safety, not performance. Two majorities always overlap, so if an entry has reached a majority, at least one of the candidates who could become the next leader holds it. Add one more condition (a vote is granted only when the candidate's log is at least as up to date as the voter's) and only a node holding that entry can become leader. So a committed entry survives in any future leader.

How it works

The log is an append-only list, and each entry carries the term at that time. That term is the only basis for comparing two logs.

index :   1        2        3        4
term  :   1        1        2        3
value : "a"      "b"      "c"      "d"

When the leader sends entries to a follower, it says, "If position 2 of your log is term 1, append these after it." The follower checks only that preceding position. The reason it does not compare the whole log is the Log Matching property: the rules imply that if the same position holds an entry of the same term, everything before it is identical. If the preceding position does not match, the follower rejects, and the leader steps back one position and asks again. Once it finds the matching point, it fills in from there with its own entries.

The leader does the counting for commits. It tallies how far each follower has received, and uses the highest position held by a majority as the commit index. There is a caveat that is almost always forgotten.

리더는 자기 임기의 항목이 과반에 닿았을 때만 커밋 번호를 올린다.
앞 임기의 항목은 그 자체로는 과반에 닿아도 커밋하지 않는다.

Without this caveat, the following happens. An entry from an earlier term was replicated to a majority and committed, then that leader dies, and another node that does not have the entry (because it holds an entry of a higher term) becomes leader and overwrites that position with a different entry. A value that was already declared final disappears. The paper explains this case with Figure 8 and adds the caveat above as the solution. The moment a new leader commits even one entry from its own term, everything before it becomes safe along with it.

The leader keeps one more value separately for each follower: "where do I start sending from next for this follower." When a new leader is elected, it optimistically sets this to the end of its own log, and reduces it by one each time it is rejected. Most followers find the matching point within one or two tries, so this optimism is normally free. In exchange, a node that has been away for a long time needs a long rewind, and that is why snapshots become necessary.

What it looks like in practice

This distinction changes the meaning of the response a client receives. If an etcd write returned success, the value has reached a majority, so it survives even if the leader dies right after. Conversely, if the response was a timeout, the write may have succeeded or may have failed, because it could have reached a majority and only the response was lost. So code that uses a distributed store must make its retries idempotent. That is why the etcd documentation puts transactions and revisions up front (etcd: why etcd).

Catching up is also a scene you see every day. When a follower restarts, the leader steps the position number it sends back one at a time to find the matching point. If the log is long, this rewind takes a long time, so real implementations also use snapshots. etcd compacts the log once a certain number of entries accumulate, and the default of --snapshot-count, which sets that number, changed from 10,000 to 100,000 in v3.2 (etcd maintenance docs).

The third thing you often see is "one slow follower." Commits proceed as long as a majority exists, so one slow node does not block the service. That is why it goes unnoticed and only shows up one day when another node restarts, because at that moment there is only one healthy node left to form a majority. Replication lag should be read not as a failure but as the exhaustion of spare capacity.

What you will do in the next lab

You write the four log rules yourself: appending, checking log matching, cutting off from the position where entries diverge, and computing the commit index. In the last two steps, you write values to the three nodes and watch the commit spread, then kill and revive one follower and confirm the process of it catching up from an empty log.