Skip to content

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

0 of 1 checks done
  1. 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.

  2. 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.

  3. 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 s
    Worker AWorker BStorage
    1. 0 sWorker ATakes the lease with fencing token 33
    2. 2 sWorker AStalls (GC pause) for 12 s
    3. 10 sLockA's lease expires
    4. 10 sWorker BTakes the lease with fencing token 34
    5. 11 sWorker BWrites the job's result
    6. 11 sStorageAccepts B's write (token 34)
    7. 14 sWorker AResumes, still believes it holds the lease, and writes
    8. 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.

  4. 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.

  5. Check

    Why isn't 'the owner stops working when it fails to renew' (self-fencing) enough on its own?

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

Write at least 60 characters (0 so far). Write it as you would say it in a design review. You will compare it against the points a strong answer makes.

Where you practise it

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.

  • 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.