Search

How Distributed Locks Work (and Why They Often Fail)

The short answer

Quick answer: A distributed lock lets processes on different machines agree that only one of them may do something at a time. It is implemented by having each process try to create a record in a shared store, such as Redis, ZooKeeper or etcd; whoever succeeds holds the lock. Because the holder might crash, the lock must expire automatically. And that is the problem: a holder that is merely paused or slow can have its lock expire and be given to someone else while it still believes it is in charge. A safe design therefore attaches an increasing fencing token to each lock, and the protected resource rejects any request carrying an old token.

Why you might want one

Inside a single process, a mutex stops two threads touching the same data. See processes vs threads. Across machines there is no shared memory, so you need a shared service to arbitrate.

Martin Kleppmann's article How to do distributed locking separates two very different purposes:

PurposeExampleIf the lock fails
EfficiencyAvoid two workers generating the same reportSome wasted work; a duplicate email
CorrectnessPrevent two processes writing the same file or charging the same accountCorrupted data, lost updates, real harm

The first can tolerate an imperfect lock. The second cannot, and needs far more care.

The basic recipe

With Redis, a simple lock looks like this:

SET lock:invoice-42 <random-value> NX PX 30000
  • NX means "only set if the key does not exist". Exactly one client succeeds.
  • PX 30000 sets an expiry of 30 seconds.
  • The random value identifies the owner.

To release, delete the key only if it still holds your value. That check-and-delete must be atomic, so it is done in a small Lua script. Without the check, a client whose lock had expired could delete a lock now held by someone else.

The Redis documentation on distributed locks describes this pattern.

Why the expiry is needed, and why it hurts

If a client takes a lock and then crashes, the lock would be held forever. So every practical distributed lock is really a lease: it is valid for a limited time and then lapses.

But a lease assumes the holder knows when its time is up. Consider:

  1. Client A acquires the lock with a 30-second lease.
  2. Client A pauses for 40 seconds. Causes include a long garbage collection pause, the virtual machine being suspended, heavy swapping, or a slow network call.
  3. The lease expires. Client B acquires the lock and starts working.
  4. Client A wakes up. It has no idea any time passed. It continues and writes to the shared resource.

Both clients now believe they hold the lock, and both write. The lock did not provide mutual exclusion.

No timeout value fixes this. A process cannot detect that it was paused between checking the lock and performing the write. Clock problems make it worse: if the holder's clock or the lock service's clock jumps, the lease may be shorter or longer than intended. See why clocks can't be trusted.

The fix: fencing tokens

The solution is to stop trusting the lock holder's opinion and let the resource enforce exclusion.

  1. Each time the lock service grants the lock, it also returns a fencing token: a number that increases with every grant.
  2. The client includes the token in every write to the protected resource.
  3. The resource remembers the highest token it has seen and rejects any write with a lower one.

Replaying the failure:

  1. Client A gets the lock with token 33, then pauses.
  2. The lease expires. Client B gets the lock with token 34 and writes. The storage records 34.
  3. Client A wakes and tries to write with token 33. The storage rejects it.

The stale client is fenced off. Correctness no longer depends on timing.

This requires the resource to support such a check. A database can do it with a conditional write, such as UPDATE ... WHERE version < 34. Many storage services offer conditional writes or version checks that serve the same purpose.

Single Redis node, and Redlock

A lock in one Redis instance has a further weakness: if that instance fails over to a replica, the lock may not have been replicated yet, because Redis replication is asynchronous. The new primary has no record of it and grants it again.

Redlock is an algorithm proposed by Redis's creator to address this. The client tries to acquire the lock on several independent Redis nodes, typically five, and considers it held if a majority succeed within a time limit.

Its safety has been publicly debated. Kleppmann's critique argues that Redlock depends on assumptions about timing, bounded clock drift, pauses and network delay, that real systems violate, and that it produces no fencing token. The practical summary:

  • For efficiency locks, a single Redis node is simple and sufficient.
  • For correctness locks, use a system built on consensus, and use fencing tokens.

Locks built on consensus

ZooKeeper and etcd replicate their state with a consensus protocol (ZAB and Raft respectively), so an acknowledged lock survives node failures and cannot be granted twice by a split cluster.

ZooKeeper's recipe

  1. Each client creates an ephemeral sequential node under a lock path. ZooKeeper appends an increasing number.
  2. The client with the lowest number holds the lock.
  3. Each other client watches the node just ahead of it and is notified when it disappears.
  4. An ephemeral node is deleted automatically when the client's session ends, so a crashed holder releases the lock.

The sequence number, or the node's transaction ID, can serve as a fencing token.

etcd offers leases and transactions, and each key has a revision number that increases with every change, which works well as a fencing token.

Even these are still leases underneath. A paused client can outlive its session, so fencing remains necessary for correctness.

Database locks

If everything you need to protect lives in one database, you do not need a separate lock service.

  • Row locks: SELECT ... FOR UPDATE inside a transaction.
  • Advisory locks: named locks the database manages for you (PostgreSQL has these built in).
  • Optimistic concurrency: a version column and a conditional update.
  • Unique constraints: let the database reject the second attempt.

These are simpler and safer, because the lock and the data are covered by the same transaction. See how database transactions handle concurrent users.

Do you need a lock at all?

Often the better answer is to remove the need:

  • Make the operation idempotent, so doing it twice is harmless. See idempotency.
  • Partition the work so each item has a single owner: one consumer per queue partition, one worker per shard.
  • Use a single writer elected as leader, with everything else read-only.
  • Use compare-and-set on the resource itself.

A lock that must be perfectly correct across machines is one of the harder things to build. Avoiding it is frequently easier than getting it right.

Checklist if you do use one

  1. Decide whether it is for efficiency or correctness.
  2. Always set an expiry.
  3. Release only if you still own the lock.
  4. For correctness, use fencing tokens checked by the resource.
  5. Keep the protected section short, well under the lease time.
  6. Handle failure to acquire: retry with backoff and jitter, or give up.
  7. Assume that, occasionally, two holders will overlap, and make that survivable.

Frequently asked questions

What is a distributed lock?

A mechanism that lets processes on different machines ensure only one of them performs a particular action at a time, using a shared coordination service.

Is a Redis lock safe?

It is fine for avoiding duplicate work. For protecting correctness, it is not sufficient on its own, because failover and expired leases can allow two holders.

What is a fencing token?

A number that increases each time a lock is granted. The protected resource rejects requests with a token older than the newest it has seen, which blocks stale lock holders.

What is the difference between a lock and a lease?

A lease is a lock with a time limit. Distributed locks are leases in practice, since the holder may crash and never release.

Conclusion

A distributed lock is easy to acquire and hard to trust. Expiry is necessary to recover from crashes, yet expiry is exactly what lets a paused holder overlap with a new one. If duplicated work is merely wasteful, a simple lock will do. If it would corrupt data, combine a consensus-backed lock with fencing tokens, or better, design the lock away.

Related articles

Sources and further reading

Usama Muneer

Usama Muneer

Coder, Blogger, Tech Speaker & Web Technologies Enthusiast. Passionate about working on open-source Programming languages & Tools while utilizing my Product Development skills.

Your experience on this site will be improved by allowing cookies Cookie Policy