Table of Contents

    leader election

    SYSTEM DESIGN • CHAPTER 15.2

    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.

    Learning objective: By the end of this article, you will understand why leaders are used, election mechanisms and their guarantees, the stale leader problem, leases and their clock dependency, failover timing trade-offs, and how to protect resources from a deposed 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.

    Electing a leader is straightforward. Ensuring the old leader knows it has been replaced is the actual problem, and it cannot be solved by asking them.

    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
    THE ASYMMETRY
    No Leader Temporary Outage

    Two Leaders Permanent Corruption
    PRIORITY RULE
    Safety dominates. Having no leader for thirty seconds is recoverable. Having two leaders for one second may corrupt state permanently.

    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
    Why the Bully Algorithm Fails Under a partition, the highest-identifier node on each side declares itself leader. Both are following the rules correctly, and both are wrong.

    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);
        }
    }
    Renew well before expiry: Renewing at one third of the lease period tolerates two consecutive failed attempts before the lease lapses. Renewing at the last moment guarantees loss on any transient hiccup.
    -- 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.

    THE DANGEROUS WINDOW
    \[ W_{\text{overlap}} = T_{\text{new leader elected}} - T_{\text{old leader stops acting}} \]
    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
    The Pause Scenario A leader begins writing, pauses for twenty seconds, and resumes. Its lease expired, a new leader was elected and made changes, and the old leader's write lands afterwards, silently overwriting them.
    THE CORE INSIGHT
    A node can never be certain it is still the leader at the moment its write reaches the resource. The resource must reject stale writers itself.

    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.

    FENCING RULE
    \[ \text{Accept write} \iff token_{\text{incoming}} \geq token_{\text{highest seen}} \]
    -- 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 };
        }
    }
    Why This Is the Only Robust Answer Fencing does not require the old leader to detect anything. It shifts enforcement to the resource, which always knows whether a newer leader has already written.
    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
    When Fencing Is Impossible Some external systems accept any authenticated caller with no version check. For these, reduce exposure through short leases, idempotency keys, and a conservative pause before acting after acquiring leadership.

    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.

    SAFE LEASE CONDITION
    \[ T_{\text{grantor waits}} > T_{\text{lease}} + \epsilon_{\text{clock}} + T_{\text{max pause}} \]
    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"
                );
            }
        }
    }
    Use monotonic time for leases: Wall-clock time can jump backwards during synchronization, appearing to extend a lease that has actually expired. Monotonic clocks never move backwards.

    Failover Timing

    The lease duration is a direct trade-off between how quickly failure is detected and how often healthy leaders are wrongly replaced.

    DETECTION TRADE-OFF
    \[ T_{\text{failover}} \approx T_{\text{lease}} + T_{\text{election}} + T_{\text{warmup}} \]

    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
    Election Churn Leases too short relative to network variance cause repeated elections. Each one halts coordinated work, so the system destabilizes precisely when it is under the most load.

    Becoming and Ceasing to Be Leader

    Transitions are where most bugs live. Both directions require deliberate handling.

    1

    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.

    2

    Before Each Action

    Verify the lease still has margin remaining. A long-running operation started while leading may complete after leadership has moved.

    3

    On Losing Leadership

    Halt in-flight work, cancel timers, close exclusive resources, and stop issuing writes. Do not attempt to finish the current operation.

    4

    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();
            }
        }
    }
    Check leadership inside loops, not only before them: A batch that takes two minutes may outlive a fifteen-second lease. Re-verifying per item bounds how long a deposed leader continues acting.

    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
    Self-Demotion on Renewal Failure A leader unable to renew its lease should stop acting before the lease expires, not after. This shrinks the overlap window rather than relying solely on fencing.

    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
    GRANULARITY RULE
    Scope leadership to the narrowest unit requiring exclusivity. A global leader for work that partitions cleanly wastes capacity for no correctness benefit.

    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
    A fencing rejection is not an error to suppress. It is proof that a stale leader attempted a write and the resource stopped it, which is the system working exactly as designed.

    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

    1

    Why is safety prioritized over liveness?

    No leader causes a recoverable outage, while two leaders can corrupt state permanently through conflicting concurrent writes.

    2

    What is the stale leader problem?

    A deposed leader continues believing it leads, because partitions and pauses give it no signal that leadership moved.

    3

    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.

    4

    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.

    5

    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.