Distributed Locks, Leases, and Fencing Tokens Explained

Updated on
10 min read

Distributed locks help multiple services coordinate work on a shared resource, but a lock granted over a network cannot make a paused or disconnected process stop running. This explainer is for developers, architects, and SREs choosing coordination patterns for workers, schedulers, and shared data. It explains how leases expire, why an old lock holder can become dangerous, and how a fencing token lets the protected resource reject stale writes.

Why Distributed Locks Are Being Discussed

As applications split work across replicas, more than one process may try to update the same account, run the same scheduled task, or rebuild the same index. A single-process mutex cannot coordinate those independent machines. Teams therefore turn to a database, coordination service, or cache to decide who may act.

The hard part is not granting a lock when everything is healthy; it is deciding what happens when a process pauses, a network request times out, or a coordinator fails over. A client may believe it still owns a lock after the system has reassigned it. Safe designs must account for both the coordination decision and the effect on the resource being protected.

What Are Distributed Locks, Leases, and Fencing Tokens?

A distributed lock is a coordination record that gives one participant at a time permission to perform some operation. The lock service is authoritative: clients must acquire the lock through it rather than relying on local state or a shared clock.

A lease is a lock or ownership grant that lasts for a bounded period. The owner renews it while working; if renewals stop, the coordinator eventually treats it as expired. Leases can recover from crashed clients, but they cannot prove that an expired client has actually stopped.

A fencing token is a monotonically increasing number associated with each new owner. The client sends that token with every write to the protected resource. The resource remembers the highest accepted token and rejects requests carrying an older one. The token—not the client’s belief that its lease remains valid—lets the resource distinguish a current owner from a stale one.

The distinction matters: a lock service can serialize acquisition, while fencing protects the downstream system from delayed work that arrives after ownership changes.

The Problem Distributed Locks Solve

Without coordination, two workers can both read that a job is available, both begin processing it, and both write a result. This can lead to duplicate billing, conflicting inventory reservations, concurrent schema changes, or a cache rebuild that publishes inconsistent output.

In a single process, a mutex or semaphore can serialize access. Across machines, each process has its own memory and may lose contact with the others. A central coordinator can establish an ordering, but network timeouts are ambiguous: the coordinator may have granted a lock even if the reply never reached the client.

This is one example of the partial failures covered in distributed system failures and network partitions. A timeout says that a response did not arrive before a deadline; it does not prove that a remote operation failed or that a former owner has stopped.

How Distributed Locks Work

A client asks the coordinator to acquire a lock associated with a resource name. A correct coordinator performs the ownership change atomically—for example, by comparing that the lock key is absent and creating it in one transaction. If another owner already holds it, the request waits, retries with backoff, or fails.

To recover from a client crash, many systems associate ownership with a lease. The coordinator expires the lease when its TTL elapses without renewal and allows another client to acquire the lock. Renewal must be bounded by the service’s own deadline: a client should stop initiating new work when it cannot confirm that it still holds ownership.

However, expiration cannot stop a process that has been paused by the operating system, garbage collection, a virtual machine suspension, or a network partition. Imagine worker A acquiring a 15-second lease and then pausing for 20 seconds. Worker B may acquire the expired lock and begin work. If A resumes and sends a delayed write, both workers may now affect the resource.

Fencing addresses that race:

  1. The coordinator grants worker A token 41.
  2. A pauses long enough for its lease to expire.
  3. The coordinator grants worker B a newer token, 42.
  4. The resource accepts B’s write and records 42.
  5. A’s delayed write arrives with 41; the resource rejects it as stale.

The protected resource must check the token atomically with the write. Merely logging a token or checking it in the client does not fence anything. If a downstream service cannot enforce an ordering token or an equivalent conditional update, the lock alone cannot guarantee safety against a stale owner.

Locking Models and Key Components

Feature Consensus-backed coordinator Redis-style TTL lock Storage-native conditional write
Acquisition Atomic compare-and-create or ordered lock queue Set a unique owner value only if the key is absent Update only if the stored version or ETag still matches
Expiration Often tied to a coordinator-managed lease Key expires after a TTL Usually no lease; conflict detected at write time
Ordering Can provide revisions or ordered contenders Random ownership value is not automatically a fencing token Storage version can act as a concurrency condition
Best fit Shared jobs or coordination across services Short-lived best-effort coordination with understood failure assumptions Preventing lost updates to one database or HTTP resource
Main risk Quorum loss can block acquisition or renewal Failover and timing assumptions can leave stale ownership ambiguous Does not by itself coordinate a long-running multi-step job

Consensus-backed systems such as etcd provide atomic transactions, revisions, and leases. Its API documentation describes transactions and revisions; the Go concurrency package provides lease-backed sessions and mutexes. Apache ZooKeeper’s lock recipe uses ephemeral sequential nodes so contenders can queue and watch the node immediately ahead of them.

Redis-style locks often use SET key owner NX PX ttl: NX avoids replacing an existing owner, while PX sets an expiry. Release should compare the stored owner value before deleting the key, so one client cannot remove another client’s lock. The Redis distributed locks guide documents this pattern and its assumptions. A random owner value protects release; it is not a monotonically increasing fencing token. Failover behavior and timing assumptions must be considered for the specific deployment.

Conditional writes can be simpler when the goal is to avoid overwriting newer data, rather than to reserve a resource during a long workflow. For an HTTP resource, the HTTP If-Match precondition in RFC 9110 lets a client make a request only if the entity tag still matches. Database version columns provide a similar compare-and-swap pattern. These techniques do not coordinate separate systems or keep a long-running task exclusive.

Real-World Use Cases

  • Scheduled jobs: Ensure only one worker performs a daily aggregation, while making the job safe to resume or retry.
  • Resource provisioning: Avoid two controllers creating or deleting the same cloud resource at once.
  • Inventory or account updates: Coordinate a multi-step reservation where duplicate execution would violate a business invariant.
  • Cache or index rebuilds: Prevent competing workers from publishing different generations of a derived dataset.
  • Leader work assignment: Allow a current leader to schedule tasks while rejecting writes from an earlier leader after failover.

Not every duplicate task needs a distributed lock. An idempotency key, unique database constraint, queue visibility timeout, or conditional update can be simpler and more robust when it matches the operation’s actual invariant.

Practical Guide: Acquire an etcd Lease and Fence Writes

For a local demonstration, make an etcd endpoint available at localhost:2379, then initialize a small Go module, install the etcd client, and check connectivity:

go mod init example.com/etcd-lock-demo
go get go.etcd.io/etcd/client/v3
etcdctl --endpoints=localhost:2379 endpoint health

This program acquires an etcd mutex through a lease-backed session and prints the revision associated with the lock. Use a protected database or service that checks this revision before accepting each write; acquiring a mutex alone does not fence a downstream resource.

package main

import (
	"context"
	"fmt"
	"time"

	clientv3 "go.etcd.io/etcd/client/v3"
	"go.etcd.io/etcd/client/v3/concurrency"
)

func main() {
	ctx, cancel := context.WithTimeout(context.Background(), 20*time.Second)
	defer cancel()

	client, err := clientv3.New(clientv3.Config{
		Endpoints:   []string{"localhost:2379"},
		DialTimeout: 5 * time.Second,
	})
	if err != nil {
		panic(err)
	}
	defer client.Close()

	session, err := concurrency.NewSession(client, concurrency.WithTTL(15))
	if err != nil {
		panic(err)
	}
	defer session.Close()

	mutex := concurrency.NewMutex(session, "/locks/rebuild-index")
	if err := mutex.Lock(ctx); err != nil {
		panic(err)
	}

	fence := mutex.Header().Revision
	fmt.Printf("Lock acquired; send fencing token %d with protected writes\n", fence)

	if err := mutex.Unlock(ctx); err != nil {
		panic(err)
	}
}

In a relational database, store the largest accepted fencing token with the protected row and include the comparison in the same atomic update. Initialize fencing_token to zero and check the number of affected rows; zero rows means the request was stale or the resource did not match:

UPDATE protected_resource
SET value = $1,
    fencing_token = $2
WHERE resource_id = $3
  AND fencing_token <= $2;

If this protects a service outside your database, that service must provide an equivalent atomic check. Ensure the new owner establishes its token at the resource before performing dependent work; otherwise, a delayed old write that arrives before any newer token has been recorded may still be accepted. In production, use authenticated TLS endpoints, monitor lease-renewal failures, set request deadlines, and test process pauses and coordinator unavailability. The distributed systems resilience guide covers retries, recovery, and failure isolation that complement this design.

Common Misconceptions

“The TTL guarantees that only one process is running.” A lease expiry changes the coordinator’s view of ownership; it does not terminate a paused process. A new owner can overlap with the old one.

“A random lock ID is a fencing token.” A random value can prove which client created a key and make safe release possible. It does not order owners. Fencing requires monotonically increasing values and enforcement by the resource that receives writes. Martin Kleppmann’s analysis of distributed locking and fencing tokens explains why this distinction matters when clients pause.

“One lock service is always the simplest safe choice.” If a database update can use a version check, or a task can be idempotent, those approaches may remove an entire dependency. The CAP theorem explainer helps frame why coordination can become unavailable when a quorum cannot be reached; the correct failure behavior depends on the invariant being protected.

Changelog and Last Updated

  • 2026-10-09: First published with an etcd lease example and fencing-token guidance.
TBO Editorial

About the Author

TBO Editorial writes about the latest updates about products and services related to Technology, Business, Finance & Lifestyle. Do get in touch if you want to share any useful article with our community.