KnowSys
⬡Build it yourself

Build a replicated key-value store with Raft

A key-value store that runs on three or five nodes, elects a leader, replicates every write through a Raft log, survives crashes and restarts, compacts with snapshots, and stays linearizable while a test harness cuts the network apart.

Serious side project⏱ 6–10 weekendsGo · Rust · Java

etcd, Consul, CockroachDB, TiKV and Kafka's KRaft mode all rest on the same small idea: a group of machines agree on one ordered log of commands, and each machine applies that log to its own copy of the state. If a majority is alive and can talk, the group keeps working. If not, it stops instead of diverging.

Raft is the algorithm most of them use, and it was designed to be understood. Figure 2 of the paper fits the core rules on one page. Getting those rules right in code, under crashes, lost messages and partitions, is a much longer story, and that's this project. You'll finish with a KV store that three or five processes serve together, plus a test harness that tries hard to break it.

01Why build this

Consensus sits under the control plane of most modern infrastructure. Once you've implemented it:

  • Split-brain stops being a vague fear. You'll know exactly which rule keeps two leaders from both accepting writes in the same term.
  • Quorum sizes make sense. You'll see why five nodes tolerate two failures, and why a fourth node adds cost without adding tolerance.
  • "Linearizable" gets a precise meaning. You'll learn the hard way that a leader serving reads from memory can return stale data.
  • etcd and ZooKeeper outages read differently. Leader churn, slow disks causing missed heartbeats, and snapshots too large to send are all things you'll have debugged.

It's also the best project here for learning to test concurrent code. The bugs are rare, timing-dependent and silent, and finding them forces you to build real tooling.

02What you're building

A single write, from client to committed, in the finished system:

How `PUT x=1` becomes durable on a three-node cluster
ClientLeaderFollower AFollower BPut(x, 1, client 7, seq 42)AppendEntries(term 3, idx 18)AppendEntries(term 3, idx 18)success, match 18OKheartbeat, commit 18
Step 1. The client sends the write to the node it thinks is leader, tagged with its client ID and a sequence number so a retry can be recognised.
1 / 6

?Why wait for a majority, and not for everyone?

Because any two majorities of the same cluster overlap in at least one node. A committed entry lives on a majority, and a new leader needs votes from a majority, so at least one voter has the entry. Raft's election rule then refuses leadership to any candidate whose log is behind that voter's. Waiting for every node would add nothing to safety and would let one slow node stall the cluster.

03Before you start

You needWhyWhere to get it
The extended Raft paper, especially Figure 2The rules, stated exactlyraft.github.io
A language with good concurrency and RPCTimers, goroutines or tasks, and message passingGo is what MIT 6.5840 uses; Rust with tokio works well
A simulated network you controlYou need to drop, delay, reorder and partition messages on demandBuild it yourself in milestone 1
Some background on replication and failure detectionWhy timeouts can't tell slow from deadChapter 27, Chapter 30
Patience with rare bugsSome failures appear once in hundreds of runsA script that runs your tests in a loop overnight

04The roadmap

Eight milestones, roughly in the order of the MIT 6.5840 labs. Keep every earlier test passing as you add later ones.

1

A network you can break

1 weekend

Before any Raft code, write the harness: nodes in one process, sending messages through a router that can drop, delay, duplicate and reorder them, and cut the cluster into groups. Seed its randomness so a failing run can be replayed.

It feels like a detour. Probably every serious Raft bug you'll find this month gets found here, so it's time well spent.

You’ll learnsimulated RPCmessage losspartitionsdeterministic seeds
Done when: A test can partition node 2 from the others, drop 20% of messages, and reproduce the same run from a seed.
2

Leader election

1–2 weekends

Every node starts as a follower with an election timer. If it hears nothing from a leader before the timer fires, it becomes a candidate, increments its term, votes for itself and asks the others. A node grants at most one vote per term. Whoever wins sends empty AppendEntries as heartbeats.

Randomise the timeout so candidates don't keep splitting the vote. Also, any message carrying a higher term makes the receiver step down to follower immediately. Most early bugs are a missing check for that.

You’ll learntermsRequestVoterandomised timeoutsheartbeats
Done when: A three-node cluster elects one leader within a couple of seconds, and a new one after the leader is partitioned away.
3

Log replication

2 weekends

The leader keeps a nextIndex and matchIndex per follower. AppendEntries carries the index and term of the entry just before the new ones, and a follower rejects it if its log doesn't match there. On rejection the leader backs up and tries again. A follower deletes conflicting entries and takes the leader's.

Advance the commit index when a majority has an entry, and only for entries from the leader's current term. That last condition is the subtle case in Figure 8 of the paper. Read it twice before you skip it.

You’ll learnAppendEntrieslog matchingcommit indexconflict resolution
Done when: Commands submitted to the leader are applied in the same order on every node, even when a follower misses a run of entries and rejoins.
4

Persistence and restarts

1 weekend

Three fields must survive a crash: the current term, who this node voted for in it, and the log. Write them to disk and fsync before replying to any RPC that changed them. A node that forgets its vote can vote twice in one term, and two leaders follow.

Backing up one entry at a time, as the paper first describes, is slow after a long outage. Add the faster version it sketches, where the follower reports the conflicting term and its first index, so the leader can skip a whole term per round trip.

You’ll learncurrentTermvotedForfsync before replycrash recovery
Done when: Killing and restarting any minority of nodes, repeatedly, during a stream of writes never loses a committed entry.
5

A key-value store on top

1–2 weekends

Each node runs a KV map and applies committed log entries to it in order. A client sends requests to the leader; followers reply with a hint of who the leader is. Only after the entry commits and is applied does the leader reply.

A client that times out will retry, and the first attempt may already be in the log. Tag each request with a client ID and sequence number, and have the state machine skip duplicates. Without that, Append runs twice.

You’ll learnstate machine replicationclient sessionsduplicate detectionleader redirects
Done when: Clients doing concurrent Put, Append and Get through leader changes see every operation applied exactly once.
6

Linearizable reads

1 weekend

A leader that has been partitioned off doesn't know it yet, and will happily answer reads from its stale map. The simplest fix is to put reads through the log like writes. That's correct, and slow.

Then try read index: record the commit index, confirm leadership with a round of heartbeats, wait until that index is applied, then answer. Leases skip the heartbeat round but depend on bounded clock drift. Know which assumption you've made.

You’ll learnstale readsread indexleader leasesclock assumptions
Done when: A test that partitions an old leader away never sees a Get return a value older than an acknowledged Put.
7

Snapshots and log compaction

1–2 weekends

Left alone, the log grows forever. When it passes a threshold, have the KV layer serialise its state, including the duplicate table, and let Raft discard every entry up to that index. Now log indexes no longer start at zero, and every array access needs an offset.

A follower too far behind asks for entries the leader no longer has. Send it the snapshot with InstallSnapshot instead. Getting indexes right after truncation is where most off-by-one bugs in this project live.

You’ll learnsnapshotsInstallSnapshotlog truncationlastIncludedIndex
Done when: The log stays below a size limit during a long run, and a node wiped clean rejoins by receiving a snapshot.
8

Real processes, real failures

1–2 weekends

Swap the simulated network for real RPC and run each node as its own process. Use iptables rules or a proxy between nodes to create partitions, and kill -9 to crash them.

Record every client operation with its start time, end time and result. Feed the history to a linearizability checker like Porcupine. If it rejects a history, the bug is real, and you have the exact interleaving that shows it.

You’ll learnRPC over TCPlinearizability checkingPorcupinefault injection
Done when: Five real processes survive a script that kills, restarts and partitions them for an hour, and a linearizability checker accepts the recorded history.

05Traps that catch everyone

SymptomCauseFix
Two leaders in the same termA vote wasn't persisted before replying, or votedFor wasn't reset on a new termPersist before replying; reset votedFor only when the term increases
Elections never settleTimeouts not randomised, or the timer resets on every message instead of only valid onesRandomise per election; reset only on a vote granted or AppendEntries from the current leader
A committed entry disappears after a leader changeCommitting an old-term entry by counting replicasOnly count replicas for entries from the current term (Figure 8)
Deadlock under loadHolding the node's lock while sending an RPCRelease the lock before any network call, then recheck the term after
A retried Append applies twiceNo duplicate detection in the state machineClient ID plus sequence number, stored in the snapshot too
Stale reads after a partition healsReads served from the leader's memory without checking leadershipRoute reads through the log, or use read index
Index out of range after snapshotsLog positions mixed up with absolute indexesWrap the log in a type that converts between the two in one place

06Stretch goals

  • Membership changes. Add and remove nodes safely, one at a time, the way Ongaro's thesis describes.
  • Sharding. Split the keyspace across several Raft groups with a configuration service that moves shards between them, as in the last 6.5840 lab. Chapter 29 covers the partitioning side.
  • Pre-vote and leadership transfer. Stop a node that rejoins after a partition from forcing a pointless election.
  • Batching and pipelining. Send many entries per AppendEntries and several in flight at once, and measure throughput before and after.
  • Deterministic simulation. Run the whole cluster in one thread with a fake clock, so every bug replays exactly from a seed.

07References worth your time

Ongaro and Ousterhout, In Search of an Understandable Consensus Algorithm

The extended Raft paper. Figure 2 is the specification you're implementing; keep it printed next to you.

MIT 6.5840, Distributed Systems

The course labs build Raft and a sharded KV store in Go with a test harness that injects failures. Probably the best structure for this project.

Jon Gjengset, Students' Guide to Raft

Written by a 6.824 teaching assistant from the bugs students actually hit. Read it before milestone 2 and again when you're stuck.

Diego Ongaro, Consensus: Bridging Theory and Practice

The PhD thesis behind Raft, covering membership changes, log compaction and client interaction in more depth than the paper.

Porcupine

Anish Athalye's fast linearizability checker in Go, used by the 6.5840 test suite. Plug your recorded histories into it in milestone 8.

Jepsen and Maelstrom

Kyle Kingsbury's analyses of real databases under partitions, and Maelstrom, a workbench for testing your own distributed system against the same kind of faults.

etcd raft library

A production Raft in Go, used by etcd and others. Read it after milestone 7 to see how a real implementation separates the protocol from I/O.

Chapters that back this project

Next projectλ a programming language→