Leases, Fencing Tokens & Distributed Locks
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.
Part 0. Start here
The problem: a job that must run exactly once
A nightly billing job must run on exactly one worker. It writes the day's 50,000 invoices in 10 chunks of 5,000. Worker A holds a 10-second lease on the job (a lock that expires unless it is renewed) and renews it every 3 seconds. These are our example's numbers; they are also exactly the settings in the example of the AWS DynamoDB Lock Client's README (withLeaseDuration(10), withHeartbeatPeriod(3)).
At second 4, A stops for 15 seconds. It might be a stop-the-world garbage-collection pause, or its VM being paused for a live migration; the length is our assumption. Its renewals stop too, because the thread that sends them lives in the same frozen process. At second 13 the lease runs out and worker B takes the job. At second 19, A wakes up with chunk 2 already built in its buffers, still believing it owns the job, and sends it. The billing store charges each invoice it records, so if that write lands, 5,000 customers are billed twice.
A's lease has expired, but A doesn't know. What, anywhere in the system, can stop A's writes from landing on top of B's?
The big picture
Synthesizing vector architecture diagram...
What to notice: A never learns it lost the lease until the store tells it. The lock service decided who should work; the store's check is what kept the data correct.
What you'll be able to do after this page
- Explain the difference between a lock, a lease and a fencing token, and the three rules that tie them together (Part 1).
- Walk through acquiring, renewing and releasing a lease, and say which state lives in the log and which lives only in the leader's memory (Part 2).
- Say which clock times a lease, why offset between machines usually doesn't matter but rate does, and what a margin can and can't cover (Part 3).
- Trace a paused holder being stopped by a fencing token, including why the new holder raises the fence before it reads (Part 4).
- Write the fence as a conditional write in DynamoDB, SQL and S3, and build a lease on a DynamoDB item without trusting anyone's clock (Part 5).
- Explain how a crashed holder's claim is found and cleaned up, and compute the earliest safe takeover (Part 6).
- Treat leader election and single-writer failover as leases with epochs (Part 7).
- Explain what happens when the lock service itself fails over, why an async Redis failover is different, and what the Redlock argument was about (Part 8).
- Decide when you don't need a lock at all (Part 9).
- Trace an acquire, a renewal and a fenced write from your code through the network, the lock service and its log to the store (Part 10).
- Map all of this to AWS services, and name the look-alikes that are not fences (Part 11).
You may have arrived from a step that relies on this: step 2.2 of the distributed unique ID generator loop (a worker froze for 40 seconds), step 2.4 of the job scheduler loop (a paused worker finished a job someone else re-ran), step 2.5 of the message queue loop (a stale worker deleted a message another worker holds), or step 2.4 of the ride-sharing loop (a late accept after the offer expired). This page is the "why" behind all four.
Part 1. The core mental model
Three words get used loosely in design reviews: lock, lease and fencing token. They mean different things, and most bugs in this area come from treating one as another. This Part pins them down, then introduces the example we follow for the rest of the page.
Lock, lease, fence
| Word | What it is | What it guarantees | What it does not guarantee |
|---|---|---|---|
| Lock | A named thing only one holder may have at a time | While the lock service is healthy, it grants the name to one holder at a time | That a holder that crashed ever lets go; that a holder knows when it has lost it |
| Lease | A lock that expires unless the holder renews it | A crashed or cut-off holder blocks everyone for at most one lease length | That the holder stops acting when the lease ends: a paused holder doesn't know it lost it |
| Fencing token | A number handed out with each grant that only ever grows | A resource that checks it refuses writes from any older holder | Anything for a resource that doesn't check it |
Synthesizing vector architecture diagram...
What to notice: the check sits at the resource, not at the worker or the lock service. The lock service only hands out numbers; the resource is where an old holder is actually stopped.
Three rules the design never breaks
| Rule | What breaks without it |
|---|---|
| Every grant carries a token from a counter that only grows (never reused, even after a failover or a restore) | A new holder can get the same number as an old one, and the resource can't tell them apart |
| The resource checks the token on every write: reject lower; equal is the current holder | A paused holder that wakes up overwrites the new holder's work |
| A holder stops acting before the lock service would give the lease away, timed on its own monotonic clock | Two holders act at the same time even without a pause, and the fence becomes the only line of defence all the time instead of in rare cases |
Not every "lease" is this lease
Two loops use the word for something else. The rate limiter loop (step 2.3) leases blocks of quota to gateways: a time-bounded grant of capacity, not of exclusive ownership, so nobody needs to be fenced. The URL shortener loop (step 1.2) "leases" blocks of counter values with an atomic add and no expiry at all. Neither needs this page's machinery.
The example we follow
Real systems hold thousands of leases, which is too many to watch. So one small story runs through the whole page: the billing job changes hands. Every trace on this page was produced by running a private reference simulator of this setup (true time, per-machine clocks, a 3-node lock service, the store), not worked out by hand.
| Setting | Our example | At real scale |
|---|---|---|
| Lock service | 3 nodes n1, n2, n3, one per Availability Zone, replicating through a consensus log (like etcd); n1 is the leader at the start | etcd or ZooKeeper with 3 or 5 nodes, or a DynamoDB lock table |
| Lock key | job/billing-0928 | one key per resource |
| Lease length (TTL) | 10 s, timed by the leader's in-memory lease tracker on its monotonic clock | the DynamoDB Lock Client README's example uses 10 s; Kubernetes core clients use 15 s; a ZooKeeper session lasts 2 to 20 times tickTime (2,000 ms in the sample config, so 4 to 40 s) |
| Renewal interval | every 3 s, about a third of the lease, so two renewal attempts can fail before the holder stops | Lock Client README 3 s; Kubernetes core clients retry every 2 s |
| Holder's safety margin | 2 s, our choice: the holder stops acting 8 s after it sent its last successful renewal | a design choice, not a rule |
| Fencing token | the lock key's create revision: the lock service's global revision counter at the moment the holder's key was created | etcd: the lock key's revision; ZooKeeper: the sequential node number or the zxid that created it |
| Revision counter at the start | 6. Other keys change in between, so tokens grow but skip: A 7, B 12, C 15 | a 64-bit counter |
| The resource | the billing store's row billing:2026-09-28: fence (highest token seen, starts at 0), state, owner_token, claimed_since (present only while claimed), chunks_done; every write is conditional | a DynamoDB item or a SQL row |
| Workers | A and B are the two instances started for tonight's run: whoever holds the lock works, the other retries every 3 s. C is a pool worker that only picks up jobs that are re-queued | |
| Raft election timeout | 1 s, randomized between 1 and 2 s per election | etcd's default is 1,000 ms, with a 100 ms heartbeat |
| Sweeper | runs every 10 s; re-queues an orphan only on a sweep at least 5 s after it first saw it | every 5 to 30 s in the loops |
| Clocks | A and C correct. B's wall clock is stepped back 3 s at 16 s (a VM restored from a snapshot, an operator, or a time daemon set to step) | EC2 instances synced to the Amazon Time Sync Service: microsecond-level accuracy on supported instances |
The story in five beats: A holds the job (Part 2), A freezes (Parts 3 and 4), B takes over (Part 4), B dies holding it (Part 6), and the lock service loses its leader while C holds it (Part 8). Each Part shows only its own slice of events; the full 24-event table is in Part 12.
What to remember from Part 1
- A lease is a lock with an expiry, so a crash can't block everyone forever.
- Expiry means a holder can lose the lock without knowing, so a lease alone never protects data.
- The resource checks a token that only grows: reject lower; equal is the current holder.
Part 2. Acquire, renew, release: the happy path
A crashed holder of a plain lock blocks everyone forever. The obvious fix, "a lock plus a timeout in the client", doesn't work: the timeout lives in the holder, so when the holder dies nobody else knows the timeout passed, and when the holder pauses its timeout pauses too. The expiry must be known to the grantor. This Part walks through the three operations that make that work.
Acquire
textACQUIRE(key K, lease length T) 1. ask the lock service for a new lease with TTL T 2. in one transaction: create K bound to that lease, only if K doesn't exist 3. the service commits both through its log (a majority of nodes has them on disk) 4. reply: lease id, and token = the revision at which K was created 5. the leader's lease tracker starts timing T on its monotonic clock, from when it applied the grant 6. the holder records sent_at (its monotonic clock, read before step 1) and its own stop time: sent_at + T - margin
For A at t = 0.0: the grant and the key creation commit at revision 7; the leader's tracker sets the deadline to 0.0 + 10 = 10.0; A's own stop time is 0.0 + 10 − 2 = 8.0.
What the lock service keeps, and where
| State | Where it lives | Why |
|---|---|---|
| The key, the lease id it's bound to, the lease's TTL | In the replicated log (majority on disk) | A new leader must know who holds what |
| The key's create revision (the token) and the global revision counter | In the replicated log | The counter must never go backwards, whichever node leads |
| Grants, releases and expiry deletes | In the replicated log | Each one changes who holds the lock |
| The lease's current deadline | Only in the leader's memory | Renewals would otherwise cost a majority disk write every few seconds per holder; etcd serves keep-alives from the leader's memory and does not log them |
The last row is the surprise. It is safe only because of what a new leader does after a failover (Part 8).
Renew and the holder's stop time
textRENEW LOOP (holder) -- runs every T/3 1. sent_at = monotonic now -- read BEFORE sending 2. send keep-alive(lease id) to the leader 3. leader: reset deadline = its now + T, in memory only; reply 4. on success: stop_at = sent_at + T - margin 5. on failure: retry (new leader if needed) until stop_at; then give up the lease locally BEFORE EVERY ACTION (holder) if monotonic now >= stop_at: stop, don't act
Synthesizing vector architecture diagram...
What to notice: the acquire goes through the log to a majority; the keep-alive stops at the leader. The holder's stop time is always 2 s (the margin) before the service's deadline.
The first events of the story, as the simulator printed them:
| # | t (s) | Event | Service deadline (leader memory) | A's stop time |
|---|---|---|---|---|
| 1 | 0.0 | A acquires: key created at revision 7, token 7 | 10.0 | 0.0 + 10 − 2 = 8.0 |
| 2 | 0.5 | A raises the store's fence 0 → 7 and marks the row RUNNING (Part 4 explains why this comes first) | 10.0 | 8.0 |
| 3 | 1.0 | A writes chunk 1 with token 7: accepted | 10.0 | 8.0 |
| 4 | 3.0 | A renews (sent 3.0, answered 3.0) | 3.0 + 10 = 13.0 | 3.0 + 10 − 2 = 11.0 |
Why does A time its lease from when it sent the renewal, not from when the answer came back?
Release
When the job is done, the holder revokes its lease. The key goes with it at once, so the next holder doesn't have to wait out the TTL. A release must be conditional on still owning the lock: revoking your own lease id only ever deletes keys bound to that lease, so a holder that lost the lock can't delete the next holder's key. The same rule applies to every lock store: Redis's documented release deletes the key only if it still holds the holder's random value; a plain delete is unsafe because a slow client can remove a lock someone else now holds.
textRELEASE(lease id) 1. finish the work, including the store's final write (with the token) 2. revoke(lease id) -- deletes only keys bound to this lease 3. the service commits the delete; the revision goes up
In the story, C releases at t = 52.0 and the key is deleted at once at revision 16 (Part 8).
Snapshot S1: A holds the job
Synthesizing vector architecture diagram...
After event 4: A holds token 7, the store's fence is 7, and chunk 1 is done.
What to remember from Part 2
- Renew at a third of the lease, so two renewal attempts can fail before the holder stops.
- The holder's stop time starts when it sent the renewal, and comes before the service's deadline.
- Release is conditional: only the current owner can release.
Part 3. Time: whose clock decides?
A lease is a promise about time, and every machine has its own clock. Two questions decide whether the promise holds: which kind of clock times the lease, and whose clock is used when one machine judges another's lease.
Two clocks in every machine
| Clock | What it is | Can jump? | Use it for |
|---|---|---|---|
Wall clock (Linux CLOCK_REALTIME) | The time of day; settable | Yes: an operator, a VM restore or a time daemon can step it forwards or backwards | Timestamps humans read; a deadline written into a shared record |
Monotonic clock (Linux CLOCK_MONOTONIC) | Time since some fixed point; not settable | No jumps. It is still sped up or slowed down by the time daemon's frequency corrections, and it does not count time the machine is suspended (CLOCK_BOOTTIME does) | Every duration: renewal timers, stop times, timeouts |
Two different clock errors
| Error | What it means | Our numbers |
|---|---|---|
| Rate (drift, and slewing) | How fast a clock runs compared with true time | Two monotonic clocks each timing the same 10 s, one 500 parts per million fast and one 500 ppm slow (a generous worst-case assumption), disagree by 10 s × 1,000 ppm = 10 ms. While a time daemon slews (corrects an offset by running the clock fast or slow), the monotonic clock is slewed too; at chrony's default maximum slew rate of one twelfth (83,333 ppm), that is up to 10 s ÷ 12 ≈ 0.83 s over 10 s. Only CLOCK_MONOTONIC_RAW is not slewed |
| Offset (skew, and steps) | How far apart two clocks read at the same moment | A step changes one machine's wall clock by any amount at once (B's 3 s below) |
In the lock-service design, the holder times its lease from its send and the leader times it from its apply, each on its own monotonic clock. Neither ever reads the other's clock, so offset doesn't matter at all; only rate does. Offset matters only when one machine's wall clock writes a deadline and another machine's clock reads it, which is exactly what a lease stored in a database does (Part 5). Then the margin must be bigger than the worst offset.
What the margin covers
| The margin covers | The margin can't cover |
|---|---|
| Rate differences between the holder's and the leader's clocks (10 ms here, up to 0.83 s during a maximum-rate slew) | A pause of any length: a paused process doesn't run its check at all |
| The slop between the holder's "am I still inside my lease?" check and the act it guards, when nothing pauses | A network delay after the write left the holder (the write can arrive any time later) |
The stop time is our first formula:
For A's renewal sent at 3.0: 3.0 + 10 − 2 = 11.0. For B's acquire sent at 13.5: 13.5 + 10 − 2 = 21.5.
The storage's clock must not decide
A tempting shortcut is to let the database judge expiry: "take over if lease_until < now()". In DynamoDB there is no now(): condition expressions offer comparisons and the functions attribute_exists, attribute_not_exists, attribute_type, begins_with, contains and size, and nothing that reads the current time. So :now is always some client's clock, and Part 5 deals with that. A single SQL database's now() is one clock, so using it for rows in that same database is fine: every judgement uses the same clock. It doesn't extend to anything outside that database.
Trace: B's wall clock steps back 3 s
B acquired at 13.5 (token 12, service deadline 23.5, stop time 21.5). At t = 16.0, B's wall clock is stepped back 3 s.
B's wall clock just jumped back 3 s. Which of B's decisions change?
Synthesizing vector architecture diagram...
What to notice: the straight line is B's monotonic clock and the line that drops at 16 is its wall clock; the flat line is the stop reading 21.5. The monotonic clock crosses it at true 21.5, the wall clock only at true 24.5, one second after the lock service's deadline of 23.5.
What to remember from Part 3
- Time a lease on a monotonic clock; wall clocks jump.
- When each machine times its own side, only the clocks' rates matter; offsets matter only when one machine's clock writes a deadline another reads, and then the margin must exceed the worst offset.
- A margin covers drift and small delays; nothing covers a pause except the resource's check.
Part 4. The pause, the zombie and the fencing token
Part 3 ended with the one thing no clock can catch: a pause between a check and an act. This Part shows the pause happening, the old holder coming back as a zombie (a process acting on a lock it no longer holds), and the fencing token stopping it.
Trace: A pauses, B takes over
Synthesizing vector architecture diagram...
What to notice: A's only check (at 3.9) passed, and it was true when A made it. Everything that saved the data happened at the store: B raised the fence at 13.6, so A's write at 19.05 met a fence of 12.
| # | t (s) | Event | Lock service | A | B | Row fence |
|---|---|---|---|---|---|---|
| 5 | 3.9 | A checks 3.9 < 11.0 and builds the chunk 2 write with token 7 | A holds, deadline 13.0 | holds | retrying | 7 |
| 6 | 4.0 | A's pause starts; its 6.0 and 9.0 renewals never happen (unrelated keys take revisions 8, 9, 10) | A holds, deadline 13.0 | frozen | retrying | 7 |
| 7 | 13.0 | n1's lease tracker finds L1 expired and commits a revoke; A's key is deleted at revision 11 | free | frozen | retrying | 7 |
| 8 | 13.5 | B acquires: key created at revision 12, token 12; deadline 23.5, B's stop time 21.5 | B holds | frozen | holds | 7 |
| 9 | 13.6 | B raises the fence 7 → 12, sets owner_token = 12, claimed_since = 13.6, and reads the old row: chunk 1 done | B holds | frozen | holds | 12 |
| 10 | 14.0 | B writes chunk 2 with token 12: accepted | B holds | frozen | holds | 12 |
| 11 | 16.0 | B's wall clock steps back 3 s; B renews at 16.5 (deadline 26.5) and 19.5 (deadline 29.5) on its monotonic clock | B holds | frozen | holds | 12 |
| 12 | 19.0 | A wakes. Its chunk 2 (token 7) reaches the store at 19.05: rejected, 7 is lower than 12. Its keep-alive at 19.1 gets "lease not found" (its own stop time of 11.0 is long past too); A drops the job | B holds | gone | holds | 12 |
Where the token is checked
At the resource, on every write and on the release of anything the resource holds for the job. The holder's own "am I still inside my lease?" check is an optimization: it stops a healthy holder from sending writes that would be refused anyway, and keeps overlap rare. It can never be the safety mechanism, because a pause can land right after it.
textFENCED WRITE (at the store, one atomic conditional write) accept only if fence <= :token -- reject lower then fence = :token, apply the data change a rejected write means: you have lost the lock; stop and discard local state
Why equal is accepted: the current holder writes many times with the same token (A wrote with 7 at 0.5 and 1.0; B wrote with 12 at 13.6, 14.0 and 20.0), so "equal" means "the current holder". The gold answer in the drill The GC Pause That Corrupted Shared Storage rejects token <= last_seen; that is right only when each token makes exactly one write.
Raise the fence first
B's first act, at 13.6, was one conditional write that raised the fence to 12 and returned the old row, before B read anything or did any work. That order matters.
Suppose B skips that first step: it reads the row at 13.6 without raising the fence, and its first real write (chunk 2) goes out at 20.0. A's chunk 2 arrives at 19.05 with token 7, and the store's fence is still 7. Is A's write safe?
Snapshots S2 and S3: the zombie is stopped
Synthesizing vector architecture diagram...
S2, after event 10: A is frozen with a write it still believes in; B holds token 12 and has already raised the fence.
Synthesizing vector architecture diagram...
S3, after event 12: A's write was refused and A let go. The row is exactly as B left it.
The same fence on S3
The drill's resource is S3, which has no "greater than" condition: a conditional PutObject can only say "only if this key doesn't exist" (If-None-Match: *) or "only if the object's ETag is still the one I read" (If-Match). The fence still works with two changes: each holder writes its output under a key that includes its token, and it commits by rewriting one small manifest object with If-Match on the ETag it last saw. Raising the fence first means the new holder rewrites the manifest with its token before it does anything else. The simulator's trace (ETags labelled E0, E1, … in order):
| t (s) | Who | Request | Result |
|---|---|---|---|
| 0.5 | A | PUT manifest If-Match: E0, body "fence 7" | 200, ETag E1 |
| 1.0 | A | PUT render/billing/t7/chunk-1 | 200 |
| 1.0 | A | PUT manifest If-Match: E1, body "fence 7, chunks: t7/chunk-1" | 200, ETag E2; A keeps E2 for its next commit |
| 13.6 | B | GET manifest | ETag E2, "fence 7" |
| 13.6 | B | PUT manifest If-Match: E2, body "fence 12, chunks: t7/chunk-1" (the raise) | 200, ETag E3 |
| 14.0 | B | PUT render/billing/t12/chunk-2, then PUT manifest If-Match: E3 | 200, ETag E4 |
| 19.05 | A | PUT render/billing/t7/chunk-2 | 200, but harmless: no manifest will ever name this key |
| 19.05 | A | PUT manifest If-Match: E2 | 412 Precondition Failed (the manifest is at E4) |
Readers trust only what the manifest names, so A's stray object changes nothing. S3 gives compare-and-set on "the version I read", not a numeric fence; the token in the manifest body is what tells a reader (and an auditor) which holder wrote each chunk. When several conditional writes race on the same key, S3 lets the first to finish succeed and fails the rest with 412.
The same rule in other systems
The stock exchange loop (step 1.4) does the same thing to a journal: the spare engine fences the 2-of-3 journal with epoch e+1 before it replays, so the old primary can't append anything the spare didn't see.
What fencing does not protect
Fencing protects exactly the resources that check the token, and nothing else. An email provider, a payment provider's API or a partner's webhook will happily accept A's late call. For those, the only safety net is an idempotency key stored with the side effect, so a duplicate call is recognised and dropped; the Idempotency & Effectively-Once Processing loop primitive covers that. This is also the answer to the first question of the drill The GC Pause That Corrupted Shared Storage: an expiring TTL lock can't protect storage because the paused holder doesn't know its lock expired, so the storage must refuse its writes by token.
What to remember from Part 4
- A paused holder wakes up believing it still holds the lock; only the resource can stop it.
- Reject lower; equal is the current holder.
- A new holder's first act is to raise the fence, before it reads or writes.
Part 5. Conditional writes are the check
A fence is useless if the check and the write are two steps: a zombie could pass the check, pause, and write after the new holder raised the fence. The check and the write must be one atomic conditional write in the store. This Part shows it in DynamoDB, SQL and S3, then builds the loops' favourite pattern: a lease that lives in a DynamoDB item, with no lock service at all.
DynamoDB
When the fence and the data are in one item, a single UpdateItem or PutItem with a ConditionExpression does it, and a failed condition returns ConditionalCheckFailedException. Our invoices are separate items, so each batch goes in a TransactWriteItems with a ConditionCheck on the fence item plus the invoice puts. A transaction holds at most 100 actions, and no two may touch the same item, so each carries 1 check + 99 puts:
- one chunk of 5,000 invoices: 5,000 ÷ 99 = 50.5, so 51 transactions;
- the whole run of 10 chunks: 10 × 51 = 510 transactions.
json{ "TransactItems": [ { "ConditionCheck": { "TableName": "billing", "Key": { "pk": { "S": "billing:2026-09-28" } }, "ConditionExpression": "fence <= :t", "ExpressionAttributeValues": { ":t": { "N": "7" } } } }, { "Put": { "TableName": "invoices", "Item": { "pk": { "S": "inv#2026-09-28#05001" }, "amount": { "N": "4200" } } } } ] }
A's request at 19.05 (token 7 against fence 12) fails as a whole, with no invoice written:
json{ "__type": "TransactionCanceledException", "CancellationReasons": [ { "Code": "ConditionalCheckFailed", "Message": "The conditional request failed" }, { "Code": "None" } ] }
The raise itself is one UpdateItem on the fence item: SET fence = :t, #s = :running, owner_token = :t, claimed_since = :now with ConditionExpression: fence <= :t, asking for the old values back (ReturnValues: ALL_OLD).
SQL
In one transaction, lock the fence row first, then write the data:
sqlBEGIN; UPDATE billing_fence SET fence = :t WHERE job_id = 'billing:2026-09-28' AND fence <= :t; -- 1 row: go on; 0 rows: we lost, ROLLBACK INSERT INTO invoices (...) VALUES (...); -- the chunk COMMIT;
The UPDATE takes the row lock, so a concurrent raise waits until we commit and then re-checks its WHERE against the committed fence (at READ COMMITTED; at stricter isolation the second raise fails with a serialization error and retries). A's transaction at 19.05 updates 0 rows and rolls back.
S3
If-None-Match: * creates an object only if the key doesn't exist; If-Match: <ETag> replaces it only if it is still the version we read. A's manifest commit at 19.05:
httpPUT /render/billing/manifest.json HTTP/1.1 Host: example-bucket.s3.amazonaws.com If-Match: "E2" HTTP/1.1 412 Precondition Failed
What each store can compare
| Store | The check | What it can compare | Failure signal |
|---|---|---|---|
| DynamoDB | ConditionExpression, alone or as a ConditionCheck in a transaction | Numbers, strings, existence: a true numeric fence | ConditionalCheckFailedException; in a transaction, TransactionCanceledException with reason ConditionalCheckFailed |
| SQL | UPDATE … WHERE fence <= :t in the same transaction as the data | Anything SQL can express: a true numeric fence | 0 rows updated |
| S3 | If-Match / If-None-Match on PutObject, CopyObject, CompleteMultipartUpload | Only "absent" or "same ETag as I read": compare-and-set, no "greater than" | 412 Precondition Failed; 409 or 404 when a delete races the write |
Shared limit: every fence raise and every ConditionCheck (and, when the lease is kept on the same item, every claim and renewal) touches the one fence item. A DynamoDB partition gives at most 1,000 write units a second, and one item always lives in one partition. Our 510 transactions a night are nothing, but a fence item shared by a busy stream of writers can become the hot key.
A lease built on a DynamoDB item
Many loops skip the lock service and keep the lease itself in a DynamoDB item. The item:
| Attribute | Meaning |
|---|---|
lock_id | Partition key: the resource, e.g. billing:2026-09-28 |
holder | Who holds it; removed on release |
epoch | The fencing token: a counter that only grows, never deleted |
lease_until | The holder's wall-clock time plus the lease, in ms: written by one clock, read by others |
rvn | A new random value on every renewal (the record version number) |
Claim, renew and release, as UpdateItem requests:
json{ "Claim": { "ConditionExpression": "attribute_not_exists(holder) OR lease_until < :claimant_now_minus_margin", "UpdateExpression": "SET holder = :me, lease_until = :my_now_plus_lease, rvn = :new_rvn, epoch = if_not_exists(epoch, :zero) + :one" }, "Renew": { "ConditionExpression": "holder = :me AND epoch = :e", "UpdateExpression": "SET lease_until = :my_now_plus_lease, rvn = :new_rvn" }, "Release": { "ConditionExpression": "holder = :me AND epoch = :e", "UpdateExpression": "REMOVE holder SET lease_until = :zero" } }
Release clears holder; it never deletes the item, because the item carries the epoch. The problem is the claim's condition: :claimant_now_minus_margin comes from the claimant's clock, while lease_until came from the holder's. That is the "one clock writes, another reads" case from Part 3.
The Lock Client's clock-free rule avoids comparing clocks at all: the claimant reads rvn, waits a full lease on its own monotonic clock, and claims only if rvn is unchanged (condition rvn = :seen). The Lock Client never stores absolute times in DynamoDB, only the lease duration. Kubernetes' Lease-based election uses the same idea: a candidate times expiry from when it last observed a change to the lease record.
Trace: D replaces B, and D's clock is 3 s fast
We replay the start of the story on a DynamoDB lock item, with worker D in B's place. A's clock is correct: it acquired at 0.0 (lease_until 10.0, rvn v1) and renewed at 3.0 (lease_until 13.0, rvn v2), then paused at 4.0. A's own stop time is still 3.0 + 10 − 2 = 11.0: that is when A's right to act ends. D's wall clock reads 3 s ahead of true time. D polls the item once a second, at 0.5, 1.5, 2.5 and so on.
D's clock is 3 s fast. With no margin, when does D take over, and is there an overlap with A's right to act? With a 5 s claim margin? With the Lock Client's rule?
Synthesizing vector architecture diagram...
What to notice: only the no-margin bar starts before A's bar ends. That half second is an overlap a pause doesn't even need; the other two rules leave a gap.
TTL is not a deadline, and it resets epochs
DynamoDB's Time to Live looks like a lease that frees itself, and it isn't one:
- TTL deletes expired items typically within a few days, not at the expiry time. Until then the item is still there and can be read and updated, so every reader must filter by the expiry attribute itself.
- Never put TTL on an item that carries the epoch counter. When the item is deleted, the next claim's
if_not_exists(epoch, :zero) + :onestarts again at 1. From the simulator: A claims with epoch 1, D with 2, a release by clearingholderkeeps 2, the next claim gets 3; then TTL deletes the item and the next claim gets epoch 1 again, a token the store has already seen. Keep the counter on an item that is never deleted, and release by clearingholder. - Restores bring old epochs back. A point-in-time restore of the lock table returns every epoch to its old value, so before using a restored table, jump every epoch by a large amount (the ID generator loop, step 3.2, adds an offset of 1,000,000).
- A multi-Region strongly consistent (MRSC) global table has no TTL at all.
What to remember from Part 5
- A conditional write is the fence: the store checks and writes in one step.
- In DynamoDB,
:nowis always a client's clock: add a margin bigger than the worst offset, or wait a full lease on your own clock and check nothing changed. - TTL deletes days later, and deleting the item resets the epoch: keep the counter on an item that is never deleted.
Part 6. Crashes and orphaned locks
A holder that crashes can't release anything. The lock service's key frees itself when the lease runs out, but the store's row doesn't: it still says RUNNING, owned by a dead worker. If nobody is left to retry, the job never finishes. This Part is about finding and cleaning up those orphans (claims whose holder is gone) without ever breaking the fence.
Trace: B dies holding the job
| # | t (s) | Event | Lock service | Billing row |
|---|---|---|---|---|
| 13 | 20.0 | B writes chunk 3 with token 12 (the sweeper at 20.0 sees a live key at revision 12 and skips) | B holds, deadline 29.5 | fence 12, chunks_done 3 |
| 14 | 21.0 | B is killed (out of memory) mid-job. Its last renewal reached n1 at 19.5 | B's key still there | unchanged |
| 15 | 21.0 to 29.5 | The lock is held by a dead process; nobody can take it until the lease runs out | held by nobody alive | unchanged |
| 16 | 29.5 | L2 expires; B's key is deleted at revision 14. The row still says RUNNING, owner 12, and A (which gave up at 19.1) and B (dead) were the only workers started for this job, so nobody will retry it | free | RUNNING, owner 12: an orphan |
| 17 | 30.0 | The sweeper finds the row through the sparse index; there is no live key under job/billing-0928. It records orphan_seen_at = 30.0 and does nothing yet | free | unchanged |
| 18 | 40.0 | The sweeper again, at least 5 s after the first sighting: still orphaned. Conditional re-queue: state = QUEUED, remove claimed_since, only if state = RUNNING AND owner_token = 12 AND fence = 12 | free | QUEUED |
| 19 | 40.2 | C picks up the job and acquires: revision 15, token 15. At 40.3 it raises the fence 12 → 15 and sets owner_token = 15, claimed_since = 40.3; the old row says 3 chunks are done, so C resumes at chunk 4 | C holds, deadline 50.2 | fence 15, RUNNING, owner 15 |
Synthesizing vector architecture diagram...
What to notice: every arrow into RUNNING is a fence raise, and the sweeper's arrow is conditional on the exact owner and fence it saw. If a new holder raised the fence in between, the sweeper's write fails and does nothing.
Two kinds of lock, two kinds of cleanup
| Where the claim lives | Who creates it | Who deletes it after a crash |
|---|---|---|
| The lock service's key | The lock service, bound to a lease | The lock service, when the lease expires: automatic |
The store's claim marker (state = RUNNING, owner_token, claimed_since) | The holder, in its fence raise | Nobody, unless something goes looking for it |
The marker must be findable. So it is written in the same conditional write that raises the fence (events 2, 9 and 19), with an attribute that puts the row in a sparse index: an index that holds only rows that have the attribute. No claim can exist without being in the index, and the sweeper reads the index instead of scanning the table.
| Store | Sparse index | Row enters it | Row leaves it |
|---|---|---|---|
| DynamoDB | A global secondary index keyed on attributes set only while claimed (a small claim_bucket number as partition key, claimed_since as sort key) | The fence raise sets them | DONE or the sweeper's re-queue removes them |
| SQL | A partial index WHERE claimed_since IS NOT NULL | The fence raise sets claimed_since | Same |
The sweeper
textEVERY 10 s: for each row in the sparse index (claimed rows only): ask the lock service: is there a live key for this job? -- no clock comparison if yes: skip if no, and not seen before: remember orphan_seen_at = now, owner and fence seen if no, and seen at least 5 s ago: conditional write: state = QUEUED, remove claimed_since only if state = RUNNING AND owner_token = :seen AND fence = :seen
The sweeper never compares timestamps from other machines: it asks the authority whether a lease is alive. It waits for a second sighting so that a holder caught between acquiring its key and raising the fence isn't swept by mistake. And it is only a liveness mechanism: the fence keeps the data safe even if the sweeper is wrong, because a wrongly re-queued job is taken by a new holder with a higher token, and the old holder's writes are then refused.
How soon can someone take over?
Our second formula:
B's last renewal reached n1 at 19.5, so the earliest safe takeover is 19.5 + 10 = 29.5. Add a claim margin only when the claimant judges expiry on its own clock (Part 5); here the lock service itself deleted the key, so no margin is needed.
B died at t = 21. When is the earliest anyone can safely take over, and why did C actually start at 40.2?
A sweeper that doesn't storm
If our own renewal path fails (the lock service is unreachable for 30 s, or a bad deploy stalls every heartbeat thread), thousands of leases expire together although their holders are healthy. A naive sweeper then re-queues all of them at once, and every job restarts. The job scheduler loop (step 2.5) adds these rules:
| Rule | Why |
|---|---|
| Re-queue at most N orphans per sweep | Restarts spread out instead of arriving as one wave |
| If more than a threshold of claims look orphaned in one sweep, pause re-queuing and alert | Mass expiry usually means our renewal path failed, not that the workers died |
| After our own renewal outage, extend every lease by the outage length before resuming | Healthy holders shouldn't lose their work to our outage |
Snapshot S4: the orphan
Synthesizing vector architecture diagram...
S4, after event 16: the lock is free, but the row still says RUNNING by 12, and nobody is left to retry. Only the sparse index makes this row findable.
What to remember from Part 6
- A crash leaves the lock held until the lease runs out: lease length is how long a dead holder blocks everyone.
- A claim that doesn't free itself must be written with an index attribute in the same write that raises the fence, and released by a conditional write.
- A sweeper restores liveness; the fence keeps the data safe even if the sweeper is wrong.
Part 7. Leader election and single writers
Many loops need exactly one active instance of something: one pacer, one relay, one matching engine, one writing Region. That is a lock on a role instead of a job, and the same machinery applies.
A leader is a lease holder
Whoever holds the lease on the role's key is the leader, and the leader's fencing token is called its epoch (or generation, or term). Every write the leader makes carries its epoch, so anything from an old leader can be refused.
textCAMPAIGN(election prefix P) 1. create key P/<me> bound to my lease -- etcd: a key; ZooKeeper: an ephemeral sequential node 2. the key with the lowest create revision (or sequence number) leads 3. otherwise: watch only the key just before mine 4. when that key is deleted: check again; if mine is now the lowest, I lead 5. my epoch = my key's revision; stamp it on every write
Watching only your predecessor matters: if every candidate watched the leader's key, one deletion would wake all of them at once (a herd), and all would race to read the prefix. With predecessor watches, each deletion wakes exactly one candidate. etcd's own lock recipe works this way: a waiter waits for the deletion of keys with a lower create revision than its own.
An election has no orphan
The main story needed a sweeper because C was not waiting for the job. In an election, the waiters queue. Side trace from the simulator, starting from the moment A's key is gone:
| t (s) | Event |
|---|---|
| 13.5 | B and C both campaign: B's key is created at revision 12, C's at 13. B leads with epoch 12; C watches revision 12 |
| 16.5, 19.5 | B renews; C renews its own lease too |
| 21.0 | B is killed |
| 29.5 | B's lease expires and its key is deleted at revision 14. C's watch fires; C's key is now the lowest, so C leads with epoch 13 at once, with no sweeper |
Synthesizing vector architecture diagram...
What to notice: C watches only B's key. When B's key is deleted at 29.5, exactly one candidate wakes, and it leads with epoch 13, higher than any epoch B used.
etcd, ZooKeeper and Kubernetes compared
| etcd election or lock | ZooKeeper recipe | Kubernetes Lease (client-go) | |
|---|---|---|---|
| The lease | A lease with a TTL, kept alive by the client | The client's session, 2 to 20 × tickTime | A Lease object with a duration: core clients default to 15 s, renew deadline 10 s, retry every 2 s |
| Queueing | Lowest create revision leads; waiters watch their predecessor | Ephemeral sequential nodes (a 10-digit counter kept by the parent); lowest leads; watch the predecessor | None: candidates poll and try to take an expired lease |
| Fencing value | The key's create revision | The node's sequence number, or the zxid of the change that created it (every change gets a zxid that orders it) | None given; its documentation says it "does not guarantee that only one client is acting as a leader (a.k.a. fencing)" |
| Who must check it | Your resource | Your resource | Your resource, with an epoch you add yourself |
| Clock used | Leader's monotonic clock for the deadline; client's for its stop time | The ensemble's session timeout | The candidate's own clock, timed from the last observed change to the Lease record; tolerant to clock offset, not to clock rate differences |
Google's Chubby lock service put the fence into the lock service itself: a holder can ask for a sequencer (the lock's name, its mode and a generation number) and hand it to the servers it writes to, which check it before acting. For servers that can't check sequencers, Chubby offers a lock-delay: when a holder vanishes, the lock can't be re-granted for a period up to a bound of one minute, which narrows the window for a zombie without closing it.
Single writers across a failover
Promoting a new writer while the old one may still be writing is the classic split brain. The recipe every loop uses:
textSINGLE-WRITER HANDOVER 1. the writer holds a lease on "I am the writer" (a flag read from a majority of copies, or an MRSC item) 2. the old writer stops when it can't renew, by its OWN monotonic clock 3. the new writer waits out the lease plus a margin, by its OWN monotonic clock 4. the new writer bumps the epoch IN THE STORE IT WILL WRITE, before its first write 5. consumers compare (epoch, version) and refuse anything from a lower epoch
The payment loop (step 3.5) uses a 10 s lease and makes the new Region wait 15 s (the lease plus a 5 s margin) before its first ledger write. From the simulator, with the old writer's last good refresh at the moment of the flip (the worst case):
Synthesizing vector architecture diagram...
What to notice: the old writer's right ends at 10 at the latest, the new writer starts at 15: a gap of at least 5 s, never an overlap. If the old Region died at the flip, nothing is written for the whole 15 s.
The standby's clock says the old leader's lease ended. Why does it still wait 5 more seconds?
Where the loops use this:
| Loop (step) | The lease | The epoch and where it is checked |
|---|---|---|
| Ad-click aggregation (3.1) | Pacer lease item with a heartbeat counter; the standby takes over when the counter is unchanged for 3 s on its own clock | epoch on the lease item, taken over with SET epoch = :seen + 1 only if epoch = :seen; every sequence number is (epoch, n), so a paused old pacer's writes are ignored |
| Transactional outbox (3.4) | The promoted database is the writer | writer_epoch bumped inside the promoted database; consumers compare (epoch, version) |
| Stock exchange (2.4) | A Raft arbiter decides the primary | Epoch e+1 committed by the arbiter; the journal refuses e |
| Payment (3.5) and hotel (3.6) | Writer flag as a 10 s lease; the new Region waits 15 s | The database's own promotion fence, plus the wait |
| Email (Round 3) | ack_until: the home Region stops when it is older than its clock minus 30 s | The standby promotes only after ack_until plus a margin on its own clock |
| Wallet (3.4) and Drive (3.4) | Per-tenant or per-shard owner | Owner epoch flipped before the new side writes |
The maps loop (step 2.3) also numbers its weight updates as "epochs", but those only order updates; nobody holds anything, so nothing needs fencing.
What to remember from Part 7
- A leader is a lease holder; its epoch is its fencing token.
- The old leader stops by its own clock before the lease ends; the new one starts after it plus a margin, so there is a gap, never an overlap.
- Stamp the epoch on every write, so anything from an old leader can be refused or parked.
Part 8. When the lock service itself fails
The lock service is a distributed system too. If a failover in it can hand the same lock to two holders, or hand out the same token twice, everything built on it breaks. This Part shows what survives a leader change, what a minority can't do, and why an asynchronously replicated cache is a different kind of lock service.
Where the grant lives
Grants, releases, expiry deletes and the revision counter are entries in the replicated log: the leader answers only after a majority of nodes has them on disk, so they survive the loss of any minority. (How the majority commit works is the Replication, Quorums & Read-Your-Writes loop primitive's subject, and Primitive #09: Consensus covers Raft.) Renewals are not in the log: the leader's lease tracker resets the deadline in memory.
Synthesizing vector architecture diagram...
What to notice: the solid path (grant) goes to the followers and waits for a majority; the dotted path (keep-alive) ends in n1's memory. If n1 dies, the grant survives on n2 and n3; the latest deadline does not.
Trace: the leader dies while C holds the job
| # | t (s) | Event | L3's deadline | C's stop time |
|---|---|---|---|---|
| 20 | 43.2, 46.2 | C renews at n1 | 53.2, then 56.2 (n1's memory) | 51.2, then 54.2 |
| 21 | 47.0 | n1, the leader, crashes. Its in-memory deadlines are gone with it | unknown to anyone | 54.2 |
| 22 | 48.1 | n2 wins the election (its randomized timeout was 1.1 s). Its lease tracker restarts every lease at the full TTL plus one election timeout from now, because it can't know when anyone last renewed: 48.1 + 1 + 10 | 59.1 | 54.2 |
| 23 | 49.2 to 49.3 | C's scheduled renewal goes to n1 and fails; C retries against n2 at 49.3: success | 59.3 | 57.3 |
| 24 | 52.0 | C writes chunk 10 with token 15, sets DONE and removes claimed_since (the row leaves the sparse index), then releases L3: the key is deleted at once, revision 16 | gone | done |
C never stopped working: its failed renewal at 49.2 was well inside its stop time of 54.2, and the retry found the new leader. The new leader's rule, restarting at full length, means a failover can only make a lease longer, never shorter, so no holder is ever cut off early. The price is on the other side: had C died at 47.0 together with n1, its lock would have lasted until 59.1 instead of 56.2. etcd has an optional lease-checkpoint feature (the LeaseCheckpoint feature gate, off by default) that makes the leader persist each lease's remaining TTL now and then, so a new leader can use it instead of the full TTL.
Synthesizing vector architecture diagram...
S5, after event 22: the grant and the token survived on n2 and n3; only the deadline was rebuilt, and it was rebuilt longer.
What a minority can't do
This is the drill The Network Partition That Elected Two Leaders. If a partition cuts a 5-node lock service into 2 and 3 nodes, the 2-node side can't commit anything: a commit needs 3 of 5, and it has 2. So it can't grant, release or expire a lock, or advance the revision counter; an old leader stuck there accepts requests it can never commit, and they time out. The 3-node side elects a leader in a higher term and carries on, and when the partition heals the old leader sees the higher term, steps down, and discards its uncommitted entries. As for size: 3 nodes tolerate 1 failure and 5 tolerate 2 (a majority is 2 of 3 or 3 of 5). Five nodes cost more machines and a commit waits for two followers instead of one; for a lock service whose writes are grants and releases, not every renewal, that cost is usually small, and the extra failure it tolerates matters most during maintenance, when one node is already down on purpose.
What if the lock service is Redis?
A single Redis primary with SET key value NX PX 10000 is a lease: fast and simple. Its failover is the problem. Redis replicates to its replicas asynchronously, so a write the primary acknowledged may never reach the replica that gets promoted. We replay events 19 to 21 on a lock kept in a primary/replica cache, with the token from INCR on a counter key. The replication lag of about 100 ms and the promotion time are assumptions:
Synthesizing vector architecture diagram...
What to notice: the failover lost both the lock and the counter's last increment, so D got the same token as C, and the store can't tell them apart.
After the failover, D got token 15, the same as C. The store checks "reject lower". What happens, and what was the real bug?
Redlock and the argument about it
Redlock, as the Redis documentation describes it, uses N = 5 independent Redis masters (no replication between them):
textREDLOCK ACQUIRE(key, random value, TTL) 1. read the current time 2. try SET key value NX PX TTL on all 5 masters in parallel, each with a short timeout 3. held only if a majority (at least 3) succeeded AND the time elapsed is less than the TTL 4. validity = the initial TTL minus the time elapsed 5. on failure: release on all 5, even the ones that seemed to fail
Three things the Redis documentation itself says: if you care about correctness you should implement fencing tokens, which applies to any distributed lock; Redis does not use a monotonic clock for key expiry, so a wall-clock shift may result in a lock being acquired by more than one process; and after a crash a node must stay down a bit longer than the maximum TTL (a delayed restart) unless it persists with fsync on every write, or its restart can let a second majority form.
The clock problem, from the simulator:
Synthesizing vector architecture diagram...
What to notice: two majorities of 5 always share a node, and here that node's clock jump let it forget the first client. The safety argument depended on node 3's clock.
| Point | Kleppmann's critique (2016) | antirez's reply | What we take from it |
|---|---|---|---|
| Fencing | Redlock has no way to generate fencing tokens; its random value doesn't grow | The unique random value can be used with check-and-set at the resource | Either way, safety comes from a check at the resource; a growing token lets the resource refuse older holders, a random value only lets it refuse other ones |
| Clocks | A clock jump on one node breaks mutual exclusion (the trace above) | Don't let operators step clocks; configure the time daemon to slew instead of jumping | The Redis docs now warn that expiry is on the wall clock; this is an operational promise, not a guarantee |
| Delays and pauses | A pause after acquiring outlives the lock | Redlock re-checks the elapsed time after acquiring; pauses after acquiring hurt every lock system | True, and Jepsen found the same for etcd locks: they allow concurrent holders even in a healthy cluster, and the fix is a revision checked at the resource |
The verdict: Redlock, like every lock, needs a check at the resource for correctness. Use Redis (single or Redlock) for efficiency locks, where a rare double run only wastes work; for correctness locks, use a consensus-backed service or a strongly consistent database that hands out a token that only grows, and check it at the resource. Kleppmann recommends a consensus system such as ZooKeeper for this.
Across Regions
| DynamoDB setup | Where the condition is checked | Safe as a lock? |
|---|---|---|
| One Region's table | One authority for the item | Yes; the lock is unavailable while that Region is |
| Global table, multi-Region eventual consistency (MREC) | On the local replica; conflicting writes are resolved by last writer wins | No: two Regions can both see the lease as expired and both "acquire"; LWW keeps one write and the other holder never finds out |
| Global table, multi-Region strong consistency (MRSC) | Strongly consistent across exactly 3 Regions (3 replicas, or 2 plus a witness); conflicting concurrent writes can fail with ReplicatedWriteConflictException | Yes, with limits: no transactions and no TTL, so the fence must sit on the same item as the data it protects (no ConditionCheck on a separate fence item), or the data must go to a store that checks the token itself |
The same holds for any store: a lock needs one linearizable authority. The Multi-Region Failover loop primitive covers detection and promotion; this page's part is the lease, margin and epoch arithmetic in Part 7.
Restores
A lock service restored from a backup forgets every token it handed out after the backup was taken. From the simulator: C got revision 15 at 40.2; if the service were restored at t = 45 from a snapshot taken at revision 14, the next key created would get revision 15 again, and C and D would share a token exactly as in the cache failover. So every restore jumps the counter: etcd's snapshot restore takes --bump-revision (with --mark-compacted, so watchers don't read history that no longer matches). With a bump of 1,000,000, D's key gets revision 1,000,015. The DynamoDB lock table rule in Part 5 is the same idea.
What to remember from Part 8
- Keep grants and tokens in a store that commits them by majority before it answers; renewals can live in the leader's memory because a new leader restarts every lease at full length.
- A lock service's own failover must never hand out a token twice or shorten a lease: async failover and restores without a bump break the first.
- Redlock, like every lock, needs a check at the resource for correctness; across Regions, a lock needs one strongly consistent authority.
Part 9. Do you need a lock at all?
A lock service is one more system to run, one more thing that can fail over, and a gap with no owner after every crash. Now that we know the resource's check is what really keeps the data safe, the honest question is whether the lock earns its place.
When the database is enough
If the data and the "lock" live in one ACID database, and the work fits in one transaction, use the database: a row lock (SELECT … FOR UPDATE) or a compare-and-set UPDATE … WHERE version = :v is both the authority and the fence, and there is no lease to expire. The hotel loop (step 1.4) holds rooms this way: the hold's expiry is judged on the database's clock (one clock, for rows in that database), a sweeper using SKIP LOCKED releases expired holds, and in the loop's words "the sweeper handles liveness (rooms come back); the conditions handle safety."
The catch is session-bound locks: an advisory lock, or a FOR UPDATE held open across minutes of work, is a lease on a database connection. If the worker pauses while the connection stays up, it still holds the lock; if the connection drops, the database releases the lock while the worker may still be working. It has the same zombie problem as any lease, so writes made outside that transaction still need a fence.
Four ways to let one worker win without a lock service
| Technique | How it works | Where the loops use it |
|---|---|---|
| Compare-and-set (optimistic concurrency) | Read a version; write only if it is unchanged; on failure, re-read and retry | Ride-sharing's accept is a conditional write on offer_epoch and status (step 2.4) |
| Idempotency keys | Store a key with the side effect in one transaction; a repeat with the same key is dropped (the Idempotency & Effectively-Once Processing loop primitive) | Every payment and notification loop |
Claim with FOR UPDATE SKIP LOCKED | Each worker locks the next unclaimed rows and skips rows others hold: "the database's lock is the claim". For a relay that fails over between databases, a writer_epoch bumped in the promoted database fences the old one | Outbox step 1.4 (claim) and 3.4 (writer_epoch) |
| Partition ownership | Split the work so each worker owns a partition; ownership itself is a lease, so it needs a fence at the resource | Kinesis shard leases (Part 11), the mobile news feed step 2.3 fenced by "only if SEQ#n doesn't exist" |
How much to lock
- Granularity: lock the smallest thing that must have one owner: a job, a partition, a row. A global lock serializes everything behind it.
- Throughput: keep-alives all land on the lock service's leader, so its load is holders ÷ renewal interval: 10,000 holders renewing every 3 s is 10,000 ÷ 3 ≈ 3,333 keep-alives a second, served from memory. Grants, releases and expiries are the expensive ones: each is a majority disk write.
- Ordering: a worker that needs several leases takes them in one fixed global order, so two workers can never wait on each other in a cycle.
The same job three ways
Three workers start chunk 2 (5,000 invoices) at the same moment. From the simulator:
| Approach | Invoices computed | Calls to the invoice writer | Applied | Thrown away |
|---|---|---|---|---|
| Lease + fence | 5,000 (only the holder works) | 5,000 | 5,000 | 0 |
Optimistic concurrency only (write if chunks_done is still 1) | 15,000 | 3 commit attempts | 5,000 (A's commit) | 10,000: B and C lose the compare and discard their work |
Idempotency keys only (inv#2026-09-28#05001 …) | 15,000 | 15,000 | 5,000 | 10,000 calls dropped as duplicates by the target |
An equal-terms comparison
| Lease + fence | Optimistic concurrency | Idempotency keys | Same-database transaction | |
|---|---|---|---|---|
| Extra state | A fence field per resource, plus the lock service | A version field per resource (the same cost as a fence field) | A key table, kept as long as a retry can arrive | None beyond the row |
| Correctness comes from | The fence check at the resource | The version check at the resource | The key check at the target | The database's own locking |
| What a paused worker does | Wakes, is refused by the fence | Wakes, loses the compare | Wakes, its repeat is dropped | If its transaction is still open, blocks others; else is refused |
| Work wasted under contention | Almost none: only the holder works | Grows with contention: every loser redoes its work (retry storms at high contention) | Every duplicate caller does all the work | Waiters block instead of wasting work |
| Needs the target's cooperation | Yes: it must store and check the token | Yes: it must check the version | Yes: it must store the key | Data and lock must be in the same database |
| Failure costs | A lock service to run; a gap of one lease after a crash; renewals | Retries | Key storage; a dedup window to size | Only works while the work fits one transaction |
Why not drop the lock and just put a version check on every write?
For a deeper look at isolation levels and what FOR UPDATE really locks, see Primitive #21: Isolation levels.
What to remember from Part 9
- If the data and the lock live in one database and the work fits one transaction, the database is the lock and the fence.
- A version check on the write is the fence; the lock only decides who should try, and optimistic concurrency wastes work under contention.
- For side effects outside your store, idempotency keys are the safety net; the lease only makes duplicates rare.
Part 10. End to end: from your code to the disk
So far "acquire", "renew" and "write with the token" were single steps. In reality each crosses several layers, and a pause or a delay can happen in any of them. Knowing which layer does what explains why a grant costs a disk write on a majority of machines and a renewal doesn't, and why the only check nothing can bypass is the one at the store.
The layers
Synthesizing vector architecture diagram...
What to notice: the store sits beside the lock service, not behind it. The worker talks to both, and a pause anywhere in the worker process sits between them.
An acquire, traced down
Event 1: A acquires at t = 0.0.
Synthesizing vector architecture diagram...
What to notice: the leader answers after its own disk and one follower's; the third node is not waited for. The grant and the key creation are two log entries in etcd (the lease grant, then the create-key transaction); both follow this path.
A keep-alive, traced down
Event 4: A renews at 3.0. The library records sent_at, the request crosses the network, and the leader resets L1's deadline to 13.0 in memory and replies. There is no log append, no replication and no fsync (what those cost on each node is in Write-Ahead Log, fsync & Group Commit). Keep-alives cost the leader CPU and network only; a slow disk hurts them only indirectly, because etcd's documentation warns that long fsync latencies can make it miss heartbeats, causing request timeouts and temporary leader loss, and an election pauses renewals (Part 8).
A fenced write, traced down
Synthesizing vector architecture diagram...
What to notice: both requests hit the same one-step check. Nothing between A's pause and the store could stop A's request from being sent; the store is the one place it can't get past.
Who does what
| Layer | Its job | What goes wrong there |
|---|---|---|
| Client library | Heartbeat thread; sent_at and the stop time on the monotonic clock; retry against a new leader; stop the app when the stop time passes | A GC pause freezes the heartbeat thread with the app; a CPU-starved thread looks like death |
| Network | Carry requests and replies | Delays, retries, partitions: a write can arrive long after it was sent |
| Leader | Validate; keep deadlines in memory; propose grants, releases and expiries | Crash: deadlines lost, rebuilt at full length by the next leader |
| Log | Append, replicate, majority on disk, commit, apply | Slow disks delay commits and heartbeats and can trigger elections |
| Store | Conditional write with the token | Nothing, if every write is conditional; everything, if one path isn't |
The sidecar trap: if a separate process (a sidecar or agent) sends the heartbeats on the app's behalf, the heartbeats prove the sidecar is alive, not the app. A frozen app keeps its lease while doing nothing, and the fence becomes the only protection all the time.
What to remember from Part 10
- A grant is real only once a majority has it on disk; a keep-alive touches only the leader's memory.
- A pause can happen at every layer; the only check that can't be bypassed is at the storage.
- Heartbeats from a sidecar prove the sidecar is alive, not the app.
Part 11. On AWS
The mechanism is the same on AWS; what changes is which piece plays the authority, which plays the fence, and which services only look like locks.
Managed services that use it
| Service | What it uses | What AWS documents |
|---|---|---|
| Amazon DynamoDB | Conditional writes (ConditionExpression), TransactWriteItems with a ConditionCheck, global tables | Condition expressions have no time function; a transaction holds up to 100 actions, none on the same item, and fails as a whole with TransactionCanceledException; a partition gives up to 1,000 write units a second; TTL deletes expired items typically within a few days; MRSC global tables span 3 Regions (3 replicas or 2 plus a witness) with no transactions and no TTL, and conflicting writes can fail with ReplicatedWriteConflictException |
| Amazon S3 | Conditional writes: If-None-Match: * and If-Match: <ETag> on PutObject, CopyObject and CompleteMultipartUpload | 412 on a failed condition; the first of several concurrent conditional writes to finish wins; 409 or 404 when a delete races the write; bucket policies can require conditional writes |
| Amazon SQS | The visibility timeout: a lease on a message (default 30 s, extendable with ChangeMessageVisibility up to 12 hours from first receipt) | A lease without a fence. A consumer whose timeout expired keeps working and its side effects land; its DeleteMessage with an old receipt handle "will succeed, but the message might not be deleted". Standard queues can also deliver a message more than once within the timeout |
| Amazon EKS | The managed control plane runs etcd (three etcd instances across three Availability Zones); controllers elect leaders with Lease objects | client-go's election "does not guarantee that only one client is acting as a leader (a.k.a. fencing)": add your own epoch and check it |
| Amazon Aurora | SQL compare-and-set, FOR UPDATE SKIP LOCKED, advisory locks, and the database's now() as one clock for its own rows | A Global Database managed failover attempts to stop writes in the old primary Region ("write fencing"), but AWS calls it best-effort: writes might be momentarily accepted there, causing split brain. The recovered old primary returns as a secondary on a new storage volume. So a Region-level writer lease and epoch (Part 7) are still needed |
| Amazon ElastiCache (Valkey or Redis OSS) | SET key value NX PX locks | Replicas are kept in sync asynchronously in the standard configuration, so a failover can lose a lock or repeat an INCR token; Redis key expiry uses the wall clock. Use it for efficiency locks only |
Running it yourself
| Option | What it is | Sizing and notes |
|---|---|---|
| etcd or ZooKeeper on Amazon EC2 | A consensus lock service you operate: one node per Availability Zone | 3 nodes tolerate 1 failure, 5 tolerate 2. Keep the election timeout well above the cross-AZ round trip (etcd's defaults: 100 ms heartbeat, 1,000 ms election timeout). Sync clocks with the Amazon Time Sync Service (microsecond-level on supported instances; the open-source ClockBound daemon reports an error bound) |
| Amazon EBS (gp3) or instance-store NVMe for the log | Where grants, releases and expiries are made durable | The lock service holds little data; what matters is fsync latency and quorum round trips, not throughput. A disk shared with noisy work is etcd's documented cause of missed heartbeats and leader loss |
| DynamoDB Lock Client (an awslabs library on DynamoDB) | Leases on a DynamoDB item, with a heartbeat thread | Never stores absolute times, only the lease duration; a claimant takes over only if the record version number is unchanged after a full lease on its own clock. README example: 10 s lease, 3 s heartbeat. Supported by the community, not by AWS. Its record version number is a random value, not a number that grows, so add an epoch of your own and check it at the resource |
| Kinesis Client Library (a library on Kinesis Data Streams, lease table in DynamoDB) | One lease per shard: leaseOwner, and leaseCounter "so that workers can detect that their lease has been taken by another worker" | A lease owner is considered failed after failoverTimeMillis (default 10,000 ms). Checkpoints are written against the lease, so write output idempotently: a worker that lost its lease can still have records in flight |
Look-alikes that are not this mechanism
| Look-alike | Why it looks like a lock | Why it isn't |
|---|---|---|
| S3 Object Lock | The word "lock" | Write-once retention (governance and compliance modes, legal holds): it stops deletes and overwrites for a period and says nothing about who may act |
| DynamoDB TTL | "The item expires" | Deletion typically within a few days, never a deadline; and it resets an epoch kept on the item (Part 5) |
| SQS receipt handle | Hides a message from other consumers | A lease without a fence: a stale consumer keeps working, and its delete reports success whether or not it deleted |
| SQS FIFO message groups | One message of a group is in flight at a time | Ordering while a message is in flight; AWS says it doesn't "lock" the group, and the next message becomes available when the visibility timeout expires |
| DynamoDB transactions | "Isolated" | Optimistic: conflicting transactions are cancelled; there is no holder and no expiry |
| DynamoDB global tables (MREC) | Replicated conditional writes | Conditions are checked on the local replica, and last writer wins: two Regions can both "acquire" (Part 8) |
| Route 53 ARC routing controls, health-check failover | They pick one "active" Region or endpoint | Routing, not a fence: the old side can still write if anything reaches it |
Aurora DSQL SELECT … FOR UPDATE | Looks like a row lock | DSQL uses optimistic concurrency: FOR UPDATE takes no lock and doesn't block; a conflicting transaction fails at commit with a serialization error |
What to remember from Part 11
- Pick the authority (a DynamoDB item, etcd or ZooKeeper, one SQL database) and put the token check in the store you write.
- DynamoDB and S3 conditions have no server time; one SQL database's
now()is a single clock, fine for rows in that database only. - Caches, SQS receipt handles, TTL, Object Lock and routing controls look like locks and are not fences.
Part 12. What you've learned
Back to the billing job
A froze with 5,000 invoices in its buffer, and nobody was billed twice. Here is how each piece did its job:
- A held a lease timed on its monotonic clock, with a stop time 2 s before the service's deadline (S1, Parts 2 and 3). B's wall clock jumped and nothing changed, because no duration was timed on it.
- A paused past its lease. B took over with a higher token, and its first act was to raise the fence to 12, so A's late chunk arrived with 7 and was refused (S2 and S3, Part 4). The check was one conditional write in the store (Part 5).
- B died holding the job. The lock freed itself at 29.5; the claim on the row did not, but it was findable through the sparse index, and a sweeper re-queued it with a conditional write; C took over with token 15 (S4, Part 6).
- The lock service lost its leader. The grant and the token were safe in the majority log; the new leader rebuilt the lease longer, and C's retry found it (S5, Part 8).
- C finished, set
DONE, left the sparse index and released the lock (S6: no holder, fence 15, rowDONE).
The whole story, event by event
| # | t (s) | Event |
|---|---|---|
| 1 | 0.0 | A acquires: key created at revision 7, token 7; deadline 10.0 (leader memory); A's stop time 8.0 |
| 2 | 0.5 | A raises the fence 0 → 7 in one conditional write: RUNNING, owner 7, claimed_since 0.5 |
| 3 | 1.0 | A writes chunk 1 with token 7: accepted |
| 4 | 3.0 | A renews: deadline 13.0 (memory only), stop time 11.0 |
| 5 | 3.9 | A checks 3.9 < 11.0 and builds chunk 2 with token 7 |
| 6 | 4.0 | A's 15 s pause starts; its 6.0 and 9.0 renewals never happen |
| 7 | 13.0 | n1 finds L1 expired and commits a revoke; A's key deleted at revision 11 |
| 8 | 13.5 | B acquires at revision 12, token 12; deadline 23.5, stop time 21.5 |
| 9 | 13.6 | B raises the fence 7 → 12, owner 12, claimed_since 13.6; reads chunk 1 done |
| 10 | 14.0 | B writes chunk 2 with token 12: accepted |
| 11 | 16.0 | B's wall clock steps back 3 s; B renews at 16.5 and 19.5 on its monotonic clock (deadline 29.5, stop time 27.5) |
| 12 | 19.0 | A wakes; its chunk 2 (token 7) arrives at 19.05: rejected; its keep-alive at 19.1: lease not found; A drops the job |
| 13 | 20.0 | B writes chunk 3 with token 12 |
| 14 | 21.0 | B is killed; its last renewal reached n1 at 19.5 |
| 15 | 21.0 to 29.5 | The lock is held by a dead process |
| 16 | 29.5 | L2 expires; B's key deleted at revision 14; the row still says RUNNING, owner 12 |
| 17 | 30.0 | Sweeper finds the orphan through the sparse index: orphan_seen_at = 30.0 |
| 18 | 40.0 | Sweeper re-queues: QUEUED, only if RUNNING, owner 12, fence 12 |
| 19 | 40.2 | C acquires at revision 15, token 15; at 40.3 raises the fence 12 → 15; resumes at chunk 4 |
| 20 | 43.2, 46.2 | C renews; deadline 56.2, stop time 54.2 |
| 21 | 47.0 | n1, the leader, crashes; in-memory deadlines lost |
| 22 | 48.1 | n2 wins the election; L3 restarted to 48.1 + 1 + 10 = 59.1 |
| 23 | 49.2 to 49.3 | C's renewal to n1 fails; the retry to n2 succeeds: deadline 59.3, stop time 57.3 |
| 24 | 52.0 | C writes chunk 10, sets DONE, removes claimed_since, releases: key deleted at revision 16 |
The cheat card
| Topic | Remember |
|---|---|
| Lease | A lock that expires; a crash blocks others for at most one lease |
| Fencing token | From a counter that only grows and is majority-committed; checked at the resource |
| The fence rule | Reject lower; equal is the current holder |
| New holder's first act | Raise the fence in one conditional write that returns the old state, then read |
| Renew | Every third of the lease; time from the send |
| Stop time | Time the last renewal was sent + lease − margin, on the monotonic clock |
| Earliest safe takeover | Last renewal the authority applied + lease (+ a margin if the claimant uses its own clock) |
| Clocks | Durations on monotonic clocks; offset matters only when one clock writes a deadline another reads |
| DynamoDB lease | :now is a client's clock: margin bigger than the worst offset, or wait a full lease and check the version is unchanged |
| TTL | Cleanup only, days late; never on the item holding the epoch |
| Orphans | Claim marker in a sparse index, written with the fence raise; sweeper re-queues conditionally |
| Election | Lowest revision leads; watch your predecessor; epoch = key revision |
| Lock-service failover | Grants and revision in the log; deadlines in leader memory; new leader restarts leases at TTL + election timeout |
| Redis | Async failover can repeat a token; efficiency locks only; Redlock also needs a resource check |
| Across Regions | One linearizable authority: one Region or MRSC with the fence on the data's item; never MREC |
Failure checklist
- Does every protected write carry the token, and does the resource check it (reject lower, equal accepted)?
- Does a new holder raise the fence before it reads or writes?
- Are all durations (renewal timers, stop times) on a monotonic clock?
- Is any lease condition comparing a timestamp written by one machine's clock with another machine's clock, and is the margin bigger than the worst offset?
- Is every host's clock offset alarmed at a threshold well below the margin?
- Can the token counter ever go backwards: a TTL on the epoch item, a deleted lock item, a restore without a bump, an async failover?
- Is every claim marker findable through an index, written in the same write as the fence raise?
- Does the sweeper re-queue conditionally, and does it pause instead of storming when many leases expire at once?
- Are side effects outside your store (emails, payments, webhooks) protected by idempotency keys?
- Are heartbeats sent by the app process itself, not a sidecar?
- Is the lock's authority one linearizable store (not an async cache, not an MREC global table)?
Think-first drills
Drill 1. A lease of 20 s, renewed every 6 s, with a margin of 4 s. Renewals sent at 0, 6 and 12 succeed, then the network drops. When does the holder stop, and when does the service free the key?
Drill 2. A store's fence is 40. Writes arrive with tokens 41, 39, 41 and 42, in that order. Which are accepted, and what is the fence after each?
Drill 3. A lock item lives in a DynamoDB global table (MREC) replicated to two Regions. Its epoch is 7. A worker in each Region sees the lease as expired at about the same moment and claims. What happens, and name two fixes.
Interview questions
| Question | Model answer |
|---|---|
| Walk me through acquiring and renewing a lease, and how leader election is the same thing. | Create a key bound to a lease with a TTL, only if absent; the service commits it through a majority log and returns the key's revision as the token. Renew every third of the lease; the leader resets the deadline in memory. The holder stops at the send time of its last successful renewal plus the lease minus a margin, on its monotonic clock. An election is the same lease on a role's key: lowest revision leads, each candidate watches its predecessor, and the leader's revision is its epoch. |
| Why doesn't a TTL lock protect storage from a paused client? | The paused client doesn't run, so it can't notice its lease expired; when it wakes, its next write goes out as if it still held the lock. A pause can land between any check and the write, so no client-side check or lease length prevents it. Only the storage can refuse, by checking a token that only grows. |
| Where is the fencing token checked, and what if the resource can't check it? | At the resource, on every write, in the same atomic conditional write as the data: reject lower, accept equal; a new holder raises the fence first. If the resource can't check it (an email provider, a webhook), use idempotency keys; the lease only makes duplicates rare. |
| How do you build a lease on DynamoDB when conditions have no server time? | Holder, epoch (from if_not_exists(epoch, 0) + 1, never deleted), lease_until from the holder's clock and a record version. Either claim when lease_until is older than your clock minus a margin bigger than the worst offset, or, like the Lock Client, wait a full lease on your monotonic clock and claim only if the record version is unchanged. Never use TTL as the deadline; fence the resource with the epoch. |
Why not just SELECT … FOR UPDATE? | If the data and the lock are in one database and the work fits one transaction, do exactly that: the row lock is authority and fence. It breaks down when the lock is held across long work or other systems: a session-bound lock is a lease on a connection with the same zombie problem, and writes outside that database still need a token. |
| Redlock: safe or not? | Not by itself for correctness: it gives no growing token, and its safety depends on bounded pauses, delays and clock behaviour (Redis's key expiry uses the wall clock). The reply that you can check-and-set the random value at the resource concedes the point: the resource check is what gives safety. Use Redis locks for efficiency; for correctness, a consensus-backed lock with a token checked at the resource. |
Where to go next
- The Idempotency & Effectively-Once Processing loop primitive (coming): protecting side effects the fence can't reach.
- The Replication, Quorums & Read-Your-Writes loop primitive (coming): how the majority commit behind every grant works. Meanwhile, Primitive #09: Consensus, Raft and Paxos.
- The Multi-Region Failover loop primitive (coming): detection and promotion around the single-writer handover in Part 7. Meanwhile, Primitive #24: Disaster recovery and multi-Region.
- Write-Ahead Log, fsync & Group Commit: what "a majority has it on disk" costs on each node.
- Primitive #15: Distributed unique ID generators: the machine-number lease and the ID fence in the ID generator loop.
- Drills: The GC Pause That Corrupted Shared Storage (answered in Parts 4 and 9) and The Network Partition That Elected Two Leaders (answered in Part 8).
- Loops: the ID generator (steps 1.4, 2.2, 2.5, 3.2), the job scheduler (steps 1.2, 1.3, 2.4, 2.5, 3.3), the message queue (steps 1.2, 2.5) and ride-sharing (steps 2.3, 2.4, 3.6).