Many systems need exactly one active coordinator at a time: one node that orders writes, one scheduler that hands out jobs, one controller that reconciles state. Leader election is how a group of peers picks that node, notices when it has gone, and picks another without ever letting two act as leader in a way that corrupts data. That last clause is the hard part. Picking a leader is easy. Being sure the old one has stopped is impossible in an asynchronous network, so a correct design assumes two nodes will sometimes both believe they lead, and makes that harmless.
This article treats leader election as an architecture with two layers. The bottom layer is a consensus group such as etcd, which elects its own leader with Raft. The top layer is your service, whose replicas elect a leader by competing for a leased key in that store. We walk through how each layer works, what the timeouts actually buy you, how long a failover takes end to end, and where fencing must live for the whole thing to be safe.
Two layers of election
Most teams never implement a Raft election themselves. They run etcd, ZooKeeper or Consul, or rely on the Kubernetes API server, which is itself backed by etcd. The consensus layer gives them a linearizable key-value store whose own leader changes are mostly invisible. On top of it, application replicas run a much simpler protocol: create a key tied to a lease, and whoever holds the oldest live key is leader.
Keeping the layers apart in your head prevents a common confusion. When etcd changes leader, your service's leader does not change; etcd refreshes lease expiries on its own leader change so a store failover alone does not expire application sessions. When your service's leader process dies, etcd does not hold an election at all; a lease simply expires and a key disappears. The two layers fail, and are tuned, independently.
Two layers of election
Most teams never implement a Raft election themselves. They run etcd, ZooKeeper or Consul, or rely on the Kubernetes API server, which is itself backed by etcd. The consensus layer gives them a linearizable key-value store whose own leader changes are mostly invisible. On top of it, application replicas run a much simpler protocol: create a key tied to a lease, and whoever holds the oldest live key is leader.
Keeping the layers apart in your head prevents a common confusion. When etcd changes leader, your service's leader does not change; etcd refreshes lease expiries on its own leader change so a store failover alone does not expire application sessions. When your service's leader process dies, etcd does not hold an election at all; a lease simply expires and a key disappears. The two layers fail, and are tuned, independently.
How a Raft election works
In Raft every node is a follower, a candidate or a leader, and time is divided into numbered terms. The leader sends heartbeats (empty AppendEntries messages) to every follower. A follower that hears nothing for its election timeout becomes a candidate: it increments its term, votes for itself and asks the others for votes. Timeouts are randomised so that one follower usually times out well before the rest and wins cleanly; the Raft paper used a range of 150 to 300 ms, and etcd defaults to a 100 ms heartbeat and a 1000 ms election timeout, randomised upward from there.
A node grants at most one vote per term, and only to a candidate whose log is at least as up to date as its own, compared first by the term of the last entry and then by its index. That rule is what makes elections safe: any entry committed on a majority is present in at least one voter of every future majority, so a candidate missing it cannot win. A candidate that collects votes from a strict majority becomes leader. With five voters that is three, so the cluster survives two failures; with four voters the majority is still three, which is why even-sized clusters add cost without adding fault tolerance.
# Follower side of RequestVote (Raft), with persistent state saved before replying.
def on_request_vote(req):
if req.term < current_term:
return Vote(term=current_term, granted=False)
if req.term > current_term:
current_term, voted_for = req.term, None # step down to follower
become_follower()
my_last_term, my_last_index = log.last_term(), log.last_index()
log_ok = (req.last_log_term > my_last_term or
(req.last_log_term == my_last_term and req.last_log_index >= my_last_index))
if log_ok and voted_for in (None, req.candidate_id):
voted_for = req.candidate_id
persist(current_term, voted_for) # fsync before answering
reset_election_timer()
return Vote(term=current_term, granted=True)
return Vote(term=current_term, granted=False)
Keeping elections rare: PreVote and CheckQuorum
Correctness is not the same as stability. Plain Raft has two weaknesses that show up in production. First, a node cut off from the rest keeps timing out and incrementing its term. When the partition heals, its high term forces the healthy leader to step down, causing a pointless election. PreVote fixes this: before incrementing its term, a would-be candidate asks whether peers would vote for it, and peers that have heard from a live leader recently say no. etcd exposes this as --pre-vote.
Second, a leader isolated from the majority can keep believing it leads, because nothing forces it to stop. CheckQuorum makes the leader step down if it has not heard from a majority within an election timeout. Together they turn most flaps into non-events. A third feature, lease-based reads, lets a leader serve reads locally while it believes its lease is valid, but that depends on bounded clock drift. The ReadIndex approach, which confirms leadership with one heartbeat round before serving a read, is safer and is etcd's default for linearizable reads.
Electing your own service on etcd
For your own service, etcd's Go client ships an election recipe in its concurrency package. Each replica opens a session, which is a lease kept alive by background heartbeats, and campaigns under a shared prefix. Campaigning writes a key under the prefix attached to the session lease and then waits until every key with a lower create revision is gone. The replica whose key is oldest is leader. If its process dies, keepalives stop, the lease expires after its TTL, the key is deleted, and the next waiter becomes leader without any extra round of voting.
sess, err := concurrency.NewSession(cli, concurrency.WithTTL(10)) // lease, kept alive in background
if err != nil { return err }
defer sess.Close()
e := concurrency.NewElection(sess, "/svc/billing/leader/")
if err := e.Campaign(ctx, hostname); err != nil { return err } // blocks until we lead
token := e.Rev() // create revision of our key: grows with every new leader
leaderCtx, cancel := context.WithCancel(ctx)
go func() { <-sess.Done(); cancel() }() // lease lost: stop leading immediately
// Writes into etcd itself can be fenced by a transaction on our key's revision.
resp, err := cli.Txn(leaderCtx).
If(clientv3.Compare(clientv3.CreateRevision(e.Key()), "=", e.Rev())).
Then(clientv3.OpPut("/svc/billing/config", newConfig)).
Commit()
if err == nil && !resp.Succeeded { return errLostLeadership } // compare failed: stop now
runLeaderLoop(leaderCtx, token) // pass the token to every external write
_ = e.Resign(context.Background()) // on graceful shutdown: hand over at once
Fencing beyond the coordination store
Look closely at what that transaction protects. The Compare on the key's create revision guarantees that a write into etcd only lands if this replica's key is still the one that won. It does nothing for a database, an object store or a payment API that the leader also writes to. A leader that pauses for 20 seconds in garbage collection or on a stalled disk can wake up after its lease expired and another replica took over, and its next write to the database will succeed unless the database checks.
The fix is a fencing token: a number that only grows across leaders, sent with every write and checked by the resource. The etcd create revision is a natural token because it is assigned by a linearizable store. The resource remembers the highest token it has accepted and rejects anything lower. If the resource cannot check tokens, the honest conclusion is that leader election alone cannot make that resource safe, and you need idempotent operations or a resource-side lock instead.
-- Resource side: accept a leader's write only if its token is not older than the newest seen.
UPDATE job_state
SET payload = :payload, fence = :token
WHERE job_id = :job_id
AND fence <= :token;
-- 0 rows updated means a newer leader exists: the caller must stop leading.
Kubernetes Lease election
Inside Kubernetes the usual mechanism is a Lease object in the coordination.k8s.io API group, used by client-go's leaderelection package and by the control plane itself. The object records holderIdentity, leaseDurationSeconds, acquireTime, renewTime and leaseTransitions. The holder renews it periodically; others poll it. Candidates do not trust the timestamps in the object, because clocks differ: each one records when it last saw the record change on its own clock, and only tries to take over once a full lease duration has passed since then.
| Parameter | kube-controller-manager default | What it controls |
|---|---|---|
| Lease duration | 15 s | How long others wait after the last observed renewal before taking over |
| Renew deadline | 10 s | How long the holder keeps retrying renewal before it gives up leading |
| Retry period | 2 s | How often candidates poll and the holder attempts renewal |
The 5-second gap between renew deadline and lease duration is the safety margin: the holder should have stopped work before anyone else may start. It is a margin, not a guarantee. A process frozen for longer than that gap can resume and act as leader for a moment, which is again why fencing matters.
Worked example: a failover time budget
Suppose a billing reconciler runs as three replicas that elect a leader through etcd with a 10-second session TTL, and etcd itself is a five-node cluster spread across three availability zones with default timings. We want to know how long reconciliation stops in each kind of failure.
| Failure | What happens | Pause in leadership |
|---|---|---|
| Leader replica shut down for a deploy | It calls Resign; the next waiter's Campaign returns | Milliseconds |
| Leader replica crashes or its node dies | Keepalives stop; lease expires; key deleted; next waiter leads | Up to the TTL, about 10 s |
| Leader replica loses network to etcd | Keepalives fail; session Done fires; it stops; lease expires | About 10 s, and the old leader stops on its own |
| etcd leader node dies | Followers time out after 1-2 s, PreVote and vote, new etcd leader | No change in app leader; etcd requests stall about 1-2 s |
| Whole AZ with 2 etcd voters lost | 3 of 5 remain, a majority; one election | Same as above |
Two lessons fall out. The TTL dominates unplanned failover, so pick it from the outage you can afford, not from habit; shortening it below a few seconds makes leadership flap on every slow keepalive. And graceful handover is nearly free, so make every shutdown path call Resign. In this example the only unsafe case is the third row: a replica whose network to etcd is slow but whose network to the database is fine can still write for a short time after losing the lease. The fencing check from the previous section is what turns that window into a rejected write instead of a double-applied invoice.
Failure modes
| Failure mode | Symptom | Mitigation |
|---|---|---|
| Election storms | Leader changes every few seconds; latency spikes | Heartbeat near network RTT, election timeout about 10x heartbeat, PreVote on |
| Disk stalls on the store | Missed heartbeats and spurious elections under load | Fast dedicated disks for the write-ahead log; alert on fsync latency |
| Stale leader after a pause | Writes from two leaders, duplicate side effects | Fencing tokens checked by every protected resource |
| Even-sized or two-zone cluster | Losing one zone loses quorum | Odd voter count spread over at least three failure domains |
| Leader work outlives lease | Long task continues after leadership is lost | Cancel work from the session Done channel; check the token per step |
| Hot leader | One replica saturated while others idle | Partition the work and elect a leader per shard instead of one global leader |
Operating elections
Measure leadership directly. Export the current leader identity, the count of leader transitions, and time since last successful renewal for every elected component; alert on transitions per hour, not on individual elections. For etcd, watch etcd_server_leader_changes_seen_total, WAL fsync duration and backend commit duration, since slow disks are the most common cause of unexplained elections.
Rehearse failover. Kill the leader replica, partition it from the store but not from its database, freeze it with SIGSTOP for longer than the TTL, and confirm in each case that exactly one replica does work afterwards and that stale writes are rejected. A failover you have only reasoned about is a hypothesis.
Trade-offs
A single leader keeps the design simple: one place to order operations, no conflict resolution. The price is a throughput ceiling and a failover pause. Shorter timeouts shrink the pause and raise the risk of false elections; longer ones do the reverse. Leasing from an existing store such as etcd or the Kubernetes API saves you from running consensus yourself but makes the store a dependency of every leader. When one leader becomes a bottleneck, elect one per shard so failures and load are spread out.
What to do next
- List every component that relies on a single leader and write down its acceptable failover pause.
- Set the session TTL or lease duration from that pause, and call Resign on every shutdown path.
- Give each protected resource a fencing check on a monotonic token, or document why it is idempotent.
- Run the coordination store with an odd number of voters across three failure domains, PreVote enabled.
- Alert on leader transitions per hour and on store fsync latency.
- Run a pause test with SIGSTOP longer than the TTL and confirm the stale leader's writes are rejected.
Keep learning: leader election algorithms compared, fencing tokens, Raft log replication, etcd and split brain.