TT Lab
Get started
Learn Learning paths Courses

There Were Two Leaders

The log grows forever

Continue in TT Lab

One-line summary

If you keep the entire agreed-upon order, restart time and disk usage grow in proportion to the length of the log. A snapshot stores "the result of applying everything up to here" as a single copy and lets you discard the log before it, and in exchange the recovery procedure and the meaning of a backup change.

Why this was needed

A consensus algorithm does not store values. It stores the order of commands. The premise of a replicated state machine is that applying that order again from the beginning always yields the same state, and raft.github.io explains this by saying that a consensus algorithm is used to agree on the commands in the servers' logs (raft.github.io).

The problem is that the log never shrinks. A cluster that receives 100 requests per second accumulates 8.6 million entries a day. More painful than a full disk is the restart time. If the log has to be reapplied from the beginning, restarting a single node can take days. Catching up is the same. When you add a new member or revive a node that was down for a long time, if the leader has millions of entries to send, that recovery is practically impossible.

How it works

The solution is simple. Store the result of applying everything up to some point in one piece, and discard the log before it. Along with what you store, you write down the last position number included and the term at that position, because you need both to attach the log that comes after.

  버린다                         남긴다
 [1..20000 로그]  →  스냅샷(last_index=20000, last_term=7) + [20001.. 로그]

This creates a new problem. If the leader wants to send entry 20001 to a follower that only has up to 19000, the leader no longer has the log to send. So real implementations provide a separate path for transferring the snapshot itself instead of the log. When the lag exceeds what the log can fill, it sends the whole snapshot, and from then on continues with the log.

In etcd this knob is --snapshot-count. Compaction happens when the raft entries held in memory reach this number, and its default changed from 10,000 to 100,000 in v3.2. Cleaning up the past revisions of keys is a separate knob, automatic compaction, which you set with --auto-compaction-mode and --auto-compaction-retention. Defragmentation (etcdctl defrag) is yet another thing, and the documentation warns that defragmenting a live member blocks reads and writes while the state is being rebuilt (etcd maintenance).

The meaning of a backup changes too. Saving an etcd snapshot is etcdctl snapshot save, and restoring is etcdutl snapshot restore. What matters is that a restore does not bring the original cluster back to life. The documentation says that a restore overwrites some of the snapshot metadata (the member ID and cluster ID) and that the member loses its former identity, and therefore that to start a cluster from a snapshot you must start a new logical cluster (etcd disaster recovery).

If you keep the three similarly named things apart, you will not get confused in operations. Log compaction shrinks the raft log, automatic compaction shrinks the past revisions of keys, and defragmentation returns the space freed that way to the disk. No matter how much you run the first two, the file size does not shrink, and no matter how much you run defragmentation, the revisions do not shrink.

What it looks like in practice

The most common incident is taking a backup and never doing a restore. What the sentence above means is that a restore is not "rolling back to yesterday's state" but standing up a new cluster. The member IDs change, so existing members must not get mixed in, and every member must be restored from the same snapshot. If you read this procedure for the first time on the day of the outage, it is already too late.

The second is the gap between the time of the snapshot and the actual data. The documentation warns that if you take a snapshot from the member/snap/db file, you "may lose data that has not been recorded yet but is in the wal folder." That is why backups are taken with snapshot save rather than by copying files.

The third is the relationship between compaction and backup frequency. The more often you compact the log, the faster restarts are, but creating a snapshot is itself a cost, so responses slow down briefly each time. Conversely, if you compact rarely, things are quiet in normal times and you pay all at once when you restart. Which is better is determined by "how often do you restart," not by the documentation's default.

What to check in the next quiz

The quiz checks what a snapshot discards and what it keeps, why the last position number and term must be written together, that log compaction, automatic revision compaction, and defragmentation are different things, and why a restore creates a new cluster.