leader election
Leader Election
Understand why systems designate a single coordinator, why electing one is harder than it appears, and why believing you are the leader is never the same as actually being the leader.
Prerequisites
Recommended Knowledge
- Consensus and quorum intersection
- Raft terms and voting restrictions
- Partial failure and timeout ambiguity
- Network partitions and split-brain
- Failure detection limits
- Replication and log ordering
- Clock skew and synchronization error
- Health checks and their weaknesses
Why Designate a Leader
Many coordination problems become dramatically simpler when a single node is responsible for a decision. Rather than negotiating every action, participants defer to one coordinator.
| Purpose | Without a Leader | With a Leader |
|---|---|---|
| Write ordering | Agreement needed per write | Leader assigns sequence |
| Scheduled jobs | Every node may run it | Only the leader executes |
| Shard assignment | Conflicting maps possible | One authoritative assignment |
| Cluster membership | Divergent views | Single source of truth |
| External integration | Duplicate calls to a partner | One caller only |
| Cache warming | Redundant work across nodes | Coordinated single pass |
Simple Analogy
A committee appoints a chair to avoid debating every decision. The difficulty is not appointing one; it is ensuring the previous chair accepts they no longer hold the role when they return from an absence.
The Two Requirements
| Property | Statement | Violation Consequence |
|---|---|---|
| Safety | At most one leader per term | Split-brain and data corruption |
| Liveness | A leader is eventually elected | System stalls indefinitely |
Two Leaders Permanent Corruption
Election Mechanisms
| Mechanism | How It Works | Safety Guarantee |
|---|---|---|
| Quorum voting | Candidate requires a majority | Strong, intersection prevents two |
| Coordination service lock | First to acquire an ephemeral node wins | Strong, delegated to the service |
| Database row lock | Conditional update on a leader row | Strong if the database is consistent |
| Lease renewal | Hold a time-bounded claim | Depends on clock assumptions |
| Bully algorithm | Highest identifier wins | Weak under partitions |
| Static assignment | Configuration names the leader | No automatic failover |
Lock-Based Election
The most common practical approach delegates the hard part to a consensus-backed coordination service. The application competes for a lock and treats acquisition as leadership.
class LeaderElection {
constructor(coordinator, config) {
this.coordinator = coordinator;
this.lockKey = config.lockKey;
this.nodeId = config.nodeId;
this.leaseSeconds = config.leaseSeconds;
this.isLeader = false;
this.fencingToken = null;
}
async campaign() {
const result = await this.coordinator.acquire({
key: this.lockKey,
holder: this.nodeId,
ttlSeconds: this.leaseSeconds
});
if (!result.acquired) {
this.isLeader = false;
return { leader: false, currentHolder: result.holder };
}
this.isLeader = true;
this.fencingToken = result.fencingToken;
this.startRenewal();
return { leader: true, token: this.fencingToken };
}
startRenewal() {
const intervalMs = (this.leaseSeconds * 1000) / 3;
this.renewalTimer = setInterval(async () => {
const renewed = await this.coordinator.renew({
key: this.lockKey,
holder: this.nodeId,
ttlSeconds: this.leaseSeconds
});
if (!renewed.ok) {
this.onLeadershipLost(renewed.reason);
}
}, intervalMs);
}
onLeadershipLost(reason) {
clearInterval(this.renewalTimer);
this.isLeader = false;
this.fencingToken = null;
this.stopLeaderWork(reason);
}
}
-- Conditional acquisition, atomic by construction
UPDATE cluster_leadership
SET holder_id = :node_id,
fencing_token = fencing_token + 1,
acquired_at = CURRENT_TIMESTAMP,
expires_at = CURRENT_TIMESTAMP + INTERVAL '15 seconds'
WHERE resource = :resource
AND (expires_at < CURRENT_TIMESTAMP OR holder_id = :node_id)
RETURNING fencing_token, expires_at;
-- Zero rows returned means another node holds valid leadership
The Stale Leader Problem
This is the defining hazard of leader election. A leader that loses contact with the cluster continues believing it leads, because nothing has told it otherwise.
| Cause | What the Old Leader Believes | Reality |
|---|---|---|
| Network partition | Peers are unreachable but I still lead | Majority elected a successor |
| Long garbage collection pause | No time has passed | Lease expired during the pause |
| Virtual machine suspension | Execution continues normally | Minutes elapsed, leadership moved |
| Disk or dependency stall | Operation is merely slow | Renewal deadline passed |
| Clock jump backwards | Lease still valid | Lease expired in real time |
Fencing Tokens
A fencing token is a monotonically increasing number issued on each leadership grant. The protected resource records the highest token it has seen and rejects anything lower.
-- Resource-side fencing: stale writers are rejected
UPDATE shard_assignment
SET assignment = :new_assignment,
fencing_token = :incoming_token,
updated_at = CURRENT_TIMESTAMP
WHERE shard_id = :shard_id
AND fencing_token <= :incoming_token;
-- Zero rows affected means a deposed leader attempted a write
class FencedResource {
constructor() {
this.highestToken = 0;
this.state = null;
}
write(value, token) {
if (token < this.highestToken) {
return {
accepted: false,
reason: "stale fencing token",
presented: token,
required: this.highestToken
};
}
this.highestToken = token;
this.state = value;
return { accepted: true, token };
}
}
| Resource Type | Fencing Implementation |
|---|---|
| Relational database | Token column with conditional update |
| Object storage | Conditional write on generation number |
| Message topic | Producer epoch rejecting older epochs |
| Internal service | Token in request header, validated server-side |
| External partner API | Idempotency key derived from the token |
Leases and Clock Assumptions
A lease grants leadership for a bounded period. Its safety rests on an assumption that deserves scrutiny: that clock drift stays within a known bound.
| Assumption | If Violated | Protection |
|---|---|---|
| Bounded clock drift | Holder believes an expired lease is valid | Monitor skew, use monotonic clocks |
| Bounded process pauses | Holder resumes after expiry | Fencing tokens |
| Grantor waits full lease | Two valid holders simultaneously | Add margin before regranting |
| Monotonic time source | Backward jump extends the lease | Use monotonic, not wall-clock, time |
class MonotonicLease {
constructor(leaseSeconds, safetyMarginMs) {
this.leaseMs = leaseSeconds * 1000;
this.marginMs = safetyMarginMs;
this.acquiredAtMonotonic = null;
}
acquire() {
this.acquiredAtMonotonic = performance.now();
}
isSafelyHeld() {
if (this.acquiredAtMonotonic === null) return false;
const elapsed = performance.now() - this.acquiredAtMonotonic;
return elapsed < (this.leaseMs - this.marginMs);
}
assertLeadership() {
if (!this.isSafelyHeld()) {
throw new LeadershipExpiredError(
"lease margin exhausted, refusing to act"
);
}
}
}
Failover Timing
The lease duration is a direct trade-off between how quickly failure is detected and how often healthy leaders are wrongly replaced.
Short Lease
- Fast detection of genuine failure
- Frequent renewal traffic
- Healthy leaders lost during latency spikes
- Election churn under load
Long Lease
- Stable under transient problems
- Extended outage after real failure
- Longer window for a stale leader to act
- Slow recovery from crashes
| Workload | Lease Preference | Reasoning |
|---|---|---|
| Interactive write path | Short | Users notice every second of outage |
| Batch coordination | Long | Stability matters more than speed |
| Cross-region cluster | Long | Latency variance causes false positives |
| Expensive warmup | Long | Failover cost exceeds detection benefit |
Becoming and Ceasing to Be Leader
Transitions are where most bugs live. Both directions require deliberate handling.
On Acquiring Leadership
Record the fencing token, load current state, verify nothing conflicts, and only then begin leader work. Starting immediately risks acting on stale assumptions.
Before Each Action
Verify the lease still has margin remaining. A long-running operation started while leading may complete after leadership has moved.
On Losing Leadership
Halt in-flight work, cancel timers, close exclusive resources, and stop issuing writes. Do not attempt to finish the current operation.
On Graceful Shutdown
Release the lease explicitly rather than waiting for expiry. This turns a lease-duration outage into a near-instant handover.
class LeaderWorker {
async runCycle() {
if (!this.election.isLeader) return;
this.election.lease.assertLeadership();
const work = await this.fetchPendingWork();
for (const item of work) {
this.election.lease.assertLeadership();
const result = await this.processItem(item, {
fencingToken: this.election.fencingToken
});
if (!result.accepted && result.reason === "stale fencing token") {
this.election.onLeadershipLost("fenced by resource");
return;
}
}
}
async shutdown() {
this.stopped = true;
await this.drainInFlight();
if (this.election.isLeader) {
await this.election.release();
}
}
}
Partition Behavior
| Side | Correct Behavior | Failure Mode If Wrong |
|---|---|---|
| Majority side | Elects a new leader and proceeds | Unnecessary outage if it refuses |
| Minority side | Steps down and stops acting | Split-brain if it continues |
| Isolated leader | Self-demotes on renewal failure | Two leaders writing concurrently |
| Rejoining node | Recognizes higher term and follows | Disrupts a stable cluster |
Scaling Beyond One Leader
A single leader is a throughput ceiling. Where the workload partitions naturally, electing a leader per partition preserves the coordination guarantee while restoring parallelism.
| Model | Scope | Trade-off |
|---|---|---|
| Single global leader | Entire system | Simple, but a hard throughput ceiling |
| Leader per shard | One data range | Scales, but more elections to manage |
| Leader per resource | One entity or job | Fine-grained, higher coordination overhead |
| Leaderless | None | No ceiling, requires conflict resolution |
Monitoring
Signals Worth Tracking
- Leadership change frequency
- Fencing token value and rate of increase
- Time with no established leader
- Lease renewal failures
- Renewal latency against the lease period
- Writes rejected by fencing
- Nodes simultaneously claiming leadership
- Leader warmup duration after election
- Process pause duration
- Clock skew across candidates
- Coordination service availability
Verification
describe("leader election safety", function () {
it("rejects writes from a deposed leader", async function () {
const first = await electLeader("node-a");
const firstToken = first.fencingToken;
await faultInjector.partition("node-a");
const second = await electLeader("node-b");
expect(second.fencingToken).toBeGreaterThan(firstToken);
const staleWrite = await resource.write("value-x", firstToken);
expect(staleWrite.accepted).toBe(false);
expect(staleWrite.reason).toBe("stale fencing token");
});
it("never permits two simultaneous leaders", async function () {
const observations = await runChaosScenario({
durationMs: 120_000,
partitions: true,
pauses: true,
restarts: true
});
for (const moment of observations) {
expect(moment.activeLeaders.length).toBeLessThanOrEqual(1);
}
});
it("self-demotes when renewal fails", async function () {
const leader = await electLeader("node-a");
await faultInjector.blockCoordinator("node-a");
await sleep(leaseMs);
expect(leader.isLeader).toBe(false);
});
});
Conditions Worth Injecting
- Leader process killed abruptly
- Leader partitioned from the cluster
- Leader paused beyond the lease period
- Coordination service unavailable
- Clock adjusted backwards on the leader
- Renewal delayed past the safety margin
- Simultaneous campaigns from several nodes
- Old leader rejoining after a long absence
Common Design Mistakes
Weak Design
- Assuming the leader knows when it is deposed
- Omitting fencing tokens entirely
- Using wall-clock time for leases
- Checking leadership only at loop start
- Building election on a non-consensus store
- Setting leases near network latency
- Renewing at the last possible moment
- Not releasing leases on shutdown
- Never testing partition and pause scenarios
Strong Design
- Enforces exclusivity at the resource
- Issues and validates fencing tokens
- Uses monotonic clocks with a safety margin
- Re-verifies leadership per operation
- Delegates election to a proven service
- Sets leases well above latency variance
- Renews at a fraction of the lease
- Releases explicitly on shutdown
- Drills failover and pause regularly
System Design Interview Discussion
| Question | What Your Answer Should Cover |
|---|---|
| How is a leader elected? | Quorum or consensus-backed lock |
| What if the old leader returns? | Fencing tokens rejecting stale writes |
| What about a long pause? | Lease expiry and resource-side enforcement |
| How long is the lease? | Detection versus churn trade-off |
| Why monotonic clocks? | Backward jumps extending expired leases |
| What happens during a partition? | Majority proceeds, minority steps down |
| Does one leader scale? | Per-shard leadership for parallelism |
| How is safety verified? | Chaos testing for concurrent leaders |
Design Checklist
Production Checklist
- Use a consensus-backed coordination service
- Issue a monotonic fencing token per grant
- Validate tokens at every protected resource
- Use monotonic clocks for lease tracking
- Reserve a safety margin before acting
- Renew at a fraction of the lease period
- Self-demote on renewal failure
- Re-verify leadership inside long loops
- Halt in-flight work on leadership loss
- Release the lease on graceful shutdown
- Set leases above observed latency variance
- Scope leadership to the narrowest unit
- Use idempotency keys for unfenceable externals
- Monitor election frequency and token progression
- Alert on fencing rejections
- Drill partitions, pauses, and abrupt kills
Knowledge Check
Why is safety prioritized over liveness?
No leader causes a recoverable outage, while two leaders can corrupt state permanently through conflicting concurrent writes.
What is the stale leader problem?
A deposed leader continues believing it leads, because partitions and pauses give it no signal that leadership moved.
How do fencing tokens solve it?
The resource records the highest token seen and rejects lower ones, so enforcement does not depend on the old leader detecting anything.
Why use monotonic clocks for leases?
Wall-clock time can jump backwards during synchronization, appearing to extend a lease that has already expired in real time.
Why elect per shard rather than globally?
A single leader caps throughput. Where work partitions cleanly, per-shard leadership preserves exclusivity while restoring parallelism.
Summary
Leader election designates one node to coordinate decisions that must not be made concurrently. It has two requirements: at most one leader at a time, and eventually some leader. Safety dominates, because an outage is recoverable while split-brain corruption often is not.
Election itself is usually delegated to a consensus-backed coordination service. The genuinely hard problem is the stale leader: a node partitioned or paused continues believing it leads, and no amount of self-checking can make it certain at the moment its write reaches a resource.
Fencing tokens are the robust answer. A monotonic token issued per leadership grant, validated at the resource, rejects writes from any deposed leader without requiring that leader to detect anything. Leases add time bounds but depend on clock assumptions that pauses and backward jumps can violate.
Lease duration trades detection speed against election churn, and leadership should be scoped to the narrowest unit requiring exclusivity so a single coordinator does not become a throughput ceiling for work that partitions cleanly.
Key Takeaway
Never trust a leader's own belief that it leads. Issue fencing tokens on every grant and validate them at the resource, track leases with monotonic clocks and a safety margin, re-verify leadership inside long operations, and test with partitions and process pauses rather than assuming failover behaves as designed.