Many distributed systems are easier to build if one node is in charge: it orders writes, assigns work, runs the scheduler or owns a partition. Leader election is the protocol by which a group of nodes agrees on who that node is, and agrees again when it fails. Textbooks present several algorithms with very different assumptions, and production systems mostly use one or two of them. Knowing why is the difference between a design that survives a network partition and one that ends with two leaders writing to the same data.

This article compares the main families: the bully algorithm, ring algorithms, quorum-and-term election as used inside consensus protocols, and election built on a coordination service with leases. For a step-by-step trace of a Raft election with PreVote and timeouts, read distributed leader election architecture; this page is about choosing between algorithms and implementing the one you choose correctly.

Advertisement

What an election must guarantee

Two properties matter, and they pull in opposite directions.

  • Safety: at most one node believes it is the leader for any given epoch, and every node that learns of a leader for that epoch learns the same one.
  • Liveness: if the leader fails, eventually a new one is chosen while enough nodes are up and can communicate.

In a fully asynchronous network, where a slow node cannot be distinguished from a dead one, electing a unique leader with guaranteed termination is as hard as consensus, and the FLP result says deterministic consensus cannot guarantee termination if even one process may crash. Practical algorithms keep safety unconditional with majorities and increasing epoch numbers, and get liveness from timeouts. Algorithms that rely on timeouts for safety, like the classic bully algorithm, are correct only under synchrony that real networks violate.

Also, a node's belief that it is leader is always slightly stale: it can be paused by garbage collection after checking, while a new leader is elected. The resources a leader touches must be able to reject a deposed leader. That is fencing, covered below.

Four families of leader election, by the assumptions they needBullyhighest live id winsRing (LCR, HS)ids circulate, max winsQuorum + termsRaft, Zab, PaxosLease on a storeetcd, ZooKeeper, k8sneeds: sync timing,perfect failure detectorneeds: known ring,no failures during runneeds: majority alive,timeouts for liveness onlyneeds: a consensus store,bounded clock driftUnsafe under partitionstwo sides can each elect a leaderAt most one leader per term / leaseminority side cannot winStill required: fencing tokens at the resourcea paused old leader can act after its term or lease endedElection decides who should lead; fencing stops anyone else from acting as if they still do.
The classic algorithms assume reliable timing or failure-free runs and can elect two leaders under a partition. Quorum and lease designs bound leadership by epoch or lease, and fencing at the resource covers the remaining window.

The bully algorithm

Garcia-Molina's bully algorithm (1982) assumes every node has a unique, totally ordered id and knows all the others, and that message delays are bounded so a timeout means a crash. The highest-id live node wins.

on detect_leader_failure() or on startup:
    start_election()

start_election():
    higher = [n for n in nodes if n.id > my.id]
    if not higher:
        become_leader(); broadcast COORDINATOR(my.id); return
    send ELECTION to every node in higher
    if no OK arrives within T:            # nobody above me is alive
        become_leader(); broadcast COORDINATOR(my.id)
    else:
        wait for COORDINATOR within T2, else start_election()

on ELECTION from lower node:
    reply OK; start_election()            # I bully it out of the running

on COORDINATOR(id):
    leader = id

Message cost is O(n2) in the worst case, when the lowest-id node detects the failure, and O(n) in the best case, when the next-highest node does. The safety story matters more. If the network splits, each side elects its own highest node and both act as leader. If the old leader was merely slow, it returns and bullies its way back, causing a second change. Fine for teaching or a tightly coupled cluster with reliable failure detection; wrong for a network you do not control.

Advertisement

Ring algorithms and their message complexity

Ring algorithms assume nodes are arranged in a logical ring and can only message their neighbours. They were studied mainly to understand the cost of agreement, and the results are still useful for reasoning about protocols that pass tokens around.

  • LCR / Chang-Roberts. Le Lann (1977) proposed circulating every id around a unidirectional ring; Chang and Roberts (1979) improved it so a node forwards an incoming id only if it is larger than its own and discards smaller ones. When a node receives its own id back, it is the maximum and becomes leader, then circulates an announcement. Worst case O(n2) messages, when ids are arranged in decreasing order along the direction of travel; average O(n log n) over random arrangements.
  • Hirschberg-Sinclair. On a bidirectional ring, each candidate probes neighbourhoods of doubling radius 1, 2, 4, ... in both directions and survives a phase only if it is the largest id in that neighbourhood. Worst case O(n log n) messages.
  • Lower bound. Comparison-based election on a ring needs Omega(n log n) messages in the worst case, so Hirschberg-Sinclair and later unidirectional algorithms (Peterson; Dolev, Klawe and Rodeh) are optimal up to constants.
# Chang-Roberts, one process. Ids are unique integers; send() goes clockwise.
def on_start():
    global participant
    participant = True
    send(("ELECT", my_id))

def on_message(kind, value):
    global participant, leader
    if kind == "ELECT":
        if value > my_id:
            participant = True; send(("ELECT", value))
        elif value < my_id and not participant:
            participant = True; send(("ELECT", my_id))
        elif value == my_id:
            leader = my_id; send(("LEADER", my_id))   # my id went all the way round
        # value < my_id and already participating: swallow it
    elif kind == "LEADER":
        leader = value
        if value != my_id:
            send(("LEADER", value))

Ring algorithms assume a stable ring and no failures during the run; a crash mid-election stalls it until a timeout restarts it, and there is no partition safety. Their lasting lesson: electing by maximum id is cheap when nothing fails, and every extra guarantee costs messages, rounds or assumptions.

Quorum-and-term election

Consensus protocols such as Raft, Zab (ZooKeeper's atomic broadcast) and Multi-Paxos embed an election that is safe under any partition. The ingredients are always the same:

  1. Epochs. Every election attempt uses a new, larger epoch number (a term in Raft, an epoch in Zab, a ballot in Paxos), persisted to disk before voting.
  2. One vote per epoch. Each node grants at most one vote per epoch, also persisted, so a node restart cannot produce a second vote.
  3. Majority quorum. A candidate needs votes from a majority of the configured members. Any two majorities intersect, so two candidates cannot both win the same epoch. The side of a partition with a minority cannot elect anyone.
  4. Log check. A voter refuses candidates whose log is less up to date than its own, so the winner holds every committed entry. This is what makes election part of replication rather than a separate service.
  5. Randomised timeouts for liveness. Raft's paper suggests election timeouts drawn from a range such as 150-300 ms so that usually one node times out first and wins before others start, avoiding repeated split votes.

The cost is that a cluster of 2f+1 nodes tolerates only f failures, and that each election needs at least one round trip to a majority. The benefit is that safety never depends on timing. If you already run a consensus-replicated system, use its built-in leadership rather than building another layer. The details of Raft's version are in Raft consensus, in depth, and the Paxos ballot equivalent in Paxos.

Election on a coordination service

Application singletons such as schedulers and Kubernetes controllers borrow consensus instead of running it: a consensus-replicated store (etcd, ZooKeeper, Consul, or the Kubernetes API) holds a leadership key tied to a lease that expires if not renewed. The store's linearizable compare-and-set gives safety; the lease gives liveness when the holder dies.

ZooKeeper recipe. Each candidate creates an ephemeral sequential znode under an election path, such as /election/n_0000000042. The candidate with the lowest sequence number is leader. Every other candidate watches only the znode immediately before its own, not the leader's, so a leader failure wakes exactly one node instead of the whole herd. Ephemeral nodes disappear when the session expires. The recipe and session semantics are covered in ZooKeeper architecture.

etcd. The Go client ships the recipe in the concurrency package. A session owns a lease kept alive in the background; Campaign blocks until this candidate holds the leader key.

import (
    "context"
    clientv3 "go.etcd.io/etcd/client/v3"
    "go.etcd.io/etcd/client/v3/concurrency"
)

func runForLeader(ctx context.Context, cli *clientv3.Client, id string) error {
    sess, err := concurrency.NewSession(cli, concurrency.WithTTL(10)) // lease, seconds
    if err != nil { return err }
    defer sess.Close()

    e := concurrency.NewElection(sess, "/services/scheduler/leader")
    if err := e.Campaign(ctx, id); err != nil { return err } // blocks until elected

    token := e.Rev() // revision of our leader key: monotonic, use it as a fencing token
    workCtx, cancel := context.WithCancel(ctx)
    go func() { <-sess.Done(); cancel() }() // lease lost: stop working immediately
    err = doLeaderWork(workCtx, token)
    _ = e.Resign(context.Background())
    return err
}

Kubernetes. client-go's leaderelection package implements the same idea on a coordination.k8s.io Lease object, configured with LeaseDuration, RenewDeadline and RetryPeriod and callbacks OnStartedLeading and OnStoppedLeading. The control-plane components default to 15 s, 10 s and 2 s. Its package documentation states plainly that it does not guarantee that only one client is acting as leader, because it relies on timestamps and clock rates rather than fencing. When OnStoppedLeading fires, the conventional and safest response is to exit the process.

Comparing the families

FamilyMessages per electionSafe under partition?Depends on timing forTypical use
BullyO(n) best, O(n2) worstNoSafety and livenessTeaching, single-host clusters
Ring (Chang-Roberts)O(n log n) average, O(n2) worstNoRestarting stalled runsToken rings, theory
Ring (Hirschberg-Sinclair)O(n log n) worstNoRestarting stalled runsTheory
Quorum + termsOne round to a majority, more on split votesYes, majority side onlyLiveness onlyInside Raft, Zab, Paxos systems
Lease on a storeA few store round trips plus renewalsYes, via the storeLease expiry (bounded drift)Application singletons, controllers

Fencing: the part that makes any of them safe

A correctly elected leader can still act after being replaced. A leader with a 10-second lease enters a 15-second GC pause just after checking its status; the lease expires, another node is elected and writes; the first node wakes and writes too. No election algorithm prevents this, because it happens after the election.

The fix is a fencing token: a number that increases with every leadership change (the Raft term, the ZooKeeper znode sequence or zxid, the etcd revision of the leader key) and that the leader sends with every request to the shared resource. The resource remembers the highest token it has seen and rejects lower ones. The pattern is covered in fencing tokens. If the resource cannot check tokens, such as a third-party API, shrink the window instead: stop work as soon as the lease is in doubt, keep leases well above the worst observed pause, and make operations idempotent so a duplicate does no harm.

Worked example: a singleton job scheduler

Three replicas of a scheduler run in Kubernetes; exactly one should enqueue jobs into a Postgres table every second. A previous design used a database advisory lock and occasionally ran two schedulers after a network blip, producing duplicate jobs.

  1. Election. Use client-go leaderelection on a Lease with the 15/10/2 defaults. Failover takes up to about 15 seconds after a leader dies, which the business accepted.
  2. Token. Store the Lease's leaseTransitions counter, which increments on each change of holder, as the fencing token. On becoming leader, the scheduler reads it and writes it to a one-row scheduler_epoch table with UPDATE ... WHERE epoch < $token.
  3. Fenced writes. Every enqueue runs in a transaction that first checks SELECT epoch FROM scheduler_epoch FOR SHARE equals its token. A deposed leader's write fails because the new leader has already raised the epoch.
  4. Idempotency. Jobs carry a unique key of schedule id plus tick time, so even a slipped duplicate is rejected by a unique index.
  5. Exit on loss. OnStoppedLeading calls os.Exit(1); Kubernetes restarts the pod as a follower.

After the change, a chaos test that froze the leader process for 30 seconds produced no duplicate jobs: the frozen leader's first enqueue on waking failed the epoch check.

Failure modes

  • Flapping. Timeouts shorter than normal latency spikes or GC pauses cause repeated elections. Set timeouts from measured tail latency, not defaults, and use PreVote or check-quorum where available. Adaptive failure detectors help; see phi accrual failure detection.
  • Split brain from non-quorum election. Bully, ring, or "whoever holds the database row" designs elect two leaders on a partition. Use a majority-based store.
  • Clock assumptions. Lease safety assumes clocks advance at similar rates. A VM paused and resumed can treat an expired lease as still valid. Fencing is the only robust protection.
  • Herd effect. Every candidate watching the leader key causes a burst of requests on failover. Watch the predecessor, or rely on the client library's backoff.
  • Silent leader. The lease is renewed but the work loop is stuck. Tie renewal to the work loop's health.

What to do next

  1. List every component that assumes it is the only active instance; each needs an explicit election.
  2. If the component is already part of a consensus group, use the group's leadership; otherwise use a lease on etcd, ZooKeeper or the Kubernetes API, not a home-grown algorithm.
  3. Pick a fencing token from the election mechanism and make every shared resource reject stale tokens.
  4. Set lease and timeout values from measured pause and latency percentiles, and document the failover time they imply.
  5. On loss of leadership, stop work immediately; exiting the process is simplest.
  6. Run a chaos test that freezes the leader longer than the lease (for example with kill -STOP) and confirm no duplicate side effects.
  7. Alert on election rate; more than a handful per day usually means timeouts are too tight.
Key takeaway: Leader election algorithms differ mainly in what they assume. Bully and ring algorithms are cheap but rely on timing or failure-free runs and can elect two leaders under a partition. Quorum-and-term election, used inside Raft, Zab and Paxos, keeps safety unconditional and uses timeouts only for liveness; lease-based election on etcd, ZooKeeper or Kubernetes borrows that safety for application services. None of them stops a paused former leader from acting, so pair election with fencing tokens checked at every shared resource.