PrepZone Logo
PrepZone

Leader Election and Failover

One writer at a time, clean failover, and fencing stale leaders — ShardPay's shard primaries.

Why this matters

  • Two leaders writing the same shard causes irreconcilable ledger state — worse than temporary unavailability.
  • Failover speed vs safety is the core trade-off; aggressive timeouts cause flapping, conservative timeouts extend outage.
  • ShardPay uses Raft-based election per ledger shard and for control-plane metadata.
  • Interviewers ask how you prevent split-brain and what happens to in-flight writes during failover.
FollowerElection timeout
CandidateRequests votes
LeaderAccepts client writes
Followers timeout, become candidates, request votes. Majority wins.

Safe failover

When ShardPay's shard leader dies, followers increment their election timer. After a randomized timeout (150–300ms), a follower campaigns for votes. A candidate needs a majority quorum to become leader, appends a new term, and rejects writes from any node with a stale term. The old leader, if it recovers from a network glitch, cannot commit — its term is lower and clients have moved on.

Election mechanics

  • Election timeout — Follower waits random interval before starting election; reduces split-vote probability when multiple nodes timeout together. ShardPay tunes base timeout to 200ms with ±50ms jitter on ledger shards in one region (RTT < 2ms).
  • Quorum (majority vote) — Candidate must receive votes from more than half of cluster members. ShardPay's 5-node shard requires 3 votes — a partitioned minority cannot elect a leader.
  • Term / epoch — Monotonic generation number; higher term wins. Every ShardPay Raft RPC carries term; stale-term responses trigger client redirect to current leader.
  • Log completeness check — Voters reject candidates whose log is behind. ShardPay ensures the new leader has all committed entries before serving writes — no lost transfers on failover.
  • Split-brain — Two nodes both believing they are leader. Prevented by quorum requirement + fencing: storage rejects writes without current-term lease.

Election lifecycle

Walkthrough: Leader crash during peak traffic

  1. Leader node-A on ledger-55 suffers JVM OOM at 14:32:01.
  2. Followers B, C, D, E stop receiving heartbeats; election timers start.
  3. C times out first at 14:32:01.220, increments term to 47, requests votes.
  4. C receives votes from D and E (3 of 5 = majority); becomes leader term 47.
  5. C sends heartbeats; B catches up and acknowledges new leader.
  6. ShardPay routing service watches metadata; clients receive updated leader endpoint by 14:32:03.
  7. Transfers that failed with NOT_LEADER retry with idempotency keys — no duplicates.
Java
// ShardPay follower — simplified election trigger
void onHeartbeatTimeout() {
    if (state == FOLLOWER) {
        state = CANDIDATE;
        currentTerm++;
        requestVotes();
    }
}

// Storage fence — reject stale leader writes
public WriteResult append(Entry entry, long claimedTerm) {
    if (claimedTerm < storage.currentTerm()) {
        return WriteResult.rejected("STALE_TERM");
    }
    return storage.append(entry);
}

Tuning and operations

Production concerns

  • Flapping — Leader repeatedly fails election or network oscillates. ShardPay requires 3 stable heartbeats before advertising new leader to external routers; adds 600ms to failover but prevents route churn.
  • Preferred leader / stickiness — Pin leadership to a node with fast disk. ShardPay sets raft.priority higher on nodes with NVMe — reduces log apply lag during normal operation.
  • Witness / observer nodes — Non-voting members for cross-AZ quorum without full data copy. ShardPay uses observers in DR region for metadata only — not in ledger write quorum.
  • Manual failover — Planned maintenance triggers graceful leadership transfer (transferLeadership) before patch. ShardPay's runbook: transfer → verify follower caught up → stop old leader → patch.
  • Fencing tokens to downstream — Leader passes monotonic fence to storage engine. ShardPay's reconciliation batch job checks fenceToken on object store writes so a zombie leader cannot overwrite S3 settlement files.

Example: Split-brain near-miss

A misconfigured network ACL isolated node-A (old leader) from B-E but clients could still reach A. Without quorum, A should step down — but a bug in an older build allowed solo writes for 90 seconds. Fix: hard gate hasWriteQuorum() before any append; deploy Jepsen regression test. Lesson: election logic and write gating must be the same code path.

Quick recall

Everything you need if you only revisit this box.

  1. Only one leader should accept writes per shard — quorum election enforces this.
  2. Randomized timeouts and terms prevent split votes and stale leader writes.
  3. Fencing tokens and term checks reject writes from recovered zombie leaders.

Test yourself

Answer these before moving on — recall is what makes it stick.