Skip to content
Rafe Uddaraj

Distributed Locks Under the Hood: Redlock, Fencing Tokens and the Split-Brain Problem

12 min readEnglishRead in Bangla
On this page

Imagine a production system where multiple Node.js servers run simultaneously. Let us say you have a background job running every night that executes generateMonthlyInvoices(). Suddenly, you notice two servers executed the same job at the exact same time, and a customer got billed twice. Or maybe a payment processed twice, or two different workers updated a shared inventory row at once. In these situations, the most obvious solution that comes to an engineer's mind is placing a Redis lock right before the critical section. At first glance, putting a Redis coordination point between Worker A and Worker B seems to solve the problem. But the real question is, what actually happens if Worker A grabs the lock but fails to finish the work before the lock's TTL expires? This single question begins a journey into one of the most complex architectures in distributed systems.

The Real Objective Behind Locking

We first need to understand what a lock actually tries to solve. Defining the engineering requirement is more important here than worrying about the code. Our main goal is to ensure only one valid worker can perform a critical operation at a time. Along with this, if a worker crashes, another worker should eventually be able to take over the job. Also, the system must not enter an unsafe state even if the infrastructure fails. These three properties are known as safety, liveness, and fault tolerance. Safety ensures two workers cannot perform conflicting operations on a protected resource at the same time. Liveness guarantees that if a worker crashes, the entire system does not stay stuck forever. Fault tolerance ensures the system behaves predictably even if a Redis node, network, or application instance fails. A lock does not just control concurrency; it also requires designing for failure recovery.

Because of horizontal scaling, running multiple instances of the same application is perfectly normal. Multiple workers can easily try to process the same job or resource. A database transaction cannot always protect the entire distributed workflow, so a Redis lock initially feels like an easy solution. But if an old worker continues its task after the lock expires, a new worker might grab the lock. The stale operation from the old worker can then cause severe damage to the entire system.

Distributed Lock Contention
Distributed Lock Contention

Mental Model Context

A process staying alive and a lock staying alive are two completely different things. A worker might be alive in memory and keep an active network connection, but its lock authority could have expired a long time ago.

From Locks to Time-Bound Leases

To understand this simply, we can use the analogy of a single access card for a warehouse. Imagine a special room inside a warehouse where only one operator can enter at a time. An operator takes an access card and starts working inside, but the security system states this card will automatically become invalid after thirty minutes. If the operator works for forty-five minutes, a second operator will get a new card after thirty minutes and enter the room. The first operator is still inside working. The security system never guaranteed that only one person would be in the room. It only guaranteed that the first operator's access would remain valid for a specific timeframe. This exact concept is the origin of the lease.

When we run the SET resource:123 worker-token NX PX 30000 command in Redis, NX ensures the key only sets if it does not already exist, and PX 30000 ensures the key expires after thirty seconds. As a reader, you might think this is the perfect solution because one person taking the lock prevents others from entering. But if Worker A crashes after taking the lock and the lock has no expiration, that lock stays there forever. Worker B will never be able to acquire the resource. This is exactly why we need a TTL or Time To Live. Using a long TTL makes crash recovery very slow, while a short TTL creates the risk of the work outliving the lease. The TTL actually gives the lock a time-bound ownership, which we call a lease. A lease means you have authority over this resource, but it is not indefinite. It is only valid for a specific duration. From the worker's perspective, it feels like acquiring a lock. However, from the system's perspective, that authority is strictly limited to the lease expiration time.

The Race Condition in Lock Release

After understanding the lease concept, we need to discuss the internal race condition of the lock release mechanism. Imagine Worker A acquires a lock and then pauses for a long time for some reason. Thirty seconds later, the lock expires, and Worker B acquires a new lock. Right at this moment, Worker A resumes and directly executes the DEL resource:123 command. A catastrophic event happens here. Worker A is not actually deleting its own lock; it is deleting Worker B's brand new lock. As a result, Worker B's lock disappears, and the system enters a completely unsafe state.

Lock Release Race Condition
Lock Release Race Condition

To avoid this type of production bug, you must use a unique ownership token or a random lock token with the lock. A random string is set as the lock value. When releasing it, the logic checks if the current lock value matches the worker's token. If they match, the lock is deleted. Otherwise, the operation is rejected. We usually use a Lua script or an atomic conditional operation to handle this release process safely:

Lua
if redis.call("get", KEYS[1]) == ARGV[1] then
return redis.call("del", KEYS[1])
else
return 0
end

But you need to keep a very subtle difference in mind here. A random lock token only identifies ownership. It cannot prevent a stale worker from executing.

Dangerous Assumption

A random lock token only guarantees a safe unlock. It never guarantees that your worker's execution context is still valid.

Redis Failover and the Split-Brain Problem

Up to this point, our architecture relied on a single Redis instance. However, Redis itself can fail in a production system. Let us assume your architecture has a Redis primary and a Redis replica. Worker A takes a lock from the primary node, but the primary node crashes immediately after. Before the lock data synchronizes or replicates to the replica, the replica node promotes itself to the new primary. In this state, Worker B sends a request to the new primary node and successfully acquires the lock. Now Worker A thinks it owns the lock, and Worker B also thinks it owns the lock. This is the most dangerous failure mode in distributed locking, known as the split-brain problem.

Redis Failover Split Brain
Redis Failover Split Brain

We need to establish a critical engineering principle here. High availability does not mean correctness. A Redis replica failover certainly helps keep your system available, but if the lock state is not perfectly preserved, this very process of increasing availability can break your system's safety. We actually want a locking mechanism where, even if the infrastructure fails, two workers can never simultaneously become valid owners.

Redlock and the Illusion of Time

The Redlock algorithm was introduced to handle the failure of a single Redis node. Redlock essentially uses multiple independent Redis instances instead of just one. Imagine you have five independent Redis nodes. A worker will try to acquire a lock on all five of these nodes. It will only consider the lock valid if it successfully acquires it on a majority, or at least three nodes. If two workers try to get the lock at the same time, their majority sets will always share at least one common node. This allows the distributed quorum concept to work. But getting a majority is not enough on its own. The time spent on the acquisition is also incredibly important.

According to the Redlock acquisition flow, the worker first records a start timestamp and tries to lock the five nodes one by one. After getting a majority, it checks how much total time it took to acquire the lock. The lock is only considered acquired if there is still validity remaining after subtracting the elapsed time. But this is exactly where the biggest hidden assumption of Redlock lies, and that is time. Redlock's validity relies entirely on time measurement. In a real production environment, countless things can happen, such as network delays, clock drift, process scheduling, garbage collection pauses, CPU starvation, VM pauses, container throttling, or network partitions. A worker might think its lock is still valid from its own perspective, but from the perspective of other parts of the distributed system, that lease might have expired a long time ago.

Redlock Vulnerability

Majority intersection alone cannot solve the entire correctness problem. Timing and network delays are always unpredictable in distributed systems.

The Most Dangerous Failure: Paused Worker

The most central failure scenario of the entire distributed locking architecture is the paused worker. Imagine Worker A successfully acquires a lock and starts working. Suddenly, Worker A suffers a long garbage collection pause, CPU starvation, VM suspension, or network stall and pauses completely. During this pause, its lock lease expires. Following the system rules, Worker B acquires a new lock and successfully finishes its job. A little while later, Worker A finishes its pause and resumes. Worker A has no idea that it was paused. It still thinks it owns the lock and tries to write its old operation to the database.

Paused Worker Timeline
Paused Worker Timeline

The real problem here is that nobody stopped the old worker. A lock expiring does not mean Worker A's process will automatically stop. A lock release or expiration only changes the system's state; it cannot physically kill a running code execution. This is a fundamental truth of distributed systems.

Fencing Tokens: The Ultimate Correctness Guarantee

We need to introduce the fencing token concept to solve this terrifying problem of the paused worker. When Worker A takes a lock, it does not just receive a lock. The lock service also gives it a monotonically increasing number, like Token 33. A while later, when the lease expires and Worker B takes a new lock, it gets Token 34. The resource itself now knows that the last accepted token it received was 34. If Worker A comes back from its pause and tries to write using Token 33, the resource will see that 33 is smaller than the current 34. It will reject the operation immediately. On the other hand, Worker B's token is 34, so the system will accept it without any issue.

Complete Failure Timeline
Complete Failure Timeline

The core idea here is that even if the old worker's process remains alive, its authority has become obsolete. It is very important to clearly understand the difference between a random lock token and a fencing token. A random lock token can only provide uniqueness, but it cannot provide ordering. A fencing token has an ordering that proves which owner is newer. The random lock token answers who owns this lock, while the fencing token answers which owner is newer.

Fencing Rule

Generating the token is a coordination problem itself. You have to create this monotonically increasing token using a database sequence, a ZooKeeper version, or a consensus-backed counter.

Database-Level Fencing Enforcement

To utilize the fencing token correctly, you must avoid a common architectural mistake. If the worker takes the token from the lock service and writes directly to the database, and the database knows nothing about the token, then the fencing token is completely useless. The correct architecture requires the database or the resource itself to reject stale operations.

Think of a practical PostgreSQL example. We can create a table with columns named id, value, and last_fencing_token. When the worker writes, it will also send its token. The query will look something like this:

SQL
UPDATE resources
SET
value = $1,
last_fencing_token = $2
WHERE id = $3
AND last_fencing_token < $2;

Here, the database itself checks if the incoming token is larger than the current token. If Worker A brings Token 103, but the database's current token is 104, the condition will not match and the update will fail. The database update must be designed so that even if a stale worker creates a bug at the application level, the database level still rejects it. You should always enforce correctness-critical validation at the authoritative resource.

Fencing Token Database
Fencing Token Database

Architectural Recap and Production Checklist

Understanding the trade-off between safety and availability during a network partition or node failure is the job of an elite architect. A distributed lock is never just a "lock acquired" boolean; it also involves a failure model. For long-running operations, workers can extend the lease with a heartbeat. However, while lease renewal can reduce premature expirations, it cannot eliminate the stale worker problem. It is highly normal for failed requests to retry in distributed systems, so you have to design locks alongside fencing, idempotency, and atomic resource updates.

Using Redis locks everywhere is not a good practice. Native database mechanisms like unique constraints, database transactions, SELECT ... FOR UPDATE, advisory locks, or optimistic concurrency control can often provide better coordination. If your only goal is to reduce duplicate work, a simple Redis lock is enough. But if correctness is critical, fencing is absolutely necessary.

You must review the following checklist before designing any distributed lock in production:

  • Which exact resource are you protecting?
  • What happens if the worker crashes or the lease expires while work is ongoing?
  • Can a stale worker still execute and write data?
  • Who is generating the fencing token, and is the resource validating it?
  • How will the system behave during a network partition, Redis failover, process pause, or client retry?
  • Can the operation be made idempotent, and can the database solve this problem on its own?

A lock dictates who can enter a critical section. A lease dictates how long that permission is valid. Redlock tries to establish a distributed lease using multiple nodes. Finally, a fencing token ensures that an old stale worker can never corrupt the resource under any circumstances. A lock does not stop a paused process, a TTL does not kill a stale worker, and a random lock ID does not establish ordering. Being able to identify and separate these layers in production distributed systems is the most important engineering skill.

All articles

Engineering Roadmap: Why 'Exactly Once' Delivery Is Mostly a Lie

A complete engineering guide on why exactly-once delivery is mathematically impossible at the network level, and how to build effectively-once architectures in production using idempotency, transactional outbox, and inbox patterns.

System DesignENBN

Get in touch

Questions about a video, an article, or working together.