Distributed Locks & Leases
The GC Pause That Corrupted Shared Storage
Test your architecture intuition: Pitch a 7-axis solution, survive two aggressive reviewer objections, and inspect the staff-level Teacher Gold Answer.
1. What It Is & Why It Exists
The Core Problem
In distributed architectures without shared physical memory, multiple independent processes concurrently reading and modifying shared resources face severe concurrency hazards:
- Financial Settlement Collisions: Two payment workers concurrently pick up the same uncaptured invoice, billing the customer twice.
- Inventory Double-Allocation: Two customer sessions purchase the last available seat on a flight simultaneously.
- Split-Brain Leader Execution: In an active-passive cluster (e.g., scheduled cron runners, database failover controllers), two nodes simultaneously believe they are the active leader, executing duplicate batch migrations and corrupting downstream state.
The Breakdown: The Unfenced Lock & The GC Pause Cliff
The Fatal Breakdown of Naive Distributed Locks
A naive lock implementation (e.g., Redis SET resource_id client_uuid NX EX 10) is fundamentally unsafe for transactional correctness across asynchronous networks without a cryptographic or monotonic fencing mechanism:
- Client 1 acquires a 10-second lock on
resource_42. - Client 1 experiences an unexpected 15-second "stop-the-world" Garbage Collection (GC) pause, hypervisor CPU throttle, or network partition.
- While Client 1 is frozen, its 10-second lease expires in the lock coordinator.
- Client 2 acquires the lock on
resource_42and commits an update to the database. - Client 1 resumes execution from its GC pause, completely unaware that its lease expired, and executes its pending database update—silently overwriting and corrupting Client 2's data.
Concrete Proof: Naive Lock vs. Lease vs. Fencing Token
Consider two workers updating a database record under high network jitter and GC pauses:
| Operational Dimension | Naive Mutex (No Expiration) | Time-Bounded Lease (No Fencing Token) | Lease with Monotonic Fencing Token |
|---|---|---|---|
| Crash Behavior | 🚨 Permanent Deadlock: If lock holder crashes, resource is locked forever. | 🟢 Deadlock-Free: Lock automatically expires after lease TTL. | 🟢 Deadlock-Free: Lock automatically expires after lease TTL. |
| GC Pause Safety | 🚨 Unsafe: Worker wakes up after TTL and writes stale data. | 🚨 Unsafe: Client 1 overwrites Client 2 during split-brain pause. | 🛡️ Safe where the token is checked: The downstream database rejects the stale write (); resources that do not check the token are not protected. |
| Network Asynchrony | ❌ Fails under delayed packets | ❌ Fails under delayed packets | 🛡️ A delayed write carrying an old token is rejected at any resource that checks it |
| Consistency Class | Best-effort / Soft lock | AP (High availability, weak safety) | Strong safety, but only for resources that check the token |
The First-Principles Solution: Leases with Fencing Tokens
To guarantee mutual exclusion across asynchronous distributed networks, two primitives must combine:
- The Distributed Lease: A time-bounded distributed lock with a hard Time-To-Live (TTL) and an active heartbeat renewal loop, preventing deadlocks when workers crash.
- Monotonically Increasing Fencing Tokens: Every time a lease is granted, the lock server issues a strictly increasing sequence number (). Downstream storage systems reject any write request carrying a token lower than the highest token previously committed; a write with an equal token comes from the current holder and is accepted, because one holder writes many times with the same token.
Synthesizing vector architecture diagram...
Both panels follow the same timeline: the lock holder freezes for 15 s in a GC pause, but its lease lasts only 10 s. In the "Broken Baseline" panel, the lease expires during the pause, client 2 acquires the lock and writes B, then client 1 wakes up, still believing it holds the lock, and writes its stale A over B. In the "Production Standard" panel, each grant carries a higher fencing token: worker 1 had 101, worker 2 gets 102 and commits with it, so the database records 102 as the highest token seen and rejects worker 1's late write with 101 as stale (409). A lease alone cannot stop a paused client from acting after it expires; only the storage checking the token can, so the check must live where the write happens.
2. Core Mechanics & Algorithmic Architecture
Formal Mathematical Formulations
1. Fencing Token Validation Invariant
Let be the fencing token attached to a storage mutation, and let be the highest token recorded in the target database row:
2. Lease Duration & Drift Safety Formula
To ensure that a worker never executes inside its critical section after a lease has expired on the coordinator, the client’s local execution window must account for the difference in clock rate between the client and the coordinator () and network round-trip latency (). Each side times the lease on its own monotonic clock, so the offset between the two clocks does not matter here; it matters only when one machine writes a wall-clock deadline that another machine reads, and then the margin must also exceed the worst offset between the two clocks (the sections Time: whose clock decides? and Conditional writes are the check of Leases, Fencing Tokens & Distributed Locks):
Where:
- is an assumed bound on the rate difference between the two clocks (for example 1,000 ppm, one clock 500 ppm fast and the other 500 ppm slow: 10 ms over a 10 s lease). While a time daemon slews a clock the rate difference can be far larger (up to one twelfth at chrony's default maximum slew rate), so choose a margin well above the drift figure.
- is the network latency of the lease acquisition RPC. It is needed only if the client starts its timer when the grant arrives; a client that starts its timer when it sends the request is already on the safe side.
3. Heartbeat Renewal Condition
To maintain lock ownership continuously without false expiration:
This gives the holder two renewal attempts (at and ) before the lease ends, so one attempt can fail to a transient network blip and the next still renews the lease.
Step-by-Step Deterministic Lifecycle Pipeline
Synthesizing vector architecture diagram...
In-Memory & Distributed State Structures
| Lock Coordinator | Underlying Storage Representation | Ownership Tracking | Fencing Mechanism |
|---|---|---|---|
| Amazon DynamoDB Lock Client | Item in dedicated DynamoDB table | ownerName (string GUID) | None built in: recordVersionNumber is a random GUID rotated on every heartbeat (it is the CAS value for renewals, not a counter). A fencing token needs a separate ADD fence :1 on the lock item, done in the same conditional UpdateItem as the acquire |
| etcd (v3) | Raft Key-Value with 64-bit Lease ID | 64-bit LeaseID with KeepAlive stream | The lock key's create revision (from etcd's global 64-bit revision counter; the lock recipe orders waiters by it). For a lock key that is never updated this equals its mod_revision |
| Apache ZooKeeper | Ephemeral Sequential znodes (/locks/res-0000000102) | Node path ownership | ZNode sequence number (-order) |
| Redis (Single Instance) | Key with random UUID string | ARGV[1] UUID verification in Lua | Must be maintained separately in an atomic integer counter (INCR); with asynchronous replication a failover can lose the last increment and hand out the same token twice |
Scenario Execution Matrix: Distributed Leases & Fencing Tokens
Under initial database row state: { id: 99, status: 'PENDING', fencing_token: 100 }:
| Step # | Event / Input | In-Memory / Distributed State | Evaluation & Transition | Outcome / Safety Enforcement |
|---|---|---|---|---|
| 1 | Worker A acquires lease (); executes work in | Row state: {id: 99, status: 'PENDING', fencing_token: 100} | UPDATE ... WHERE fencing_token <= 101PASS | Commit Succeeded (HTTP 200) Row state updated to {status: 'SETTLED', fencing_token: 101}; lease released |
| 2 | Worker B acquires lease (); hits 20s stop-the-world JVM GC | Lock coordinator detects heartbeat timeout; expires lease and issues to Worker C | Worker C completes in , commits : Row state becomes {status: 'CONFIRMED', fencing_token: 103} | Worker C commits successfully; fencing ceiling advances to 103 |
| 3 | Worker B wakes up from GC pause; attempts to commit | Row state: {status: 'CONFIRMED', fencing_token: 103} | UPDATE ... WHERE fencing_token <= 102REJECT | Stale Mutation Aborted (HTTP 409) 0 rows affected; Worker B write discarded, zero data corruption |
| 4 | Worker D acquires lock (TTL 5s); hangs for 7s; Worker E acquires lock | Redis lock key holds Worker E UUID; Worker D attempts lock release | Lua script checks: if redis.call('get', KEYS[1]) == ARGV[1]owner_D != owner_E mismatch | Release Aborted (Return 0) Worker E lock untouched; prevents accidental lock revocation |
3. Data Migration, Anti-Entropy & Consistency Protocols
Consensus-Backed Replicated State Machines (etcd / ZooKeeper)
Distributed lock coordinators must guarantee Linearizable Consistency (CP). If a coordinator loses network connectivity or experiences a partition, it must reject writes rather than risk dual leaders:
- Quorum Consensus: In etcd (Raft) and ZooKeeper (Zab), granting a lease, creating or deleting a lock key and advancing the revision counter require acknowledgment from a strict majority (). Renewals are cheaper: etcd's leader resets a lease's deadline in its own memory, and a newly elected leader restarts every lease at its full TTL, so a failover can make a lease longer but never shorter (the section When the lock service itself fails of Leases, Fencing Tokens & Distributed Locks).
- Leader Step-Down: etcd runs Raft with CheckQuorum on, so a leader that has not heard from a majority for an election timeout steps down, and a leader cut off in a minority stops acting as leader. It could not grant leases there anyway, because a grant needs a majority commit.
Synthesizing vector architecture diagram...
Start from a five-node lock cluster that a network partition has just split into two groups. In the "Majority Partition" panel, nodes 1-3 can still reach each other, so the leader collects 3 of 5 acknowledgements, which is a quorum, and keeps granting and renewing leases. In the "Minority Partition" panel, nodes 4 and 5 can only reach 2 of 5, below the quorum of 3, so they refuse to grant anything. Because any two majorities of 5 must share at least one node, the two sides can never both grant the same lease; the minority stays unavailable until the partition heals, which is the price of safety.
The Redis Failover Anomaly (Why Redis is an AP Lock)
In standard Redis Sentinel or Redis Cluster setups, replication between Primary and Replicas is asynchronous:
- Client 1 acquires lock on Primary node (
SET lock_key uuid NX EX 10). Primary returnsOK. - Before the write replicates to the Replica, the Primary crashes or loses power.
- Sentinel promotes the Replica to Primary.
- The new Primary has no record of
lock_key. - Client 2 requests the lock on the new Primary and is granted it!
- Result: Both Client 1 and Client 2 hold the lock simultaneously.
- Production Takeaway: Standard Redis is suitable for optimistic advisory locks or deduplication, but unacceptable for transactional financial safety unless paired with fencing tokens checked at the database. Those tokens must come from a store that commits each increment by majority or synchronously before answering (etcd, ZooKeeper, a strongly consistent database such as a single-Region DynamoDB table), not from an
INCRon the same asynchronously replicated Redis, which the same failover can roll back so that two holders get the same token.
3. Graceful Draining, Cooperative Cancellation & Lease Abandonment
- Cooperative Cancellation Token Architecture: When a lock-holding application container receives
SIGTERMor its background renewal daemon thread detects consecutive heartbeat failures (), it trips an in-processCancellationToken. Critical section threads must evaluate this token prior to issuing mutations to downstream databases, preempting zombie execution before network dispatch. - Graceful Lease Yield vs. Automatic Expiration: On clean container shutdown, the client issues an atomic compare-and-delete release (
EVALSHAmatching its owner GUID). If a worker process crashes unannounced, mutual exclusion does not hang: the coordinator lease automatically expires after , and downstream database fencing tokens (accept only ) reject any delayed in-flight writes from the crashed worker once the next holder has raised the fence. That is why the next holder's first act must be one conditional write that raises the fence, before it reads or writes anything else.
Unlock Complete Architecture & Production Runbooks
You have explored the free architectural preview (~43%). Spend 1 Coin to unlock the remaining 7 production deep-dive sections for a full 24 hours.