Leases and fencing tokens
Ownership that expires unless renewed, plus a token that lets the rest of the system reject an owner that has lost its claim without knowing it.
Distribution
Learn it
Some work must have exactly one owner at a time: a job, a document's sequencer, a leader. Owners crash, so ownership must be reclaimable. But a crashed owner looks exactly like a slow one: a garbage-collection pause, a network blip or a frozen VM looks like death. Take ownership from a slow owner and, when it wakes, you have two.
Leases handle the crash: ownership is granted until a time (
leased_until = now() + 30s) and renewed by heartbeats. If renewals stop, the lease lapses and someone else may claim it, with no failure detector or cooperation from the dead process. Compute expiry with one clock (the database's or coordinator's), so skew between machines doesn't matter.Freeze worker A past its lease and see what its late write does, first without fencing, then with it.
A lease that runs out while its holder is asleep. The lease lasts 10 s and A renews it every 3 s, but A freezes at 2 s (a GC pause, a VM migration, a slow disk). Change how long A is frozen, and whether storage checks fencing tokens. 12 sWorker AWorker BStorage- 0 sWorker ATakes the lease with fencing token 33
- 2 sWorker AStalls (GC pause) for 12 s
- 10 sLockA's lease expires
- 10 sWorker BTakes the lease with fencing token 34
- 11 sWorker BWrites the job's result
- 11 sStorageAccepts B's write (token 34)
- 14 sWorker AResumes, still believes it holds the lease, and writes
- 14 sStorageAccepts A's write and overwrites B's result
Two workers acted as the owner. A's lease expired during the pause, B took over, and storage accepted A's late write anyway. The job's result is now whatever A wrote last.
Fencing handles the slow owner. Every grant comes with a token that increases (an epoch, a version), every write carries it, and the system holding the data rejects stale tokens:
UPDATE jobs SET status = 'done' WHERE id = $1 AND lease_token = $2;The old owner's late write matches zero rows, and it learns it lost. The resource enforces exclusivity; the lease only advises it.
Check
Why isn't 'the owner stops working when it fails to renew' (self-fencing) enough on its own?Writes that can't carry a token (files in object storage, calls to third parties) should go to per-owner locations or be deduplicated, so overlap between two owners can't corrupt them.
Quick reference
The same ideas, condensed for revision.
How it goes wrong
- Lease without fencing
- A paused owner resumes after expiry and overwrites the new owner's work.
- Lease too short
- Normal pauses cause spurious handovers and duplicate work.
- Lease too long
- Recovery after a real crash takes the full duration.
- Clock comparison across machines
- A worker with a fast clock thinks leases have expired early.
Instead, consider
- Consensus-based leadership (Raft, etcd, ZooKeeper)
- You need one leader for a cluster with strong guarantees; these systems provide leases and fencing epochs.
- Idempotent, overlapping work
- Two owners doing the same work is harmless, so exclusivity does not need enforcing.
In practice
- Lease columns in a jobs table
- leased_until + lease_token, claimed with a conditional update.
- Queue visibility timeouts
- A lease on a message; extend it for long work.
- etcd/ZooKeeper sessions and revisions
- Leases with monotonically increasing revisions usable as fencing tokens.
It assumes
- The resource being protected can check a token on each write (a conditional update, a versioned API).
- Lease duration comfortably exceeds the renewal interval plus expected pauses.
- Expiry is judged by a single clock, not compared across machines.
Explain it in your own words
Where you practise it
- A reliable video processing pipeline
Trace the happy path end to end · A worker dies 14 minutes in · Two workers, one job · Write the claim and the completion
- Notifications across email, push and in-app
- A look-aside cache at Facebook's scale
A stale value that never leaves · Write the lease-aware read · Defend 'best-effort eventual consistency'
- A real-time collaborative editor
Further reading
Engineers describing it in systems they run.
- Making multiplayer more reliable
Figma · Darren Tsung · Post, Oct 2022
Closing the gap between periodic saves: a durable journal of small changes, and one owner per document's history.
- Scaling Memcache at Facebook
Meta (Facebook) · Rajesh Nishtala and others · Paper, Apr 2013
The reference on running a look-aside cache hard: leases, invalidation from the commit log, failover without hammering the database, and consistency across regions.
Related concepts
- Concurrency control
Making read-decide-write sequences safe when other actors may change the same data in between: locks, conditional writes and constraints.
- Message queues
A durable buffer between producers and consumers that hands each message to one consumer at a time and redelivers it unless acknowledged.
- Idempotency
Designing an operation so that performing it twice has the same effect as performing it once, which is what makes retries safe.
- Ordering
There is no global 'now' in a distributed system. Order exists only where something assigns it, so decide which order you need and who assigns it.
- Asynchronous processing
Accepting a request, recording the work durably, and doing it later in another process, so the work can outlive the request.
- Partitioning
Splitting data or work by key so each part is handled independently: scaling out, and giving each key a single owner.