Clustering with Raft
Run three brokers as one cluster and every partition is replicated: a write is acknowledged once a majority of the brokers hold it, and the cluster keeps serving if one broker fails.
How it fits together
Each partition is its own Raft group, with its own leader. All groups on a broker share one listener and one connection to each other broker; every frame names the group it belongs to.
- Node
- Elects a leader, replicates the log, decides what is committed.
- Transport
- One authenticated connection to each peer, carrying every partition's messages.
- Apply task
- Writes committed entries into the partition's segment files and answers the waiting producer.
Start a three-node cluster
Start each broker with its own id, the address its peers reach it on, every other member, and the same secret. The membership is exactly the brokers you list; there is no discovery.
export ES_RAFT_SHARED_SECRET='a long random string, the same on every node' es-broker --data-dir /var/lib/es --bind 10.0.0.1:9000 --auth required \ --raft-node-id 1 --raft-bind 10.0.0.1:6000 \ --raft-peer 2=10.0.0.2:6000 \ --raft-peer 3=10.0.0.3:6000
es-broker --data-dir /var/lib/es --bind 10.0.0.2:9000 --auth required \ --raft-node-id 2 --raft-bind 10.0.0.2:6000 \ --raft-peer 1=10.0.0.1:6000 \ --raft-peer 3=10.0.0.3:6000
es-broker --data-dir /var/lib/es --bind 10.0.0.3:9000 --auth required \ --raft-node-id 3 --raft-bind 10.0.0.3:6000 \ --raft-peer 1=10.0.0.1:6000 \ --raft-peer 2=10.0.0.2:6000
Then create topics on every broker with the same name and partition count. Each partition elects a leader within a few hundred milliseconds.
The broker refuses a configuration that cannot work: a node id of 0, a peer listed twice or sharing an address, a node listing itself, or a Raft port reachable from the network without a shared secret.
A single node
A broker with --raft-node-id and no peers is a cluster of one. It elects itself and commits immediately,
which is useful for trying the replicated write path locally:
es-broker --data-dir ./data --raft-node-id 1 --raft-bind 127.0.0.1:6000
Writing and reading
- Produce to the leader. A produce sent to a follower fails with a hint naming the leader; send it there instead.
GET /raftshows each partition's leader. - Acknowledgement means a majority has it. The leader appends the entry durably, replicates it, and answers once a majority of members have it on disk and it is applied.
- Reads come from what a broker has applied. A follower can be slightly behind the leader; reads are not linearizable.
What happens to a write
Propose. The leader turns the request into one log entry holding all its records, the offset of the first, and their timestamps. Every replica will store the same records at the same offsets.
Replicate. The leader appends the entry to its log and sends it to each follower, in batches of at most 1,024 entries or 4 MiB. A follower checks the entry follows on from what it has, appends it durably, and confirms.
Commit. Once a majority holds an entry from the leader's current term, it is committed, along with everything before it.
Apply. Each broker writes committed records into its partition files under the leader's offsets. The producer that proposed the entry gets its offsets; it is matched by a proposal id, so it is never handed another leader's entry by mistake.
A new leader appends an empty entry first, so entries left uncommitted by its predecessor are committed without waiting for the next client write.
Snapshots and catching up
The partition files already hold the data, so the Raft log does not need to keep it forever. Every
snapshot_after_applies entries (1,024 by default) a broker fsyncs the partition and drops the log entries
it has applied.
A broker that was down long enough, or a new one, may need entries the leader has already dropped. The leader then streams it the records it is missing straight from its own partition, in 1 MiB chunks, each acknowledged with the follower's new position. When the follower has everything up to the snapshot, ordinary replication resumes.
Transport and security
Brokers authenticate each other with the shared secret, without sending it: each side proves it knows the secret against a random challenge from the other, so a recorded handshake cannot be replayed. The dialing broker checks it reached the node it meant to; the listening broker checks the dialer is a configured member.
After the handshake every frame carries a MAC over a sequence number, so frames cannot be forged, replayed or reordered, and a message must name the authenticated broker as its sender. Messages from outside the membership, or with an implausibly large term, are ignored.
Without a secret the same handshake runs with no authentication, and the broker only allows that on a loopback
address. Brokers from before this transport (handshake RAFT rather than RAF2) cannot join;
upgrade every member together. Raft state from older versions is converted on first start.
Durability of Raft state
- The term and vote are in a small file replaced atomically whenever they change.
- The log is append-only: one write and one fsync per batch of entries.
- Nothing that depends on a change is sent before the change is on disk.
- If writing fails, the node stops instead of answering, so it can never vote twice in a term or confirm entries it does not have.
Checking a cluster
curl -s -H "authorization: Bearer $ADMIN_KEY" 10.0.0.1:9000/raft
/raft lists every replicated partition with its role, term, leader and commit index, and whether this
broker is replicating: as leader, a majority has caught up; as follower, it knows the leader. Ask every member; the
leader's answer proves a majority exists, and each follower's proves it is part of it.
Limitations
- No leader forwarding. Clients must send produces to the leader.
- Static membership. Adding or removing a broker means restarting every broker with the new peer list.
- No linearizable reads. A follower serves what it has applied.
- No pre-vote. A rejoining broker cannot win an election against a healthy leader, but its higher term can make the leader step down once.
- Topics are created per broker. Create each topic, with the same partition count, on every member.