Consistent Hashing Explained
How consistent hashing works, why it beats modulo hashing when servers change, how virtual nodes balance load, and where it’s used in caches and databases.
How distributed locks and leases work, why process pauses break naive locks, fencing tokens, Redis, ZooKeeper, etcd and database locks, and alternatives.

A distributed lock lets one process at a time, across many machines, do something that must not happen concurrently: run a scheduled job, process a file, or update a shared resource. Because the holder can crash, locks are granted as leases that expire. The hard part is that a process can pause (for garbage collection, a slow disk or a network stall), outlive its lease and carry on as if it still held it. For correctness, pair leases with fencing tokens, and prefer a consensus-backed store or your database over ad-hoc schemes.
Kleppmann’s widely cited analysis splits locks into two kinds:
Be honest about which you have, because it decides the design.
A lease is a lock with a time limit. If the holder crashes, the lease runs out and someone else can take over, so the system never deadlocks on a dead process. Long-running holders renew the lease periodically.
The catch: the holder can’t be sure its lease is still valid. A 30-second garbage-collection pause, a VM migration or a slow network can mean it resumes after expiry, while another process already holds the lock. Checking the clock just before writing doesn’t help, because the pause can happen right after the check.
The robust fix is a fencing token: every time the lock is granted, the lock service hands out a number that only ever increases. The holder sends the token with every write, and the storage system rejects any write carrying a token lower than one it has already seen.
| Step | Client A | Client B | Storage |
|---|---|---|---|
| 1 | Gets lock, token 33 | ||
| 2 | Pauses (GC) | ||
| 3 | Lease expires | Gets lock, token 34 | |
| 4 | Writes with token 34 | Accepts; highest seen is 34 | |
| 5 | Wakes, writes with token 33 | Rejects: 33 is lower than 34 |
Google’s Chubby lock service described a similar idea, called sequencers. Fencing requires the protected resource to check tokens, which is easy for a database row with a version column and harder for third-party APIs.
If the protected data lives in one relational database, use it: row locks with SELECT … FOR UPDATE, advisory locks such as PostgreSQL’s pg_try_advisory_lock, or a lock table with a unique constraint and an expiry column. Locks and data then share one source of truth. Our guide to transaction isolation levels covers row locking and optimistic concurrency, which often removes the need for a separate lock.
For efficiency locks, a single Redis instance works well. Set a key only if it doesn’t exist, with an expiry and a unique random value:
SET lock:nightly-report token-8f3a NX PX 30000
Release it only if the value still matches, using a small script, so you never delete a lock someone else now holds. Redis also documents Redlock, an algorithm across several independent Redis nodes; its safety depends on timing assumptions that Kleppmann and others have challenged, so don’t rely on it where correctness is at stake without fencing.
These are consensus-based coordination services designed for exactly this problem. Locks are tied to a client session or lease: ZooKeeper uses ephemeral sequential nodes that vanish if the client’s session dies, and etcd attaches keys to leases. Both expose ever-increasing revision or version numbers that can serve as fencing tokens, and both stay consistent through node failures, at the cost of running a quorum cluster.
| Option | Safety | Operational cost | Good for |
|---|---|---|---|
| Database locks | Strong, within one database | Low (you already run it) | Protecting that database’s data |
| Single Redis instance | Weak: fails with the instance | Low | Efficiency locks, deduplicating jobs |
| Redlock | Debated; timing-dependent | Medium | Efficiency locks across nodes |
| ZooKeeper / etcd | Strong (consensus) | Medium to high | Leader election, correctness locks |
For efficiency, usually yes. For correctness, a single instance can lose the lock on failover, and pauses can outlive leases, so add fencing tokens or use a consensus-based system.
A lease is a lock with an expiry. In distributed systems, locks should almost always be leases, so a crashed holder can’t block everyone for ever.
Long enough to cover normal processing with a margin, and short enough that recovery after a crash is acceptable. Renew it for longer jobs rather than setting a huge timeout.
Every article is edited by a human and checked against our editorial policy. Spotted a mistake? Tell us.
How consistent hashing works, why it beats modulo hashing when servers change, how virtual nodes balance load, and where it’s used in caches and databases.
The CAP theorem and its extension PACELC explained: consistency, availability and partition tolerance, what the trade-offs mean in practice and common myths.
How CDNs work: edge locations, request routing, cache keys, Cache-Control headers, invalidation, origin shielding, security features and edge compute.