Skip to main content
ExplainerDistributed SystemsApache ZooKeeper· 6 min read· in Technology

How Fencing Tokens Prevent Split-Brain Writes in Distributed Locks

A distributed lock cannot guarantee data safety on its own because a paused client can wake up and overwrite a new lock holder's work. By attaching a strictly increasing sequence number to every write, storage systems can automatically reject delayed operations from expired leases.

By Wei Zhang

In short

  • Time-based distributed locks cannot guarantee data safety because paused clients can wake up and write after their leases expire.
  • Fencing tokens solve this by attaching a strictly increasing sequence number to every lock acquisition.
  • The storage system actively checks these tokens, rejecting any write that carries a lower number than the highest one previously processed.

Fencing tokens prevent split-brain data corruption by attaching a strictly increasing sequence number to every lock acquisition, allowing the underlying storage system to reject writes from clients whose leases have expired. When a system experiences a stop-the-world garbage-collection pause, a client can freeze, lose its lock, and wake up believing it still holds exclusive access.[1]

This vulnerability exists in any distributed architecture that relies on time-bound leases for mutual exclusion. A lock server grants a client permission to modify a shared resource for a specific duration, such as 30 seconds.[2]

If the client crashes, the lease expires, preventing the resource from being locked forever. However, time-to-live expirations introduce a fatal flaw when combined with the realities of modern runtime environments.[1]

Applications written in languages like Java or Go periodically halt all execution threads to reclaim memory. These garbage-collection pauses usually last 5 to 50 milliseconds, but under heavy load, they can stretch for 120 seconds or more.[1]

The split-brain write hazard

During a severe pause, the client process is completely suspended and cannot communicate with the network. The lock server observes the silence, assumes the client has failed, and allows the 30-second lease to expire.[3]

The lock is then safely granted to a second client, which begins processing the same task. The catastrophic failure occurs when the first client's garbage-collection pause finally ends.[1]

A garbage collection pause can easily exceed the duration of a standard distributed lock lease.

From the perspective of the suspended process, no time has passed, and it remains entirely unaware that its lease was revoked. It resumes execution exactly where it left off and issues a write command to the database.[1]

Because the second client is also actively writing to the same resource, the system enters a split-brain state. Two independent processes now hold what they both believe to be exclusive access to a single piece of data.[5]

"The lock has a timeout, which is always a good idea, otherwise a crashed client could end up holding a lock forever," distributed systems researcher Martin Kleppmann noted in his foundational 2016 analysis. "However, if the GC pause lasts longer than the lease expiry period, it may go ahead and make some unsafe change."[1]

Enforcing sequence with fencing tokens

The lock service itself is powerless to prevent this collision. By the time the delayed client wakes up and transmits its write request, the lock server is no longer involved in the transaction.[3]

The client communicates directly with the storage layer, which has no native awareness of distributed lock ownership. This architectural blind spot means that adding more robust consensus protocols to the lock server does not solve the underlying race condition.[1]

Whether the lock is managed by a single instance or a highly available cluster, the storage engine remains vulnerable to delayed packets. The protection must be enforced at the destination.[1]

The storage layer uses the fencing token to reject writes from clients whose leases have expired.

To bridge this gap, engineers implement fencing tokens—a strictly monotonic sequence number generated by the lock service. Every time a client successfully acquires or renews a lock, the service returns an integer that is guaranteed to be higher than any previously issued token.[1]

The client is required to append this fencing token to every write request it sends to the shared storage system. The storage engine, in turn, maintains a high-water mark recording the largest token it has successfully processed.[5]

Generating monotonic sequences

This simple integer comparison forms an impenetrable barrier against stale operations. When the first client wakes up from its extended garbage-collection pause, it attempts to write using its original, lower token, such as 33.[1]

The storage system compares this incoming number against the higher token, such as 34, already submitted by the second client. Recognizing the sequence violation, the database outright rejects the delayed write.[1]

This mechanism shifts the ultimate responsibility for data integrity from the volatile lock lease to the durable storage layer. The lock server coordinates the sequence, but the database enforces the physical boundary.[3]

"You should implement fencing tokens," the official Redis documentation states regarding distributed locking patterns. "This is especially important for processes that can take significant time and applies to any distributed locking system."[2]

Not all distributed lock implementations natively support the generation of strictly increasing tokens. Standard Redis setups, for example, require developers to manually script an atomic increment operation alongside the lock acquisition.[2]

The limits of time-based assumptions

Dedicated consensus systems are explicitly designed to provide these monotonic guarantees out of the box. Apache ZooKeeper assigns a 64-bit sequential transaction identifier, known as a zxid, to every state change.[3]

Clients can seamlessly extract this identifier and use it as a cryptographically safe fencing token. Similarly, the etcd key-value store, which underpins Kubernetes cluster coordination, tracks every modification through a 64-bit global revision number.[4]

When a client acquires a distributed lease in etcd, the creation revision of that specific key serves as an automatically incrementing token that satisfies the fencing requirement.[4]

Documented system stalls frequently exceed the default lease timeouts used by naive locking implementations.

Hazelcast implements this pattern directly within its FencedLock API, abstracting the sequence generation away from the application logic. The framework replicates the lock state across a consensus group and automatically increments the token each time the lock transitions.[5]

"FencedLock tracks liveness of lock holders via a session mechanism that works in a unified manner," Hazelcast engineers explained in a technical breakdown of the feature. "It allows third-party systems to participate in the locking protocol."[5]

Implementing storage-side checks

The necessity of fencing tokens exposes a fundamental truth about distributed systems: wall-clock time is an illusion. Relying on a 30-second lease assumes that 30 seconds on the lock server perfectly matches 30 seconds on the client.[1]

In reality, clock drift and CPU scheduling make time highly subjective. A hypervisor managing virtual machines can freeze an entire guest operating system to migrate it to a different physical host.[1]

Network latency introduces identical hazards without requiring a process pause. A client might generate a write request well within its valid lease window, only for the packet to get trapped in a degraded network switch.[1]

In a documented incident at GitHub, network packets were delayed in transit for approximately 90 seconds before being delivered. If a system relied purely on a standard 30-second lock lease, those delayed packets would have silently overwritten newer data.[1]

Illustration: The ultimate responsibility for data integrity rests with the storage layer, not the lock server.

For strict correctness, the storage layer must actively participate in the fencing protocol. Relational databases like PostgreSQL implement this easily by adding a last-token column to the target table, allowing the application to include a conditional clause in its update statement.[1]

If the query affects zero rows, the client immediately knows its token was rejected and its lock has expired. By removing time from the equation and relying entirely on strictly ordered sequences, fencing tokens allow distributed architectures to maintain absolute consistency.[6]

How we did this

Method
Comparison of default lock lease durations against documented extreme system stall events to derive the actual vulnerability window for split-brain corruption.
What we found
Because maximum system stall times mathematically exceed standard lock lease timeouts by a factor of three, relying purely on time-to-live expiration guarantees data corruption under load, making storage-side fencing mandatory.
What we worked from
  • Default distributed lock lease duration: 30 seconds — Redis
  • Documented network/GC stall duration: 90 seconds — Martin Kleppmann
Limits of this analysis
Stall durations vary wildly by runtime environment; highly tuned C++ or Rust applications without garbage collection face narrower vulnerability windows primarily limited to network delays.

Definitions

Fencing Token
A strictly increasing sequence number issued by a lock service to prevent delayed clients from overwriting newer data.
Split-Brain
A catastrophic state where two independent processes simultaneously believe they hold exclusive access to the same shared resource.
Garbage Collection Pause
A temporary halt in application execution while the runtime environment reclaims unused memory.
Time-To-Live (TTL)
The maximum duration a distributed lock lease remains valid before the server automatically revokes it.
Monotonic Sequence
A series of numbers that only ever increases and never repeats, ensuring absolute chronological ordering.

Questions & answers

What triggers a stop-the-world garbage collection pause?

Memory management routines in languages like Java or Go must occasionally halt all application threads to safely remove unused objects. While usually brief, these pauses can stretch for minutes under heavy heap pressure.

Can I use a standard Redis lock for data correctness?

Not without custom scripting. A standard Redis lock provides a time-based lease but does not natively generate the strictly increasing sequence numbers required to prevent split-brain writes.

How does the database know a token is invalid?

The database stores the highest fencing token it has successfully processed. If a delayed client attempts to write using a token lower than this stored high-water mark, the database rejects the operation.

Do network delays cause the same split-brain problem?

Yes. Even if the client process never pauses, a network switch can delay a write packet long enough for the lock lease to expire, resulting in the same data corruption when the packet finally arrives.

Analysis by camp

Distributed Systems Theorists

Argue that time-based leases are fundamentally unsafe for correctness and require monotonic sequence numbers.

Researchers in this camp, notably Martin Kleppmann, emphasize that wall-clock time is an illusion in distributed architectures. They argue that any system relying on a time-to-live expiration is mathematically guaranteed to fail under sufficient load. From this perspective, fencing tokens are not an optional enhancement but a mandatory requirement for any lock that protects persistent data.

Pragmatic Implementers

Acknowledge theoretical risks but argue naive locks are sufficient for efficiency optimizations.

Maintainers of caching layers and in-memory data stores often advocate for a pragmatic approach to locking. They argue that if a lock is only used to prevent two workers from performing the same expensive computation—such as generating a daily report—a rare split-brain collision causes no permanent damage. In these efficiency-driven scenarios, the overhead of implementing storage-side fencing tokens outweighs the benefits.

Consensus Framework Developers

Advocate for building sequence generation directly into the infrastructure layer.

Engineers working on consensus protocols like Apache ZooKeeper and etcd focus on abstracting the complexity of sequence generation away from the application. They design their systems to automatically issue strictly increasing identifiers, such as zxids or revision numbers, with every state change. This camp believes the infrastructure should provide foolproof primitives so developers do not have to manually script atomic counters.

Distributed Systems Theorists 40%Pragmatic Implementers 30%Consensus Framework Developers 30%
Distributed Systems Theorists
Argue that time-based leases are fundamentally unsafe for correctness and require monotonic sequence numbers.
Pragmatic Implementers
Acknowledge theoretical risks but argue naive locks are sufficient for efficiency optimizations.
Consensus Framework Developers
Advocate for building sequence generation directly into the infrastructure layer.

Perspectives this story doesn't cover

  • Database Engine Maintainers
  • Application Developers

Sources

Source coverage

6 outlets

3 viewpoints surfaced

Distributed Systems Theorists 40%Pragmatic Implementers 30%Consensus Framework Developers 30%
  1. [1]Martin KleppmannDistributed Systems Theorists

    How to do distributed locking

    Read on Martin Kleppmann →
  2. [2]RedisPragmatic Implementers

    Distributed Locks with Redis

    Read on Redis →
  3. [3]Apache ZooKeeperConsensus Framework Developers

    ZooKeeper Recipes and Solutions

    Read on Apache ZooKeeper →
  4. [4]etcdConsensus Framework Developers

    etcd API Reference

    Read on etcd →
  5. [5]HazelcastConsensus Framework Developers

    Long Live Distributed Locks

    Read on Hazelcast →
  6. [6]Factlen Editorial Team

    Synthesis by Factlen editorial team

    Read on Factlen Editorial Team →

Comments

Stay informed

Every angle. Every day.

Get Technology stories with full source coverage and perspective breakdowns, free every day.