Short answers to what readers ask most about this topic.
01How does leader election work in distributed systems?
Nodes agree on one of themselves to do work that must not run twice. In practice the winner takes a lease, a lock with an expiry, in a shared store such as a database row, a Redis key or etcd, and keeps renewing it. If it stops renewing, the lease expires and another node takes over with a higher fencing token.
02What is the difference between leader election and a distributed lock?
A distributed lock protects one critical section for a short time, while leader election gives one node a long-lived role that it must keep renewing. Both are usually built from the same lease mechanism. The failure mode is identical, a stale holder acting after it lost the lock, so both need fencing tokens when correctness matters.
03What is a fencing token and why do I need one?
A fencing token is a number that increases every time a lease is granted. The leader sends it with each write and the resource rejects any token lower than one it has already seen. You need it because a paused leader can wake after its lease expired and still believe it leads.
04Can I use Redis SET NX PX for leader election?
Yes, for best-effort cases such as a cron job where an occasional double run is harmless. Store a random token and release with a script that checks it. Redis gives you no fencing token and expiry depends on the Redis clock, so do not use it alone to protect data that must never be written by two nodes.
05How does Kubernetes leader election work?
Components such as kube-controller-manager and kube-scheduler compete for a Lease object in the coordination.k8s.io API group, and only the holder is active. The controller manager defaults are a 15 second lease duration, a 10 second renew deadline and a 2 second retry period. Your own controllers can use the same Lease mechanism.
Leader Election in Distributed Systems: Leases and Fencing
How leader election works in distributed systems: a Postgres lease with a fencing epoch, Redis and etcd alternatives, Kubernetes Lease objects and split-brain.
Leader election makes exactly one node run work that must never run twice, such as a singleton job or a primary database. The simplest safe design is a lease: a row or key with an expiry that the leader keeps renewing, plus a fencing token that rises on every takeover so a stale leader's writes are rejected.
Run a NestJS app as two Docker replicas instead of one and the nightly job that sweeps unpaid invoices quietly starts running twice. Nothing crashes. Both copies read the same rows, both send the same reminder, and the customer gets two messages.
This post answers the question behind that bug: how does leader election work in distributed systems? It builds a runnable lease in Postgres with real psql output, compares it with Redis, etcd and the Kubernetes Lease object, and spends most of its time on the part people skip, which is what happens when the old leader does not know it has been replaced. The sources are Kleppmann's essay on distributed locking, the Kubernetes and etcd docs, the Redis SET reference and the Raft paper.
Why does a distributed system need a single leader?
Some work is only correct when one actor does it at a time. Running it on every node is not redundancy, it is duplication with side effects. Three common cases:
A singleton job: a cron-style sweeper, a report generator or a queue-draining task that must not be sent twice.
A primary database: one node accepts writes and the replicas copy it, so two primaries means two diverging histories.
A coordinator: the node that assigns partitions, shards or work to the others and must hand out each assignment once.
Leader election is the protocol that chooses that one node, lets everyone else find out who it is, and picks a replacement when it disappears. The hard part is never the choosing. It is the replacement, because a node you cannot reach might be dead, or might just be slow, and the cluster cannot tell the difference. I wrote about a database flavour of this in the replication post, and about a sketch of an advisory-lock leader in the distributed job scheduler design post; this one goes under both.
How does a lease in a shared store elect a leader?
A lease is a lock with an expiry date. A node writes its name into a shared place together with a deadline, and it stays leader only while it keeps pushing the deadline forward. If it stops, the deadline passes and anyone else may take over. Expiry is what turns a lock into something that survives a crashed holder. The place can be a database row, a Redis key or a consensus store:
Where the lease lives
How it expires
What to watch for
Postgres row
An expires_at column compared with the database now()
One clock only, and the epoch column gives you a fencing token for free
Postgres advisory lock
Released when the session ends, no deadline
A paused but still connected node keeps the lock, and there is no token
Redis SET NX PX
The key expires after the PX milliseconds
Fast and simple, but no fencing token and the TTL runs on the Redis clock
etcd, ZooKeeper, Consul
A lease or session with a TTL kept alive by the client
Consensus-backed and ordered, but a separate cluster to run
The Redis reference documents the minimal form: SET with NX (only if the key is absent) and PX (expire in milliseconds), and it says to store a random token and release with a script that checks the token first, so a late client cannot delete the key a newer client created. The same page notes that this plain SET pattern is discouraged in favour of Redlock, a position Kleppmann disputes, which is the next section's theme.
# Redis: acquire. NX = only if absent, PX = expire in milliseconds.
SET lock:invoice-sweeper node-a-7f3c NX PX 15000
# Redis: release. Compare the token first, or you delete someone else's lock.
# (Lua script, run with EVAL ... 1 lock:invoice-sweeper node-a-7f3c)
if redis.call("get", KEYS[1]) == ARGV[1] then
return redis.call("del", KEYS[1])
else
return 0
end
How do you implement lease election in TypeScript with Postgres?
For a stack that already has Postgres, such as NestJS and Postgres on a single VPS, a table with four columns is the whole coordination service. The trick is one UPDATE whose WHERE clause says who may take the lease. Postgres locks the row, so two nodes racing on the same statement cannot both match. The epoch column goes up by one on every takeover and never on a renewal; that number is the fencing token.
CREATE TABLE leader_lease (
name text PRIMARY KEY,
holder text,
epoch bigint NOT NULL DEFAULT 0, -- the fencing token
expires_at timestamptz NOT NULL DEFAULT '-infinity'
);
INSERT INTO leader_lease (name) VALUES ('invoice-sweeper');
import { Pool } from "pg";
const pool = new Pool({ connectionString: process.env.DATABASE_URL });
const NAME = "invoice-sweeper";
const LEASE_SECONDS = 15;
const TICK_MS = 5_000; // one third of the lease: two missed ticks are survivable
const GIVE_UP_MS = 10_000; // stop leading before a rival can legally take over
// Compare-and-swap: the WHERE clause is the whole election.
// Only one concurrent UPDATE can match, because Postgres re-checks the
// predicate after taking the row lock.
const ACQUIRE = "UPDATE leader_lease SET holder = $2, epoch = epoch + 1, " +
"expires_at = now() + make_interval(secs => $3) " +
"WHERE name = $1 AND (holder IS NULL OR holder = $2 OR expires_at < now()) " +
"RETURNING epoch";
// Renewal must NOT bump the epoch, and must name the epoch we believe we hold.
const RENEW = "UPDATE leader_lease SET expires_at = now() + make_interval(secs => $3) " +
"WHERE name = $1 AND holder = $2 AND epoch = $4 AND expires_at > now() " +
"RETURNING epoch";
export function runElection(
nodeId: string,
lead: (epoch: number, signal: AbortSignal) => Promise<void>,
): () => void {
let epoch: number | null = null;
let lastRenewOk = 0;
let abort = new AbortController();
const stepDown = () => {
epoch = null;
abort.abort(); // the job must observe this signal and stop
abort = new AbortController();
};
const tick = async () => {
try {
if (epoch === null) {
const r = await pool.query(ACQUIRE, [NAME, nodeId, LEASE_SECONDS]);
if (r.rowCount === 1) {
epoch = Number(r.rows[0].epoch); // bigint arrives as a string
lastRenewOk = Date.now();
void lead(epoch, abort.signal);
}
return;
}
const r = await pool.query(RENEW, [NAME, nodeId, LEASE_SECONDS, epoch]);
if (r.rowCount === 1) lastRenewOk = Date.now();
else stepDown(); // someone else holds a newer epoch
} catch {
// Database unreachable: we cannot prove we still lead. Stop on OUR clock
// before the lease can expire on the database's clock.
if (epoch !== null && Date.now() - lastRenewOk > GIVE_UP_MS) stepDown();
}
};
const timer = setInterval(tick, TICK_MS);
void tick();
return () => {
clearInterval(timer);
if (epoch !== null) stepDown();
};
}
I ran these exact statements in psql against a local PostgreSQL 16.15, with node-a and node-b standing in for two app instances. This is the output, trimmed to the result rows:
-- 1. node-a acquires the empty lease
holder | epoch | live
--------+-------+------
node-a | 1 | t
-- 2. node-b tries while node-a holds it
holder | epoch
--------+-------
(0 rows)
-- 3. node-a renews (same epoch)
holder | epoch
--------+-------
node-a | 1
-- (expiry forced with: UPDATE leader_lease SET expires_at = now() - interval '1 second')
-- 4. node-b takes over after expiry
holder | epoch
--------+-------
node-b | 2
-- 5. node-a wakes up and renews with its stale epoch 1
holder | epoch
--------+-------
(0 rows)
Step 2 is the election working: node-b's UPDATE matched zero rows because node-a's lease was live. Step 5 is the part that matters most. node-a, resurrected with its old belief, tried to renew epoch 1 and matched zero rows, so it learns it was replaced at its next tick instead of carrying on as a second leader. The TypeScript loop turns zero rows into stepDown().
What is split-brain, and how do fencing tokens stop it?
Split-brain means two nodes both believe they are the leader. A lease reduces the chance but cannot remove it, because the leader checks the lease, then does the work, and the gap between those two moments is unbounded. Kleppmann's example is a client that holds the lock, then stalls in a long garbage collection pause, resumes after its lease expired and writes to storage as if nothing happened. He points out that checking the expiry just before the write does not help, because a pause can strike at any point.
His fix is a fencing token: a number that increases every time the lock is granted. The client sends it with every write and the storage layer rejects any write whose token is lower than one it has already seen. In the schema above the epoch is that number. The key requirement is that the check lives inside the resource being protected, in the same statement as the write, not in the leader's own code:
CREATE TABLE sweep_state (id int PRIMARY KEY, last_epoch bigint NOT NULL, note text);
INSERT INTO sweep_state VALUES (1, 2, 'written by node-b');
-- Wrong: a stale leader (epoch 1) writes unconditionally and wins by arriving last.
UPDATE sweep_state SET note = 'stale write by node-a' WHERE id = 1;
-- Right: the write carries its epoch and the resource refuses anything older.
UPDATE sweep_state SET last_epoch = 1, note = 'stale write by node-a'
WHERE id = 1 AND last_epoch <= 1 RETURNING *;
-- (0 rows) rejected: the table has already seen epoch 2
UPDATE sweep_state SET last_epoch = 2, note = 'sweep by node-b'
WHERE id = 1 AND last_epoch <= 2 RETURNING *;
-- id | last_epoch | note
-- ----+------------+-----------------
-- 1 | 2 | sweep by node-b
A lease alone is not mutual exclusion. If the thing you protect cannot compare a token, for example an external email API, you only have best-effort exclusion, so make the action idempotent too. Kleppmann says a lock used purely for efficiency can be loose, but one used for correctness needs a token.
How do clock skew and process pauses break a lease?
A lease depends on time, and time is the least reliable thing in a cluster. Kleppmann notes that Redis expiry relies on the system clock, which can jump, so a key can expire sooner or later than intended. Process pauses are worse because they need no clock fault at all. Here is a worked timeline with a 15 second lease and a 20 second stall:
t = 0: node-a holds epoch 1 and its lease runs until t = 15.
t = 1: node-a freezes, for example a stop-the-world pause or a VM migration, for 20 seconds.
t = 15: the lease expires and node-b takes over with epoch 2.
t = 21: node-a wakes up still believing it leads and writes with epoch 1. Without a fencing check this write lands; with one it is rejected.
The timing constants are arithmetic, not luck. With a 15 second lease and a tick every 5 seconds, a leader gets three attempts per lease window (at 0, 5 and 10 seconds). If it has had no successful renewal for 10 seconds it steps down by itself, which leaves a 5 second gap before a rival may legally take over at 15. Failover takes 15 to 20 seconds in the worst case, because the lease must expire first and the rival only notices at its next 5 second tick. With the Postgres version, skew between app servers does not matter, since every deadline is written and compared with the database's own now().
Choose the lease duration from the failover time you can tolerate, then set the renewal tick to a third of it. Shorter leases fail over faster but make a slow database round trip look like a dead leader.
When do you need etcd, ZooKeeper or Raft instead of a database row?
A row in one Postgres is only as available as that Postgres. If the database is down nobody can renew, so every leader steps down, which is safe but it halts the work. A consensus store keeps its state on several replicas and agrees on an order. The etcd documentation describes leases that carry a time to live and are kept alive by the client, keys that die with their lease, a store-wide revision counter that rises on every change, and transactions that give compare-and-swap. That revision counter is a ready-made fencing token, and Kleppmann recommends a consensus system such as ZooKeeper plus fencing tokens whenever the lock protects correctness.
Under those systems sits a leader election of their own. The classic bully algorithm lets the live node with the highest ID win, and it assumes you can reliably detect who is alive. Raft elects a leader using terms and randomised election timeouts, and a candidate needs votes from a majority. The term number plays the same role as our epoch: any message carrying an older term is ignored. The Raft paper is the primary source, and the consensus mechanics deserve their own post.
My rule of thumb: if you run one Postgres anyway and a 20 second failover is fine, use the row. If you already operate etcd or ZooKeeper, or the leader decides who writes to a primary database, use the consensus store, because there the cost of two leaders is data loss.
How does Kubernetes do leader election?
Kubernetes uses the lease pattern in its own control plane. The Leases page lists leader election among the uses of the Lease object: in a high availability setup, kube-controller-manager and kube-scheduler run several instances while only one is active, and your own controllers can do the same. The kube-controller-manager reference lists the knobs and their defaults: a 15 second lease duration, a 10 second renew deadline and a 2 second retry period, with leases as the lock type.
# kube-controller-manager defaults, from the Kubernetes reference
kube-controller-manager \
--leader-elect=true \
--leader-elect-lease-duration=15s \
--leader-elect-renew-deadline=10s \
--leader-elect-retry-period=2s \
--leader-elect-resource-lock=leases
# The lock object lives in the coordination.k8s.io/v1 API group:
apiVersion: coordination.k8s.io/v1
kind: Lease
spec:
holderIdentity: <pod-name>_<uuid>
leaseDurationSeconds: 15
renewTime: <timestamp>
Read the defaults as the same arithmetic as before. The renew deadline of 10 seconds is shorter than the 15 second lease, so a leader stops leading before a rival may take over, and retries every 2 seconds give it about five attempts inside that window. The Lease page also advises naming your lease after your component, for example example-foo, so two operators do not collide on one object.
Which approach should you choose? A checklist
Work through these in order and stop at the first that applies:
Is it fine if the job sometimes runs twice? Make it idempotent and skip election entirely.
Do you already run Postgres and accept a 15 to 20 second failover? Use the lease row with an epoch.
Are you on Kubernetes? Use the Lease API through a client library rather than building your own.
Does a second leader corrupt data, such as a database primary? Use etcd or ZooKeeper and a fencing token.
Whatever you pick, can the protected resource reject a stale token? If not, you only have a best effort lock.
A leader is a lease, not a title: it must be renewed, it can be lost without the holder noticing, and so every write needs a token that the resource itself checks. Build the lease in the store you already run, stop leading on your own clock before the lease ends, and reach for a consensus store only when two leaders would cost you data.