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: consensus inside the store, leases for your service on topCoordination store (etcd / Raft), 5 votersNode 1leader, term 7Node 2followerNode 3followerNode 4followerNode 5followerheartbeats every 100 ms; election timeout 1 s, randomisedPreVote + CheckQuorum keep elections rareYour service: N replicas, one activeReplica Aholds key, rev 812Replica BwaitingReplica Cwaitingsession lease, TTL 10 scampaignProtected resourcerejects older tokenswrite + token 812Metricstransitions, renew lagElection prefix /svc/leader/lowest create revision winsThe store elects its own leader with Raft; your replicas elect theirs by racing for a leased key in the store.Only the fencing token carried to the resource makes a stale leader harmless.
Leader election in two layers: the store elects its own leader with Raft, replicas of your service compete for a leased key, and fencing protects resources.

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.

Parameterkube-controller-manager defaultWhat it controls
Lease duration15 sHow long others wait after the last observed renewal before taking over
Renew deadline10 sHow long the holder keeps retrying renewal before it gives up leading
Retry period2 sHow 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.

FailureWhat happensPause in leadership
Leader replica shut down for a deployIt calls Resign; the next waiter's Campaign returnsMilliseconds
Leader replica crashes or its node diesKeepalives stop; lease expires; key deleted; next waiter leadsUp to the TTL, about 10 s
Leader replica loses network to etcdKeepalives fail; session Done fires; it stops; lease expiresAbout 10 s, and the old leader stops on its own
etcd leader node diesFollowers time out after 1-2 s, PreVote and vote, new etcd leaderNo change in app leader; etcd requests stall about 1-2 s
Whole AZ with 2 etcd voters lost3 of 5 remain, a majority; one electionSame 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 modeSymptomMitigation
Election stormsLeader changes every few seconds; latency spikesHeartbeat near network RTT, election timeout about 10x heartbeat, PreVote on
Disk stalls on the storeMissed heartbeats and spurious elections under loadFast dedicated disks for the write-ahead log; alert on fsync latency
Stale leader after a pauseWrites from two leaders, duplicate side effectsFencing tokens checked by every protected resource
Even-sized or two-zone clusterLosing one zone loses quorumOdd voter count spread over at least three failure domains
Leader work outlives leaseLong task continues after leadership is lostCancel work from the session Done channel; check the token per step
Hot leaderOne replica saturated while others idlePartition 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

  1. List every component that relies on a single leader and write down its acceptable failover pause.
  2. Set the session TTL or lease duration from that pause, and call Resign on every shutdown path.
  3. Give each protected resource a fencing check on a monotonic token, or document why it is idempotent.
  4. Run the coordination store with an odd number of voters across three failure domains, PreVote enabled.
  5. Alert on leader transitions per hour and on store fsync latency.
  6. 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.

Key takeaway: Leader election is two problems: choosing a leader, which consensus solves, and stopping the old one, which no timeout can guarantee. Use a store such as etcd to choose, size the lease from the pause you can afford, hand over explicitly on shutdown, and carry a monotonic fencing token to every resource so a leader that wakes up late is refused rather than obeyed.