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.
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
- Leader
node-Aonledger-55suffers JVM OOM at 14:32:01. - Followers
B,C,D,Estop receiving heartbeats; election timers start. Ctimes out first at 14:32:01.220, increments term to 47, requests votes.Creceives votes fromDandE(3 of 5 = majority); becomes leader term 47.Csends heartbeats;Bcatches up and acknowledges new leader.- ShardPay routing service watches metadata; clients receive updated leader endpoint by 14:32:03.
- Transfers that failed with
NOT_LEADERretry with idempotency keys — no duplicates.
// 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.priorityhigher 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
fenceTokenon 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.
- Only one leader should accept writes per shard — quorum election enforces this.
- Randomized timeouts and terms prevent split votes and stale leader writes.
- Fencing tokens and term checks reject writes from recovered zombie leaders.
Test yourself
Answer these before moving on — recall is what makes it stick.