Table of Contents

    distributed locks

    SYSTEM DESIGN • CHAPTER 15.4

    Distributed Locks

    Understand why mutual exclusion across machines is fundamentally different from mutual exclusion within a process, when a distributed lock is genuinely required, and why most systems that use one should not.

    Learning objective: By the end of this article, you will understand efficiency versus correctness locking, why a lock cannot guarantee exclusion without resource cooperation, common implementations and their flaws, the Redlock controversy, and the alternatives that avoid locking entirely.

    Prerequisites

    Recommended Knowledge

    • Leases and expiry semantics
    • Fencing tokens and resource-side validation
    • Leader election and the stale holder problem
    • Consensus and quorum intersection
    • Partial failure and timeout ambiguity
    • Process pauses and clock skew
    • Idempotency and duplicate handling
    • Optimistic and pessimistic concurrency control

    Why This Is Not Like a Mutex

    A mutex within a process works because the operating system knows definitively which thread holds it, and a thread cannot vanish while holding it without the kernel noticing.

    Property Local Mutex Distributed Lock
    Holder identity Known with certainty Believed, never confirmed
    Holder death Detected by the kernel Indistinguishable from slowness
    Release guarantee Enforced on thread exit Requires expiry
    Communication failure Impossible Routine
    Acquisition cost Nanoseconds Network round trip
    Exclusion guarantee Absolute Probabilistic without fencing

    Simple Analogy

    A sign on a door reading "occupied" works only if everyone agrees to read it, nobody falls unconscious inside, and the sign cannot be left behind by someone who has already gone. Distributed locks violate all three routinely.

    A distributed lock does not create mutual exclusion. It creates a shared belief about mutual exclusion, which is a considerably weaker thing.

    Efficiency Versus Correctness

    This distinction determines everything that follows. Before choosing an implementation, decide which of these two problems you are actually solving.

    Efficiency Locking

    • Avoids duplicated work
    • Occasional overlap is merely wasteful
    • Correctness does not depend on it
    • A simple implementation suffices

    Correctness Locking

    • Prevents data corruption
    • Overlap causes permanent damage
    • Correctness depends entirely on it
    • Requires fencing at the resource
    Use Case Category Consequence of Overlap
    Cache warming Efficiency Wasted compute
    Scheduled report generation Efficiency Duplicate report, overwritten
    Sending a notification Efficiency User receives it twice
    Financial ledger write Correctness Balance corrupted
    Inventory decrement Correctness Overselling
    File compaction Correctness Data loss
    Schema migration Correctness Inconsistent database state
    THE DECIDING QUESTION
    If two holders acted simultaneously, would you lose money or data? If yes, a lock alone is insufficient and you need fencing. If no, a simple lock is fine.

    Why Locks Alone Cannot Guarantee Exclusion

    The lock service may be perfectly correct and still fail to provide exclusion, because the gap between holding a lock and acting on it is not instantaneous.

    THE VULNERABLE GAP
    Acquire Lock Arbitrary Delay Write to Resource
    The Pause Sequence Client A acquires the lock, then pauses for garbage collection. The lock expires. Client B acquires it and writes. Client A resumes, still believing it holds the lock, and writes over B's change.
    Delay Source Typical Magnitude Detectable by the Client
    Garbage collection pause Milliseconds to seconds No, time appears continuous
    Virtual machine suspension Seconds to minutes No
    Disk or dependency stall Seconds Only as slowness
    Network retransmission Seconds No
    CPU scheduling starvation Milliseconds to seconds No
    Page fault storm Seconds No
    No timeout value fixes this: Process pauses have no upper bound. Extending the lease reduces the probability of overlap but never eliminates it, and a longer lease worsens recovery time.

    Fencing Is the Actual Solution

    Since the client cannot know whether it still holds the lock, the resource must decide. A fencing token makes the resource the arbiter of exclusion.

    RESOURCE-SIDE ENFORCEMENT
    \[ \text{Accept} \iff token_{\text{presented}} \geq token_{\text{last accepted}} \]
    -- Lock service issues a monotonic token per acquisition
    UPDATE distributed_locks
    SET holder_id     = :client_id,
        fencing_token = fencing_token + 1,
        expires_at    = CURRENT_TIMESTAMP + (:ttl_ms * INTERVAL '1 millisecond')
    WHERE lock_name = :lock_name
      AND (expires_at < CURRENT_TIMESTAMP OR holder_id = :client_id)
    RETURNING fencing_token;
    
    -- The resource rejects anything stale
    UPDATE inventory
    SET quantity      = quantity - :amount,
        fencing_token = :token
    WHERE sku           = :sku
      AND fencing_token <= :token
      AND quantity      >= :amount;
    async function performGuardedWork(lockService, resource, work) {
        const acquisition = await lockService.acquire({
            lockName: work.lockName,
            clientId: work.clientId,
            ttlMs: work.ttlMs
        });
    
        if (!acquisition.acquired) {
            return { done: false, reason: "lock held elsewhere" };
        }
    
        try {
            const outcome = await resource.apply(work.payload, {
                fencingToken: acquisition.fencingToken
            });
    
            if (!outcome.accepted) {
                return {
                    done: false,
                    reason: "fenced",
                    presented: acquisition.fencingToken,
                    required: outcome.requiredToken
                };
            }
    
            return { done: true, result: outcome.result };
        } finally {
            await lockService.release({
                lockName: work.lockName,
                clientId: work.clientId,
                fencingToken: acquisition.fencingToken
            });
        }
    }
    What Changes With Fencing The lock becomes an optimization that prevents wasted contention. The token becomes the correctness guarantee. Overlap still occurs, but it becomes harmless.
    When the Resource Cannot Fence Many external systems accept any authenticated write with no version check. For these, exclusion cannot be guaranteed at all, and you must rely on idempotency keys and accept residual risk.

    Implementation Approaches

    Backend Mechanism Safety Under Failover
    Consensus service Ephemeral node or compare-and-swap Safe, replicated durably
    Relational database Conditional update on a lock row Safe if replication is synchronous
    Database advisory lock Session-scoped lock primitive Lost on connection or failover
    Single cache node Set-if-not-exists with expiry Unsafe, grant can be lost
    Replicated cache quorum Majority of independent nodes Disputed, see below
    Object storage Conditional write on generation Safe, higher latency

    The Single-Node Cache Hazard

    Asynchronous Replication Loses Locks A client acquires a lock on the primary. The primary fails before replicating. A replica is promoted with no record of the lock, and a second client acquires the same lock. Both believe they hold it.
    async function acquireSimpleLock(cache, lockName, clientId, ttlMs) {
        const acquired = await cache.set(lockName, clientId, {
            notExists: true,
            expiryMs: ttlMs
        });
    
        return { acquired, clientId };
    }
    
    async function releaseSimpleLock(cache, lockName, clientId) {
        const script = `
            if redis.call('get', KEYS[1]) == ARGV[1] then
                return redis.call('del', KEYS[1])
            else
                return 0
            end
        `;
    
        return await cache.eval(script, 1, lockName, clientId);
    }
    Release must be conditional: Deleting the key unconditionally can release a lock another client now holds, if yours expired first. The comparison and delete must be atomic.

    The Redlock Debate

    Redlock acquires a lock on a majority of independent cache nodes, aiming to survive individual node failure without full consensus. Whether it is safe has been actively disputed.

    Concern Argument Against Counter-Argument
    Clock dependency Skew or jumps break the majority guarantee Clocks are adequately synchronized in practice
    Process pauses A pause invalidates the holder's belief True of every lease-based lock
    No fencing tokens Provides no monotonic counter Tokens can be layered on separately
    No durable log Nodes may lose state on restart Persistence can be configured
    Not formally proven No consensus-grade safety proof Sufficient for efficiency locking
    The Resolution The disagreement largely dissolves along the efficiency versus correctness line. For efficiency locking, Redlock is adequate. For correctness locking, no lock protocol suffices without fencing, so the debate is somewhat beside the point.

    Lock Duration and Renewal

    DURATION CONSTRAINT
    \[ T_{\text{lock}} > T_{\text{work p99}} + T_{\text{pause p99}} + T_{\text{acquire RTT}} \]
    Duration Risk Mitigation
    Shorter than the work Lock expires mid-operation Renew, or size the lock correctly
    Much longer than the work Crashed holder blocks others Release explicitly on completion
    Renewed indefinitely Hung holder never yields Cap total renewals and total hold time
    class RenewingLock {
        constructor(service, config) {
            this.service = service;
            this.config = config;
            this.held = false;
            this.renewalsRemaining = config.maxRenewals;
        }
    
        async acquire() {
            const result = await this.service.acquire({
                lockName: this.config.lockName,
                clientId: this.config.clientId,
                ttlMs: this.config.ttlMs
            });
    
            if (!result.acquired) return result;
    
            this.held = true;
            this.token = result.fencingToken;
            this.acquiredAt = performance.now();
            this.startRenewal();
    
            return result;
        }
    
        startRenewal() {
            this.timer = setInterval(async () => {
                if (this.renewalsRemaining <= 0) {
                    return this.abandon("maximum hold time exceeded");
                }
    
                const renewed = await this.service.renew({
                    lockName: this.config.lockName,
                    clientId: this.config.clientId,
                    token: this.token,
                    ttlMs: this.config.ttlMs
                });
    
                if (!renewed.ok) {
                    return this.abandon(renewed.reason);
                }
    
                this.renewalsRemaining -= 1;
            }, this.config.ttlMs / 3);
        }
    
        heldSafely() {
            const elapsed = performance.now() - this.acquiredAt;
            return this.held && elapsed < (this.config.ttlMs - this.config.marginMs);
        }
    
        abandon(reason) {
            clearInterval(this.timer);
            this.held = false;
            this.onLockLost(reason);
        }
    }
    CAP TOTAL HOLD TIME
    Unbounded renewal converts a crashed-holder problem into a hung-holder problem. A process stuck in a loop renews forever and blocks everyone else permanently.

    Contention and Fairness

    A lock serializes access, so throughput on the protected resource becomes independent of how many clients exist. Contention is not a tuning problem; it is the design.

    SERIALIZED THROUGHPUT
    \[ \text{Throughput} \leq \frac{1}{T_{\text{acquire}} + T_{\text{work}} + T_{\text{release}}} \]
    Problem Symptom Response
    High contention Most attempts fail to acquire Narrow the lock scope
    Starvation Some clients never succeed Queue-based fair acquisition
    Thundering herd All retry on release simultaneously Jittered backoff
    Convoy effect One slow holder delays a long queue Bound work performed under the lock
    Lock held during I/O Hold time dominated by waiting Move I/O outside the critical section
    Narrow the Scope First Locking per entity rather than globally usually eliminates contention entirely. One lock per product is very different from one lock for the inventory service.

    Alternatives Worth Preferring

    Most distributed locks exist because a simpler mechanism was overlooked. These alternatives avoid the coordination entirely.

    1

    Optimistic Concurrency

    Read a version, compute, and write conditionally on that version being unchanged. No lock is held, and conflicts are detected rather than prevented.

    2

    Database Transactions

    If all contended state lives in one database, its isolation machinery already provides exclusion, tested far more thoroughly than anything you will build.

    3

    Partitioned Ownership

    Route all operations for a key to one owner. Exclusion becomes a routing property rather than a runtime negotiation.

    4

    Idempotent Operations

    If repeating an operation is harmless, concurrent execution is also harmless, and no exclusion is required at all.

    5

    Atomic Primitives

    Conditional updates, atomic increments, and compare-and-swap perform exclusion inside a single operation with no coordination window.

    -- Optimistic concurrency: no lock, conflict detected on write
    UPDATE documents
    SET content = :new_content,
        version = version + 1,
        updated_at = CURRENT_TIMESTAMP
    WHERE document_id = :document_id
      AND version     = :expected_version;
    
    -- Atomic conditional decrement: exclusion without a lock
    UPDATE inventory
    SET quantity = quantity - :amount
    WHERE sku      = :sku
      AND quantity >= :amount;
    Situation Prefer Reason
    Single database involved Transaction Isolation already guarantees it
    Low contention expected Optimistic concurrency No coordination on the common path
    Work partitions by key Ownership routing Exclusion by construction
    Operation is repeatable Idempotency Overlap causes no harm
    Single value mutation Atomic primitive No window to exploit
    Genuinely cross-system exclusion Lock plus fencing No simpler option exists
    The best distributed lock is the one you found a way not to need. Every lock adds a coordination dependency, a failure mode, and a throughput ceiling.

    Failure Modes

    Failure Consequence Mitigation
    Holder pauses past expiry Two clients act concurrently Fencing tokens
    Lock service unavailable No client can proceed Decide fail-open or fail-closed explicitly
    Lock service failover Grant lost, lock double-issued Consensus-backed durable service
    Unconditional release Releases another client's lock Compare holder atomically on release
    Missing release on crash Lock held until expiry Always set an expiry
    Clock jump on holder Believes lock is still valid Monotonic clocks with margin
    Unbounded renewal Hung holder blocks forever Cap total hold duration
    Deadlock across locks Circular waiting Consistent acquisition order
    Decide the unavailability policy deliberately: When the lock service is down, does work stop or proceed unprotected? For correctness locking the answer must be stop; for efficiency locking, proceeding is often better.

    Monitoring

    Signals Worth Tracking

    • Acquisition success and failure rate
    • Time spent waiting to acquire
    • Lock hold duration distribution
    • Locks expiring before release
    • Renewal failures and abandonment causes
    • Writes rejected by fencing
    • Contention rate per lock name
    • Locks approaching the renewal cap
    • Lock service latency and availability
    • Clients repeatedly failing to acquire
    A lock that expires before release is a warning that your duration is undersized. It means work regularly outlives the protection it was given.

    Verification

    describe("distributed lock safety", function () {
        it("rejects work from a client that paused past expiry", async function () {
            const clientA = await acquireLock("resource-x", "client-a");
            const staleToken = clientA.fencingToken;
    
            await faultInjector.pauseProcess("client-a", ttlMs * 2);
            const clientB = await acquireLock("resource-x", "client-b");
    
            expect(clientB.fencingToken).toBeGreaterThan(staleToken);
    
            const staleWrite = await resource.apply(payload, {
                fencingToken: staleToken
            });
    
            expect(staleWrite.accepted).toBe(false);
        });
    
        it("does not release a lock held by another client", async function () {
            const clientA = await acquireLock("resource-x", "client-a");
            await sleep(ttlMs + 100);
    
            const clientB = await acquireLock("resource-x", "client-b");
            await releaseLock("resource-x", "client-a");
    
            const stillHeld = await lockService.holder("resource-x");
            expect(stillHeld).toBe("client-b");
        });
    
        it("maintains at most one holder under chaos", async function () {
            const timeline = await runChaosScenario({
                durationMs: 180_000,
                clients: 8,
                pauses: true,
                partitions: true,
                serviceFailover: true
            });
    
            for (const moment of timeline) {
                expect(moment.acceptedWriters.length).toBeLessThanOrEqual(1);
            }
        });
    });

    Common Design Mistakes

    Weak Design

    • Assuming a lock guarantees exclusion
    • Acquiring without any expiry
    • Releasing without checking ownership
    • Using a single-node cache for correctness
    • Holding a lock across slow I/O
    • Renewing without a total hold cap
    • Locking globally where per-key would do
    • Reaching for a lock before considering alternatives
    • Never testing pause and failover

    Strong Design

    • Classifies efficiency versus correctness first
    • Pairs correctness locks with fencing
    • Always sets an expiry
    • Releases conditionally and atomically
    • Uses a consensus-backed lock service
    • Keeps critical sections short
    • Caps renewals and total hold time
    • Scopes locks to the narrowest key
    • Prefers transactions and optimistic concurrency

    System Design Interview Discussion

    Question What Your Answer Should Cover
    Do you need a distributed lock? Alternatives considered and rejected
    Efficiency or correctness? Consequence of simultaneous holders
    Can the lock guarantee exclusion? No, pauses require fencing
    What backs the lock service? Durability and failover behaviour
    How long is the lock held? Work and pause percentiles
    What if the service is down? Fail-open or fail-closed policy
    How is contention managed? Lock granularity and hold time
    How is safety verified? Chaos testing with pauses and failover

    Design Checklist

    Production Checklist

    • Exhaust simpler alternatives before locking
    • Classify the lock as efficiency or correctness
    • Pair correctness locks with fencing tokens
    • Validate tokens at the protected resource
    • Use a consensus-backed lock service
    • Always set an expiry on acquisition
    • Size duration above work and pause percentiles
    • Release conditionally, comparing the holder
    • Renew at a fraction of the duration
    • Cap total renewals and hold time
    • Track remaining margin with a monotonic clock
    • Keep slow I/O outside the critical section
    • Scope locks to the narrowest key
    • Define behaviour when the service is unavailable
    • Acquire multiple locks in a consistent order
    • Monitor contention, expiry, and fencing rejections
    • Test pauses, partitions, and service failover

    Knowledge Check

    1

    Why is a distributed lock weaker than a mutex?

    The holder's liveness cannot be determined, communication can fail, and there is a gap between holding the lock and acting on it.

    2

    What separates efficiency from correctness locking?

    Whether simultaneous holders merely waste work or actually corrupt data. Only the latter requires fencing.

    3

    Why can no timeout make a lock safe?

    Process pauses are unbounded, so a holder may resume after any expiry believing it still holds the lock.

    4

    Why must release be conditional?

    If your lock expired and another client acquired it, an unconditional delete releases their lock rather than yours.

    5

    Why cap total hold time?

    Unbounded renewal lets a hung process retain the lock forever, turning a recoverable crash into a permanent block.

    Summary

    A distributed lock is fundamentally weaker than a local mutex. Holder liveness cannot be determined, communication fails routinely, and an arbitrary delay can occur between acquiring the lock and acting on the resource it protects.

    The essential first question is whether the lock provides efficiency or correctness. If simultaneous holders merely duplicate work, a simple implementation suffices. If they corrupt data, no lock protocol is sufficient on its own, because process pauses have no upper bound.

    Correctness requires fencing tokens validated at the resource, which shifts enforcement to the only party that can reliably decide. This reframes the lock as an optimization preventing wasteful contention, with the token providing the actual guarantee.

    Most distributed locks exist because a simpler mechanism was overlooked. Database transactions, optimistic concurrency, partitioned ownership, idempotent operations, and atomic primitives each avoid the coordination dependency, throughput ceiling, and failure modes that a lock introduces.

    Key Takeaway

    A lock expresses intent; only the resource can enforce exclusion. Establish whether you need efficiency or correctness, exhaust the simpler alternatives first, and if a lock is genuinely required, back it with a consensus-backed service, always set an expiry, release conditionally, cap total hold time, and validate fencing tokens where the write actually lands.