Some work must run on exactly one node at a time: a scheduler, a partition owner or a coordinator. Leader election picks that node, and the hard part is not choosing a leader but guaranteeing that a slow or paused old leader cannot keep acting after a new one is chosen.
Before you start
You should understand replication, timeouts and basic consensus ideas. This article is conceptual.
Step-by-step walkthrough
Step 1: Choose a mechanism
A strongly consistent store using a lease, or a consensus protocol such as Raft, provides election. A lease is a key with a TTL that the leader renews; if it stops renewing, the lease expires and another node can claim it. Consensus gives stronger guarantees at the cost of a quorum.
Step 2: Add fencing tokens
A paused leader may believe it still leads after its lease expired. Every write it makes should carry a monotonically increasing fencing token, and the resource rejects writes with a token older than the last seen. This turns “am I still leader” into an enforceable check at the storage layer.
Step 3: Design for the pause, not the crash
A crashed leader is easy; a paused one (a long GC pause or a stalled network) is dangerous because it resumes believing it leads. Fencing handles this by making stale writes fail. Without it, two leaders can write, causing split brain and data corruption.
Worked scenario
The resource rejects an old leader’s write using the token.
lease token 5: leader A wins, writes with token 5
A pauses; lease expires
leader B wins with token 6, writes with token 6
A resumes, writes with token 5 -> storage rejects 5 < 6, A steps downWalk through the example
B holds the higher token, so the storage refuses any write from A carrying the stale token. A discovers it is no longer the leader and stops. The fencing token converts an unreliable local belief into a check the shared resource can enforce.
Common mistake
Relying on a lease alone without fencing, so a paused leader can still write after expiry. Another is trusting wall clocks across nodes, since clock skew makes lease durations unreliable; use the store’s own time or a consensus term instead.
Verify the behavior
Pause the leader artificially past its lease and confirm the new leader takes over. Have the paused leader attempt a write and confirm it is rejected by the token. Assert only one node’s writes are accepted at a time.
Interview exercise
Why is a paused leader more dangerous than a crashed one?
Answer and reasoning
A crashed leader stops acting, so the system notices and elects a replacement cleanly. A paused leader resumes with stale beliefs and may issue writes as if it were still in charge, creating a second writer and corrupting state. Fencing tokens and leases with short renewals are what make the pause safe.
Continue learning
Compare coordination in Service discovery and replication in Replication. Read the Raft consensus documentation and try the System design interview questions.