Skip to content

CONCEPT Cited by 2 sources

Distributed lease

A distributed lease is a time-bounded ownership claim over a resource (a connection, a partition, a shard, a "leader" role, a file) that a process must continuously renew to keep. If the holder stops renewing — because it crashed, was partitioned, or was killed — the lease expires and another process may take over. The lease is the standard way to grant exclusive ownership in a system where owners fail independently and no external observer can reliably tell "dead" from "slow."

Definition

A lease has three moving parts:

  1. An owner identity — who holds it (lease_owner = a worker/node ID).
  2. An expiry time — when the claim lapses if not renewed (lease_expires_at_ms).
  3. A renewal (heartbeat) — the holder periodically pushes the expiry forward.

The invariant it buys you is at-most-one live owner: as long as the lease duration exceeds the maximum time a new owner will wait before assuming the old one is dead, two processes cannot both believe they hold a live lease simultaneously. This is the same safety property as leader election — a lease is, in effect, a leader election scoped to a single resource, and per-resource leases give you fine-grained single-ownership across a whole fleet rather than one global leader.

How it is enforced

Every lease transition (acquire, renew, release, reclaim) is a compare-and-set so that concurrent contenders resolve atomically. The dominant implementations:

  • Conditional write on a storage item (patterns/conditional-write). The acquire condition is "no lease exists or it has already expired"; the renew condition is "I am still the owner." A racing contender's write fails the precondition and it backs off. This is how the AWS DynamoDB-lease pattern works (ConditionExpression on update_item), and how Fly.io's Litestream redesign enforces one active replication writer per destination (S3/Tigris conditional writes — "CASAAS").
  • Purpose-built lease/coordination services — ZooKeeper ephemeral znodes, etcd leases, Chubby locks, Consul sessions. These bundle the heartbeat, the TTL, and the atomic swap into the coordination system itself. The trade-off the conditional-write approach makes is removing that dependency: the lease lives in a database or object store the application already uses.

Failure handling: the two paths

A well-designed lease system handles the full failure spectrum with two complementary mechanisms:

  • Unexpected termination (crash / partition / SIGKILL). The dead owner cannot release, so the lease simply expires after its duration. A reconciliation loop — any healthy peer periodically scanning for expired-but-still-wanted leases — reclaims it via the same conditional acquire. Worst-case recovery time is bounded by lease duration + reconciliation interval. In the AWS pattern: 20 s + up to 60 s ≈ 80 s.
  • Graceful shutdown (rolling deploy / scale-in / SIGTERM). The owner explicitly releases — the fast path. In the AWS pattern the worker sets lease_expires_at_ms = 0, making the lease look already-expired so peers reclaim it on their next reconciliation cycle instead of waiting for natural expiry. A shared shutdown signal (e.g. an asyncio.Event) fans out to every held lease so releases happen in parallel, keeping shutdown time roughly constant regardless of how many leases the process holds. Both paths converge on the same outcome — a healthy peer takes over — but graceful release is dramatically faster.

Tuning the durations

The two dials are lease duration and heartbeat (renewal) interval, and the essential ratio is heartbeat ≪ lease so several renewals can fail before the lease lapses. The AWS pattern uses 20 s lease / 5 s heartbeat (a 4:1 ratio → 4 renewal attempts before expiry) and a 60 s reconciliation interval. The trade-offs:

  • Shorter lease → faster failover but more sensitivity to transient network hiccups (false expirations) and higher renewal write cost.
  • Longer lease → cheaper and more hiccup-tolerant but slower failover.
  • Reconciliation interval trades recovery speed against read cost.

Clock skew: the sharp edge

A lease is only safe if all participants agree on time. When expiry is evaluated against a now value supplied by the calling process (as with DynamoDB conditional writes — the database does not inject a server-side clock), every holder's local clock must be reasonably synchronized. AWS Fargate tasks in one Region get Amazon Time Sync (skew within a few ms — negligible against a 20 s lease). Off managed infrastructure, you must verify NTP and, where skew can't be bounded, increase the lease duration by the maximum expected skew to prevent a slow clock from prematurely expiring a live lease. (The classic hardening beyond this is a monotonically increasing fencing token handed out with each lease, so a delayed write from a previously-evicted owner is rejected by the resource itself; the AWS pattern relies on the conditional write and clock sync rather than an explicit fencing token.)

When to reach for a lease

  • A resource must have exactly one owner at a time and owners fail independently (a WebSocket/upstream connection, a shard, a scheduled singleton job, a single-writer replication stream).
  • You want automatic failover without an operator in the loop.
  • You'd rather not run (or depend on) a separate coordination cluster — the lease can ride on a database/object store you already have, via conditional writes.

A lease is not the right tool for one-shot task dispatch — that's a queue's job (and a queue like SQS has no notion of continuous ownership, only delivery). Leases and queues compose well: the queue delivers a fast "new work" hint, and the lease decides who actually owns it.

Seen in

Last updated · 766 distilled / 2,225 read