Lesson 035 · Phase 2, Mechanisms

Distributed Locks and Their Sharp Edges

A lock promises you something about a resource, and almost every implementation of one actually promises you something about a client instead.

22 min read

Lesson 35 · 35 published · 90 planned

On this page
The systems in this lessonUsed here: Marlow Books, Stagefront, and Galewatch.

Made up for this course and reused from lesson to lesson so their numbers become familiar. None of them exist. All three

Marlow Books · A small online bookshop
Four people, one server and one Postgres database. About 40 requests a second on a normal day and ten times that in the week before Christmas. The one box that the early lessons stress until it breaks.
Stagefront · An event ticketing service
Quiet most of the time, then a stadium show goes on sale at 10:00 and two hundred thousand people press the same button in the same minute. Oversold seats are a lawsuit, so correctness matters as much as speed.
Galewatch · Telemetry for wind farms
Nine hundred turbines, a reading every two seconds, over links that drop for hours in bad weather and come back with a backlog. Dashboards that lag by seconds, reports that scan a year.

Marlow Books is an online bookshop run by four people, and it exists only in this course. In February its founder took the first week off since the shop opened, and while they were away the catalogue put two copies of a novel on the shelf that did not exist.

Three weeks earlier, on the wet Saturday in January that lesson 030 opens with, they had drawn the whole shop on one page so somebody else could cope for a week. Two items off that page shipped before the train. The first was the one lesson 030 argued for hardest: move the stock decrement into the same commit as the pending order row and its idempotency key, so a copy is held at the moment of intent rather than when the card clears. That is lesson 015's shape carried back at last from Stagefront, the ticketing service where two hundred thousand people press the same button at ten in the morning. It also created a refund path whose job is putting a held copy back on the shelf.

The second item was a lock.

Lesson 019's sweeper runs every two minutes over orders pending for longer than two minutes, asks the payment provider what happened to each idempotency key, and either completes the order or releases it, refunding the charge if one landed. Lesson 027 watched it jam: a dozen questions to a provider that answered none, at the ten second timeout each, came to a hundred and twenty seconds, which is the sweeper's whole schedule. The crontab sits on both boxes, copied with the image back in lesson 005. A run longer than its own interval gets started again on the other machine while it is still going, and nobody is watching.

So the founder put a lock in front of it, with the two lines everybody uses.

SET lock:sweeper <hostname> NX PX 120000
  ... run the sweep ...
DEL lock:sweeper

NX sets the key only if it does not exist, so one box gets an answer and the other gets nothing. PX 120000 expires it after a hundred and twenty seconds, so a dead holder cannot block the next run forever. That number was not careless: the founder matched it to the schedule so a dead holder costs at most one skipped cycle, which is a defensible rule and the reason this is a lesson and not a scolding.

On the Wednesday of that week the provider went slow rather than dark. Answers crossed ten seconds at ten to nine in the evening, and two things started happening at once. Checkouts began hitting lesson 019's ten second client timeout and leaving pending rows behind, about one in five of them at lesson 027's rate of an order every fifty seconds, which is seventy two orders in the hour that followed. And the sweeper's own questions crossed the same line: ten seconds each, timed out, no answer, so nothing it asked about could be settled and every row went round again on the next run.

So the list grew by a row about every four minutes, nothing ever came off it, and a run costs ten seconds a row. Twelve rows is a hundred and twenty seconds, which is the whole lease, so somewhere around twenty to ten the sweep stopped fitting inside the lock it was holding, and from then on each run overlapped the next by a little more than the one before. Nothing marked that moment anywhere, and nothing went wrong either: a sweep that settles nothing does no damage, however many copies of it are running.

What changed at 21:50 is that the provider started answering. Slow is not steady, and the replies began landing just inside the ten second timeout instead of just outside, so for the first time that hour the sweep was finishing rows instead of collecting them. A finished row puts a held copy back on the shelf. Here is the cycle where that cost something.

21:50:00, box B's cron fires. Fourteen rows on the list, and it starts down them.

21:52:00, two things happen in the same second. Box B, a hundred and twenty seconds in, has answered twelve of the fourteen and is starting the thirteenth. The key it set expires, because that is what PX 120000 means. Box A's cron fires, finds no key, sets it, and starts its own pass over the two rows still saying pending.

Both asked about the same order. Both got the same answer, that the money had never moved. Both marked it abandoned, which is harmless, and both ran update books set stock = stock + 1 on the same ISBN, which is not. One held copy went back on the shelf twice.

Then box B finished at 21:52:20 and ran DEL lock:sweeper, which deleted box A's key rather than its own, box B's having expired twenty seconds earlier. That half of it had been going on since the overruns started, because a run that outlives its lease releases whatever key the next run is holding. What was new at 21:52 was that two sweeps settled the same row.

Nobody logged anything, because nothing had failed. Both boxes believed they held the lock and both were telling the truth as they understood it. On the Friday a temporary member of the packing staff went to the shelf for a title the catalogue said had two copies and found none, which is how the shop learned. It is the second time this course has watched the catalogue promise copies the shelf did not have, and lesson 009 was the first. The founder read about it on the train home.

What a lock promises, and who it promises it to

A distributed lock is mutual exclusion between processes that do not share memory. You want at most one of them touching some resource at a time, and they sit on different machines, so no kernel and no runtime can see both.

Hold that against a mutex inside one process, which you have used and never worried about. It works because the lock and the thing it protects share an address space and a scheduler, and the holder cannot outlive the lock: if the holder dies the address space goes with it, mutex included. No clock anywhere.

Move the resource to another machine and every one of those properties goes. The lock service cannot see the resource, so it cannot promise you anything about it. All it can promise is something about clients: that it has not told two of them yes at once, by its own reckoning of time, which is a weaker claim than no two of them acting at once. A lock is a promise about a resource, and almost every implementation makes a promise about a client instead. Everything below is either a way of narrowing that gap or a way of living with it.

The implementation choice comes second, though, behind a question whose clearest version is Martin Kleppmann's, from an argument this lesson reaches shortly. Why do you want the lock?

If two runs would waste work, you have an efficiency lock. One of twelve pods should rebuild the cache rather than all twelve, two runs cost CPU and a bill, and a lock that fails occasionally costs you a duplicate you can price and shrug at. If two runs would corrupt something, you have a correctness lock, and a lock that fails occasionally gives you the hook.

Lesson 032 gave you the same question from the other side: a decision needs one owner when its second copy contradicts the first rather than repeating it. Answer it before you shop, because the two answers justify completely different amounts of machinery. The sweeper looked like an efficiency problem, where two runs cost two API calls, and it was a correctness problem all along. What made it one was a single stock + 1 four lines from the end.

Marlow's other scheduled job needs none of this. The monthly publisher payout runs at 2 am on the first Sunday, every box tries to insert one row into job_runs keyed on the job name and the date, and whoever writes the row first runs the job. Lesson 032 spent a section on why that election has never picked two winners: it asks who was first, which a unique constraint answers correctly forever, and it never asks who is alive, so it has no clock in it.

Now look at what the sweeper did to that key. Seven hundred and twenty runs a day do not fit in (job_name, run_date), which has a resolution of one day, and lesson 019 put the sweeper in that table for its deadline rather than its exclusion: one row saying the job ran today, which is what lesson 018's watcher reads. A deadline row is not an election. The lock the shop needed was the one it already had, with a finer key. Write (job_name, run_slot) where the slot is the two minute bucket, insert with on conflict do nothing, and you have the payout's exclusion in the database the sweep was going to need anyway, with no expiry and no Redis.

That fixes one of the hook's two problems. A per-slot insert stops both boxes starting the same slot, which is the entire crontab-copied-with-the-image failure. It does nothing about 21:52, where the run that overran belonged to the previous slot and the new slot's key was free. Keep the two apart, because most teams reaching for a distributed lock are solving the first and most of the sharp edges belong to the second.

The four things you might actually type

In the order they get chosen, which is close to the reverse of how often they are right.

A row or an advisory lock in the database you already have. Lesson 015 taught the transaction. select ... for update holds a row until the transaction ends, and the database releases it whether you remember to or not. Postgres also has locks attached to no row: pg_try_advisory_lock(hashtext('sweeper')) answers true or false at once, holds the lock for the life of the connection, and lets go when the connection closes.

Read that last clause twice, because it is the property the other three lack. No expiry, so nothing to renew and no number to guess. A run that overruns its slot still holds the lock, so the next slot gets false and skips, which is what the unique insert could not do.

One caveat, and it is worth more than most of the rest. A process that exits releases the lock at once, because the server watches the socket close. A box that loses power releases it when the server notices the socket is gone, and with default TCP keepalives that can be hours. Set tcp_keepalives_idle or tcp_user_timeout on that connection and look at the number you chose, because the clock did not disappear. It moved into the kernel, where nobody argues about it.

Three sharp edges. The lock holds a database connection for the whole run, so a hundred and forty second sweep is one connection out of Marlow's max_connections of 100, of which lesson 011 counted ninety six reserved by application pools: fine for one cron job, not for twelve pods. Session-scoped advisory locks are a trap behind a pooler in transaction mode, where you do not know which backend you got and the release can land on a different one than the acquisition did. Marlow has no pooler, which is the only reason this is safe there. And for queue-shaped work you want FOR UPDATE SKIP LOCKED, which is one lock per row, letting eight workers share a table without meeting.

SET key value NX PX in Redis. The hook's version, and the one on every blog. The missing line is in the release:

-- release, as one Lua script so GET and DEL cannot be split
if redis.call("GET", KEYS[1]) == ARGV[1] then
  return redis.call("DEL", KEYS[1])
end
return 0

The value has to be a fresh random string per acquisition, not the hostname the founder used, and the release has to compare it. Comparing in the client is not enough: between your GET and your DEL the key can expire and somebody else can take it, so the two have to be one operation.

Now the part that matters more than the fix. That script stops box B deleting box A's lock. It does not stop box A holding box B's. At 21:52:00 box B was working twenty seconds past an expiry it could not see, and a random value changes nothing about that: the second sweeper was already running and the shelf was already wrong before any DEL happened. Checking the value before you release stops you deleting somebody else's lock. It does not stop somebody else holding yours. The two bugs look like one, and the famous fix only touches the cheap one.

The expiry is where the rest lives, and lesson 032 took that apart: a lease is a promise to stop, kept by the machine that is losing it, reading its own clock, which a process that is paused, swapping, blocked on a socket or waiting out a collection does not have. Lesson 034 then showed that clock stepping backwards six seconds and handing a losing holder time it does not own. One failure here belongs to Redis specifically: a master with a replica replicates asynchronously, so a failover can lose the last few writes, lock keys among them, and two clients hold the same lock with nothing anywhere reporting trouble. That window is what lessons 011 and 012 were about in Postgres, and nothing about a cache makes it smaller.

Redlock, across several independent Redis nodes. The answer to the failover problem, and where the interesting argument lives. Five Redis masters replicating nothing to each other, the same SET key value NX PX to all five, and you count. If a majority say yes and the time spent collecting those answers is comfortably under the expiry, you hold the lock for the expiry minus that time minus an allowance for clock drift. Otherwise you release everywhere and have nothing. The majority is lesson 033's, for lesson 033's reason: two majorities of one set must share a member.

In February 2016 Martin Kleppmann wrote that this is not safe enough to build correctness on, and Salvatore Sanfilippo, who wrote Redis, replied that the criticism assumed more than it was entitled to. Both were right about what they were each talking about, and the disagreement is genuinely about which assumptions you may make, which is why it is worth reading ten years later instead of being settled.

Kleppmann's case is that the algorithm's safety rests on timing: clocks that do not jump, pauses that are bounded, network delays that are bounded. Break one of the three on one node and mutual exclusion goes, and none of the three is a property you can buy. Then the deeper point, which survives everything else: without a fencing token no lock of any kind is safe for correctness, so arguing about which lock to use is arguing about the wrong layer.

Sanfilippo's case is that the pause argument is not specific to Redlock. Every lease-based lock has it, the consensus-backed ones people recommend instead included, so it is no reason to prefer one over another. The clock argument he treated as operational rather than a design flaw, since a daemon that slews small corrections rather than stepping does not produce the jump the attack needs. Lesson 034 is the reason to take that half seriously and also the reason not to relax: Galewatch, which collects a reading every two seconds from each turbine on the wind farms it monitors, measured a fleet whose worst clock was forty one seconds out after two years of reporting cheerfully.

My own view, and it is a view rather than a finding. I would not run Redlock, and not because the argument went against it. Five masters is five pieces of infrastructure bought to produce a lock that still needs a fence before it is safe for anything that matters. If the resource can hold a fence, one Postgres row would have done. If it cannot, five masters do not help.

A key in a consensus cluster: etcd or ZooKeeper. Lesson 033 built the machinery. The lock is a key, created with a lease you keep alive, held on a majority before the cluster tells you anything, so there is no asynchronous replica to lose it on a failover.

What you get here and nowhere else is the token. Every write in etcd carries a cluster-wide revision that only goes up, because it is a position in a Raft log, and lesson 033's definition of committed is why that position cannot move: a majority has already stored the entry there. ZooKeeper's equivalent is the sequence number on an ephemeral sequential node, handed out by the cluster, and the node deletes itself when your session dies. Either way the lock service hands you a number lesson 032 would recognise on sight: minted by the only party allowed to mint it, monotonic, impossible for a stale holder to invent. Kleppmann's post recommends exactly this.

The edges: the session timeout is still a lease, so your old holder still has to stop itself, and every lock operation is a consensus round, which lesson 033 priced at half a millisecond inside one data centre and two hundred across the world. A lock per request is a design error you will feel.

Lock What makes it safe What breaks it
Postgres row or advisory lock and resource on one machine a pooler; a held connection
SET NX PX nothing, by itself the expiry; async failover
Redlock a majority of independent nodes clock jumps; long pauses
etcd or ZooKeeper consensus, plus a real token the session lease; a round trip

Two of those four middle columns hold without an argument about assumptions: the one where the lock and the resource are the same machine, and the one where the lock service hands you a number the resource can check.

The only safe lock is one the resource enforces

Take lesson 032's fencing token as given: a number that only goes up, minted by the only thing allowed to mint it, carried on every statement, compared by the resource against the highest it has already obeyed. The resource keeps the score.

Lesson 032 sorted resources in a single paragraph and moved on, because fencing works only where the resource can remember. That paragraph is the whole decision, and it was filed as a footnote to the fence. Do the sorting first instead, before you choose a lock, because it decides whether you need one.

A SQL database can keep score, and you have already shipped a fence into one, though not in the statement you would guess. Lesson 015's update books set stock = stock - 1 where isbn = ? and stock > 0 is not a fence, and lesson 019 said so in four words: atomic is not idempotent. It closes the window between reading and writing, and subtracting one does not destroy stock > 0 on any copy but the last. The fence you already own is lesson 019's first pile, a state transition guarded by the state it is transitioning out of, where the resource compares the request against its own state and refuses what it has already handled. That is a fencing token's whole job with the token thrown away, and the sweeper's correct version is that pile pointed at a status column.

update orders set status = 'abandoned'
 where id = $1 and status = 'pending';
-- only if that returned one row:
update books set stock = stock + 1 where isbn = $2;

One transaction, the increment conditional on the first statement's row count. The second sweeper's update matches zero rows, so it never increments, and the shelf is right however many sweepers run. The row's own status is a fence with two values. Lesson 030 warned about the mirror image of this in its own words, that putting the decrement in the second transaction lets a retrying sweeper take a second copy off the shelf. February got it from the other direction: the fix that made the checkout correct created the increment nothing was guarding.

An object store can keep score where it offers a conditional write, and writing the same key twice with the same bytes is harmless anyway, which is why lesson 019 put each publisher's CSV under a key of (publisher, month) rather than reaching for a lock. Lesson 044 owns the internals.

A payment provider cannot. Nor can an email provider, somebody else's API, or a bare file on a mounted volume. There is no token to carry and nowhere to compare it, and no amount of care on your side changes that, which is lesson 019's third pile arriving from a new direction. The lock buys a smaller window and not a guarantee. A lock without a fenceable resource is a performance optimisation wearing a safety hat, and if you need more than that the answer is lesson 019 or a declared at-most-once somebody has agreed to in writing.

Which is the honest end of the argument. Lesson 019 found a recursion it could not escape: a claimed row that gets stuck needs an expiry, an expiry is a lease, and a lease needs the work behind it to be safe twice. Across machines that gets worse, because the expiry and the work now sit on separate hardware reading separate clocks, with lesson 034's error bars between them. You do not get out of it by finding a better lock. You get out by making the work safe to repeat, at which point the lock stops being load bearing and becomes what it should have been all along, which is a way of not doing work twice.

Where a lock is still the right answer

None of that makes locks a mistake, and a lesson that left you refusing to use one would have taught the wrong thing.

The efficiency case, declared as what it is. Twelve pods, one cache rebuild, a Redis key with an expiry, no fence and no renewal. If it fails you rebuild a cache twice, and a comment saying "this is an efficiency lock, a double run is harmless, here is why" is worth more than any implementation detail. Half the locks in the world are this and the other half wish they were.

The case where the lock and the resource are the same machine, which is Marlow's whole shape: one Postgres, an advisory lock, the work inside the transaction. Lesson 032 put it in a line for the payout job, that the arbiter and the victim are the same database. When that is true, take it, and stop reading.

And the case where the lock service mints a number and the resource checks it, etcd on one side and a version column on the other. You can tell it is the real thing because the lock goes boring: if the token is checked, a second holder cannot do damage, so the lock is only saving you a wasted attempt.

Which leaves the strongest argument in this course for not reaching for a lock, and it belongs to Stagefront. The on-sale minute is the worst contention in these lessons, two hundred thousand people on one button at 10:00 and oversold seats a lawsuit, and there is no lock in it. Lesson 015 settled the seat claim as a committed row in one short transaction: the claim is data, nothing is held while the buyer hunts for their card, and a second transaction turns the hold into a sale. The decision and the resource never separated, so there was nothing to coordinate. Lesson 033 ran the same refusal against consensus, and Stagefront escapes a lock by the route it escaped that. Lesson 013 said sharding spreads load across data rather than across time. Here is its lock version. A lock moves a decision away from the data, and the cheapest fix is usually to move it back.

Galewatch is the counterexample where you cannot. Lesson 032's mover held a claim row in one database while deleting rows from three others, which is why a correct lease still cost a year of one turbine's history, and the mover_fence row that would have caught it was proposed in April and has never been built. Not laziness: the resharding is done, so nothing exercises the fence, and a fence nobody has watched refuse a statement is a fence you should assume does not work. The next divisor change is when they find out.

Recap

A lock is a promise about a resource, and almost every implementation makes a promise about a client instead. A mutex works because the lock and the thing it guards share an address space and a scheduler, and the holder cannot outlive the lock. Across machines the lock service cannot see your resource, so all it can tell you is that it has not said yes twice.

Ask why you want it before you ask which one to use. An efficiency lock that fails costs duplicate work; a correctness lock that fails corrupts data. Marlow's sweeper looked like the first and was the second, and what made it the second was one stock + 1.

The lock you need is often the one you already have, with a finer key. (job_name, run_slot) gives seven hundred and twenty runs their own election with no expiry anywhere. It stops two boxes starting the same run, and does nothing about a run that outlives its slot, which are two different problems.

Checking the value before you release stops you deleting somebody else's lock. It does not stop somebody else holding yours. The famous fix for SET NX PX repairs the cheaper of its two bugs, and the expensive one is lesson 032's lease with lesson 034's clock underneath it.

The row's own status is a fence with two values. Guard the write with the state it transitions out of and the second worker matches zero rows, which is lesson 019's first pile doing a fence's job without a token.

A lock without a fenceable resource is a performance optimisation wearing a safety hat. Sort what you are protecting by whether it can keep score. A database can, an object store with a conditional write can, a payment provider never will, and for that one you wanted lesson 019.

A lock moves a decision away from the data, and the cheapest fix is usually to move it back. Stagefront has the hardest contention in this course and no lock anywhere in it, because the claim is a committed row in the database that owns the seat.

Check your understanding

  1. Take the hook's timeline and replace the Redis lock with pg_try_advisory_lock on box A's Postgres, changing nothing else. Say what happens at 21:52:00, what the catalogue says on the Friday, and name the one new way this version can fail that the Redis version could not.

  2. A colleague proposes (job_name, run_slot) for the sweeper and calls the problem solved. Agree with the part that is solved, then write the sequence of events that still oversells a copy, and say what you would add.

  3. You inherit a service that takes a Redis lock, calls a shipping provider's API to book a courier, and releases the lock. Sort the resource, say what a fencing token would be worth here, and give the design you would ship.

  4. Somebody wants Redlock across five nodes for a nightly reconciliation job that writes to one Postgres database. Make the strongest case for them, then make yours, and say which single fact about the job would decide it.

  5. A lock is held through a call that normally takes 300 milliseconds and occasionally takes 30 seconds, and the lease is 10 seconds. Give three ways to size this, say why each is wrong, then say what you would change about the work instead.

  6. Galewatch's mover needs a fence before the next divisor change, and the resharding that would exercise it is over. Describe how you would test the fence without a real migration, and what you would want to see in a log line to believe it works.

Next lesson

036 Exactly Once Is a Lie: Delivery Guarantees in Practice. Today ended on work that is safe to repeat because the lock could not be made safe; next lesson takes the same honesty to a broker that promises to deliver your message exactly one time.

Finished reading?

Marking a lesson done keeps your place on the course index. It is stored only in this browser.

Tip: use the ← and → keys to move between lessons.