How to do distributed locking
Published by Martin Kleppmann on 08 Feb 2016.
As part of the research for my book, I came across an algorithm called Redlock on the Redis website. The algorithm claims to implement fault-tolerant distributed locks (or rather, leases [1]) on top of Redis, and the page asks for feedback from people who are into distributed systems. The algorithm instinctively set off some alarm bells in the back of my mind, so I spent a bit of time thinking about it and writing up these notes.
Since there are already over 10 independent implementations of Redlock and we don’t know who is already relying on this algorithm, I thought it would be worth sharing my notes publicly. I won’t go into other aspects of Redis, some of which have already been critiqued elsewhere.
Before I go into the details of Redlock, let me say that I quite like Redis, and I have successfully used it in production in the past. I think it’s a good fit in situations where you want to share some transient, approximate, fast-changing data between servers, and where it’s not a big deal if you occasionally lose that data for whatever reason. For example, a good use case is maintaining request counters per IP address (for rate limiting purposes) and sets of distinct IP addresses per user ID (for abuse detection).
However, Redis has been gradually making inroads into areas of data management where there are stronger consistency and durability expectations – which worries me, because this is not what Redis is designed for. Arguably, distributed locking is one of those areas. Let’s examine it in some more detail.
What are you using that lock for?
The purpose of a lock is to ensure that among several nodes that might try to do the same piece of work, only one actually does it (at least only one at a time). That work might be to write some data to a shared storage system, to perform some computation, to call some external API, or suchlike. At a high level, there are two reasons why you might want a lock in a distributed application: for efficiency or for correctness [2]. To distinguish these cases, you can ask what would happen if the lock failed:
- Efficiency: Taking a lock saves you from unnecessarily doing the same work twice (e.g. some expensive computation). If the lock fails and two nodes end up doing the same piece of work, the result is a minor increase in cost (you end up paying 5 cents more to AWS than you otherwise would have) or a minor inconvenience (e.g. a user ends up getting the same email notification twice).
- Correctness: Taking a lock prevents concurrent processes from stepping on each others’ toes and messing up the state of your system. If the lock fails and two nodes concurrently work on the same piece of data, the result is a corrupted file, data loss, permanent inconsistency, the wrong dose of a drug administered to a patient, or some other serious problem.
Both are valid cases for wanting a lock, but you need to be very clear about which one of the two you are dealing with.
I will argue that if you are using locks merely for efficiency purposes, it is unnecessary to incur the cost and complexity of Redlock, running 5 Redis servers and checking for a majority to acquire your lock. You are better off just using a single Redis instance, perhaps with asynchronous replication to a secondary instance in case the primary crashes.
If you use a single Redis instance, of course you will drop some locks if the power suddenly goes out on your Redis node, or something else goes wrong. But if you’re only using the locks as an efficiency optimization, and the crashes don’t happen too often, that’s no big deal. This “no big deal” scenario is where Redis shines. At least if you’re relying on a single Redis instance, it is clear to everyone who looks at the system that the locks are approximate, and only to be used for non-critical purposes.
On the other hand, the Redlock algorithm, with its 5 replicas and majority voting, looks at first glance as though it is suitable for situations in which your locking is important for correctness. I will argue in the following sections that it is not suitable for that purpose. For the rest of this article we will assume that your locks are important for correctness, and that it is a serious bug if two different nodes concurrently believe that they are holding the same lock.
Protecting a resource with a lock
Let’s leave the particulars of Redlock aside for a moment, and discuss how a distributed lock is used in general (independent of the particular locking algorithm used). It’s important to remember that a lock in a distributed system is not like a mutex in a multi-threaded application. It’s a more complicated beast, due to the problem that different nodes and the network can all fail independently in various ways.
For example, say you have an application in which a client needs to update a file in shared storage (e.g. HDFS or S3). A client first acquires the lock, then reads the file, makes some changes, writes the modified file back, and finally releases the lock. The lock prevents two clients from performing this read-modify-write cycle concurrently, which would result in lost updates. The code might look something like this:
// THIS CODE IS BROKEN
function writeData(filename, data) {
var lock = lockService.acquireLock(filename);
if (!lock) {
throw 'Failed to acquire lock';
}
try {
var file = storage.readFile(filename);
var updated = updateContents(file, data);
storage.writeFile(filename, updated);
} finally {
lock.release();
}
}
Unfortunately, even if you have a perfect lock service, the code above is broken. The following diagram shows how you can end up with corrupted data:

In this example, the client that acquired the lock is paused for an extended period of time while holding the lock – for example because the garbage collector (GC) kicked in. The lock has a timeout (i.e. it is a lease), which is always a good idea (otherwise a crashed client could end up holding a lock forever and never releasing it). However, if the GC pause lasts longer than the lease expiry period, and the client doesn’t realise that it has expired, it may go ahead and make some unsafe change.
This bug is not theoretical: HBase used to have this problem [3,4]. Normally, GC pauses are quite short, but “stop-the-world” GC pauses have sometimes been known to last for several minutes [5] – certainly long enough for a lease to expire. Even so-called “concurrent” garbage collectors like the HotSpot JVM’s CMS cannot fully run in parallel with the application code – even they need to stop the world from time to time [6].
You cannot fix this problem by inserting a check on the lock expiry just before writing back to storage. Remember that GC can pause a running thread at any point, including the point that is maximally inconvenient for you (between the last check and the write operation).
And if you’re feeling smug because your programming language runtime doesn’t have long GC pauses, there are many other reasons why your process might get paused. Maybe your process tried to read an address that is not yet loaded into memory, so it gets a page fault and is paused until the page is loaded from disk. Maybe your disk is actually EBS, and so reading a variable unwittingly turned into a synchronous network request over Amazon’s congested network. Maybe there are many other processes contending for CPU, and you hit a black node in your scheduler tree. Maybe someone accidentally sent SIGSTOP to the process. Whatever. Your processes will get paused.
If you still don’t believe me about process pauses, then consider instead that the file-writing request may get delayed in the network before reaching the storage service. Packet networks such as Ethernet and IP may delay packets arbitrarily, and they do [7]: in a famous incident at GitHub, packets were delayed in the network for approximately 90 seconds [8]. This means that an application process may send a write request, and it may reach the storage server a minute later when the lease has already expired.
Even in well-managed networks, this kind of thing can happen. You simply cannot make any assumptions about timing, which is why the code above is fundamentally unsafe, no matter what lock service you use.
Making the lock safe with fencing
The fix for this problem is actually pretty simple: you need to include a fencing token with every write request to the storage service. In this context, a fencing token is simply a number that increases (e.g. incremented by the lock service) every time a client acquires the lock. This is illustrated in the following diagram:

Client 1 acquires the lease and gets a token of 33, but then it goes into a long pause and the lease expires. Client 2 acquires the lease, gets a token of 34 (the number always increases), and then sends its write to the storage service, including the token of 34. Later, client 1 comes back to life and sends its write to the storage service, including its token value 33. However, the storage server remembers that it has already processed a write with a higher token number (34), and so it rejects the request with token 33.
Note this requires the storage server to take an active role in checking tokens, and rejecting any
writes on which the token has gone backwards. But this is not particularly hard, once you know the
trick. And provided that the lock service generates strictly monotonically increasing tokens, this
makes the lock safe. For example, if you are using ZooKeeper as lock service, you can use the zxid
or the znode version number as fencing token, and you’re in good shape [3].
However, this leads us to the first big problem with Redlock: it does not have any facility for generating fencing tokens. The algorithm does not produce any number that is guaranteed to increase every time a client acquires a lock. This means that even if the algorithm were otherwise perfect, it would not be safe to use, because you cannot prevent the race condition between clients in the case where one client is paused or its packets are delayed.
And it’s not obvious to me how one would change the Redlock algorithm to start generating fencing tokens. The unique random value it uses does not provide the required monotonicity. Simply keeping a counter on one Redis node would not be sufficient, because that node may fail. Keeping counters on several nodes would mean they would go out of sync. It’s likely that you would need a consensus algorithm just to generate the fencing tokens. (If only incrementing a counter was simple.)
Using time to solve consensus
The fact that Redlock fails to generate fencing tokens should already be sufficient reason not to use it in situations where correctness depends on the lock. But there are some further problems that are worth discussing.
In the academic literature, the most practical system model for this kind of algorithm is the asynchronous model with unreliable failure detectors [9]. In plain English, this means that the algorithms make no assumptions about timing: processes may pause for arbitrary lengths of time, packets may be arbitrarily delayed in the network, and clocks may be arbitrarily wrong – and the algorithm is nevertheless expected to do the right thing. Given what we discussed above, these are very reasonable assumptions.
The only purpose for which algorithms may use clocks is to generate timeouts, to avoid waiting forever if a node is down. But timeouts do not have to be accurate: just because a request times out, that doesn’t mean that the other node is definitely down – it could just as well be that there is a large delay in the network, or that your local clock is wrong. When used as a failure detector, timeouts are just a guess that something is wrong. (If they could, distributed algorithms would do without clocks entirely, but then consensus becomes impossible [10]. Acquiring a lock is like a compare-and-set operation, which requires consensus [11].)
Note that Redis uses gettimeofday, not a monotonic clock, to
determine the expiry of keys. The man page for gettimeofday explicitly\
says that the time it returns is subject to discontinuous jumps in system time –
that is, it might suddenly jump forwards by a few minutes, or even jump back in time (e.g. if the
clock is stepped by NTP because it differs from a NTP server by too much, or if the
clock is manually adjusted by an administrator). Thus, if the system clock is doing weird things, it
could easily happen that the expiry of a key in Redis is much faster or much slower than expected.
For algorithms in the asynchronous model this is not a big problem: these algorithms generally ensure that their safety properties always hold, without making any timing\ assumptions [12]. Only liveness properties depend on timeouts or some other failure detector. In plain English, this means that even if the timings in the system are all over the place (processes pausing, networks delaying, clocks jumping forwards and backwards), the performance of an algorithm might go to hell, but the algorithm will never make an incorrect decision.
However, Redlock is not like this. Its safety depends on a lot of timing assumptions: it assumes that all Redis nodes hold keys for approximately the right length of time before expiring; that the network delay is small compared to the expiry duration; and that process pauses are much shorter than the expiry duration.
Breaking Redlock with bad timings
Let’s look at some examples to demonstrate Redlock’s reliance on timing assumptions. Say the system has five Redis nodes (A, B, C, D and E), and two clients (1 and 2). What happens if a clock on one of the Redis nodes jumps forward?
- Client 1 acquires lock on nodes A, B, C. Due to a network issue, D and E cannot be reached.
- The clock on node C jumps forward, causing the lock to expire.
- Client 2 acquires lock on nodes C, D, E. Due to a network issue, A and B cannot be reached.
- Clients 1 and 2 now both believe they hold the lock.
A similar issue could happen if C crashes before persisting the lock to disk, and immediately restarts. For this reason, the Redlock documentation recommends delaying restarts of crashed nodes for at least the time-to-live of the longest-lived lock. But this restart delay again relies on a reasonably accurate measurement of time, and would fail if the clock jumps.
Okay, so maybe you think that a clock jump is unrealistic, because you’re very confident in having correctly configured NTP to only ever slew the clock. In that case, let’s look at an example of how a process pause may cause the algorithm to fail:
- Client 1 requests lock on nodes A, B, C, D, E.
- While the responses to client 1 are in flight, client 1 goes into stop-the-world GC.
- Locks expire on all Redis nodes.
- Client 2 acquires lock on nodes A, B, C, D, E.
- Client 1 finishes GC, and receives the responses from Redis nodes indicating that it successfully acquired the lock (they were held in client 1’s kernel network buffers while the process was paused).
- Clients 1 and 2 now both believe they hold the lock.
Note that even though Redis is written in C, and thus doesn’t have GC, that doesn’t help us here: any system in which the clients may experience a GC pause has this problem. You can only make this safe by preventing client 1 from performing any operations under the lock after client 2 has acquired the lock, for example using the fencing approach above.
A long network delay can produce the same effect as the process pause. It perhaps depends on your TCP user timeout – if you make the timeout significantly shorter than the Redis TTL, perhaps the delayed network packets would be ignored, but we’d have to look in detail at the TCP implementation to be sure. Also, with the timeout we’re back down to accuracy of time measurement again!
The synchrony assumptions of Redlock
These examples show that Redlock works correctly only if you assume a synchronous system model – that is, a system with the following properties:
- bounded network delay (you can guarantee that packets always arrive within some guaranteed maximum delay),
- bounded process pauses (in other words, hard real-time constraints, which you typically only find in car airbag systems and suchlike), and
- bounded clock error (cross your fingers that you don’t get your time from a bad NTP\ server).
Note that a synchronous model does not mean exactly synchronised clocks: it means you are assuming a known, fixed upper bound on network delay, pauses and clock drift [12]. Redlock assumes that delays, pauses and drift are all small relative to the time-to-live of a lock; if the timing issues become as large as the time-to-live, the algorithm fails.
In a reasonably well-behaved datacenter environment, the timing assumptions will be satisfied most of the time – this is known as a partially synchronous system [12]. But is that good enough? As soon as those timing assumptions are broken, Redlock may violate its safety properties, e.g. granting a lease to one client before another has expired. If you’re depending on your lock for correctness, “most of the time” is not enough – you need it to always be correct.
There is plenty of evidence that it is not safe to assume a synchronous system model for most practical system environments [7,8]. Keep reminding yourself of the GitHub incident with the 90-second packet delay. It is unlikely that Redlock would survive a Jepsen test.
On the other hand, a consensus algorithm designed for a partially synchronous system model (or asynchronous model with failure detector) actually has a chance of working. Raft, Viewstamped Replication, Zab and Paxos all fall in this category. Such an algorithm must let go of all timing assumptions. That’s hard: it’s so tempting to assume networks, processes and clocks are more reliable than they really are. But in the messy reality of distributed systems, you have to be very careful with your assumptions.
Conclusion
I think the Redlock algorithm is a poor choice because it is “neither fish nor fowl”: it is unnecessarily heavyweight and expensive for efficiency-optimization locks, but it is not sufficiently safe for situations in which correctness depends on the lock.
In particular, the algorithm makes dangerous assumptions about timing and system clocks (essentially assuming a synchronous system with bounded network delay and bounded execution time for operations), and it violates safety properties if those assumptions are not met. Moreover, it lacks a facility for generating fencing tokens (which protect a system against long delays in the network or in paused processes).
If you need locks only on a best-effort basis (as an efficiency optimization, not for correctness), I would recommend sticking with the straightforward single-node locking algorithm for Redis (conditional set-if-not-exists to obtain a lock, atomic delete-if-value-matches to release a lock), and documenting very clearly in your code that the locks are only approximate and may occasionally fail. Don’t bother with setting up a cluster of five Redis nodes.
On the other hand, if you need locks for correctness, please don’t use Redlock. Instead, please use a proper consensus system such as ZooKeeper, probably via one of the Curator recipes that implements a lock. (At the very least, use a database with reasonable transactional\ guarantees.) And please enforce use of fencing tokens on all resource accesses under the lock.
As I said at the beginning, Redis is an excellent tool if you use it correctly. None of the above diminishes the usefulness of Redis for its intended purposes. Salvatore has been very dedicated to the project for years, and its success is well deserved. But every tool has limitations, and it is important to know them and to plan accordingly.
If you want to learn more, I explain this topic in greater detail in chapters 8 and 9 of my\ book, now available in Early Release from O’Reilly. (The diagrams above are taken from my book.) For learning how to use ZooKeeper, I recommend Junqueira and Reed’s book [3]. For a good introduction to the theory of distributed systems, I recommend Cachin, Guerraoui and\ Rodrigues’ textbook [13].
Thank you to Kyle Kingsbury, Camille Fournier, Flavio Junqueira, and Salvatore Sanfilippo for reviewing a draft of this article. Any errors are mine, of course.
Update 9 Feb 2016: Salvatore, the original author of Redlock, has posted a rebuttal to this article (see also HN discussion). He makes some good points, but I stand by my conclusions. I may elaborate in a follow-up post if I have time, but please form your own opinions – and please consult the references below, many of which have received rigorous academic peer review (unlike either of our blog posts).
References
[1] Cary G Gray and David R Cheriton: “ Leases: An Efficient Fault-Tolerant Mechanism for Distributed File Cache Consistency,” at 12th ACM Symposium on Operating Systems Principles (SOSP), December 1989. doi:10.1145/74850.74870
[2] Mike Burrows: “ The Chubby lock service for loosely-coupled distributed systems,” at 7th USENIX Symposium on Operating System Design and Implementation (OSDI), November 2006.
[3] Flavio P Junqueira and Benjamin Reed: ZooKeeper: Distributed Process Coordination. O’Reilly Media, November 2013. ISBN: 978-1-4493-6130-3
[4] Enis Söztutar: “ HBase and HDFS: Understanding filesystem usage in HBase,” at HBaseCon, June 2013.
[5] Todd Lipcon: “ Avoiding Full GCs in Apache HBase with MemStore-Local Allocation Buffers: Part 1,” blog.cloudera.com, 24 February 2011.
[6] Martin Thompson: “ Java Garbage Collection Distilled,” mechanical-sympathy.blogspot.co.uk, 16 July 2013.
[7] Peter Bailis and Kyle Kingsbury: “ The Network is Reliable,” ACM Queue, volume 12, number 7, July 2014. doi:10.1145/2639988.2639988
[8] Mark Imbriaco: “ Downtime last Saturday,” github.com, 26 December 2012.
[9] Tushar Deepak Chandra and Sam Toueg: “ Unreliable Failure Detectors for Reliable Distributed Systems,” Journal of the ACM, volume 43, number 2, pages 225–267, March 1996. doi:10.1145/226643.226647
[10] Michael J Fischer, Nancy Lynch, and Michael S Paterson: “ Impossibility of Distributed Consensus with One Faulty Process,” Journal of the ACM, volume 32, number 2, pages 374–382, April 1985. doi:10.1145/3149.214121
[11] Maurice P Herlihy: “ Wait-Free Synchronization,” ACM Transactions on Programming Languages and Systems, volume 13, number 1, pages 124–149, January 1991. doi:10.1145/114005.102808
[12] Cynthia Dwork, Nancy Lynch, and Larry Stockmeyer: “ Consensus in the Presence of Partial Synchrony,” Journal of the ACM, volume 35, number 2, pages 288–323, April 1988. doi:10.1145/42282.42283
[13] Christian Cachin, Rachid Guerraoui, and Luís Rodrigues: Introduction to Reliable and Secure Distributed Programming, Second Edition. Springer, February 2011. ISBN: 978-3-642-15259-7, doi:10.1007/978-3-642-15260-3
If you found this post useful, please support me on Patreon so that I can write more like it!
To get notified when I write something new, follow me on Bluesky or Mastodon, or enter your email address:
Martin Kleppmann’s Blog | Substack
Martin Kleppmann’s Blog
Distributed systems, databases, and information security
Over 7,000 subscribers
Subscribe
By subscribing you agree to Substack's Terms of Use, our Privacy Policy and our Information collection notice
I won't give your address to anyone else, won't send you any spam, and you can unsubscribe at any time.
tempest.services.disqus.com
tempest.services.disqus.com is blocked
This page has been blocked by an extension
- Try disabling your extensions.
ERR_BLOCKED_BY_CLIENT
Reload
This page has been blocked by an extension
Disqus Comments
We were unable to load Disqus. If you are a moderator please see our troubleshooting guide.
G
Join the discussion…
Comment
Log in with
or sign up with Disqus or pick a name
Disqus is a discussion network
- Don't be a jerk or do anything illegal. Everything is easier that way.
Read full terms and conditions
This comment platform is hosted by Disqus, Inc. I authorize Disqus and its affiliates to:
- Use, sell, and share my information to enable me to use its comment services and for marketing purposes, including cross-context behavioral advertising, as described in our Terms of Service and Privacy Policy, including supplementing that information with other data about me, such as my browsing and location data.
- Contact me or enable others to contact me by email with offers for goods or services
- Process any sensitive personal information that I submit in a comment. See our Privacy Policy for more information
Acknowledge I am 18 or older
I'd rather post as a guest
-
Discussion Favorited!
Favoriting means this is a discussion worth sharing. It gets shared to your followers' Disqus feeds, and gives the creator kudos!
- Tweet this discussion
- Share this discussion on Facebook
- Share this discussion via email
-
Copy link to discussion
- Newest
- +
- Flag as inappropriate
N
Hi Martin,
Thank you for your post.
The fencing algorithm would seem to allow Client 1 to write without holding the lease if it wrote with token 33 before Client 2 wrote with token 34 (I made a crude diagram below). I'm interested in enforcing that the client with the expired lease is not allowed to write to storage at all even if it writes before the client holding the lease.
1. Are there any algorithms to enforce this?
2. Are there any libraries or storage solutions that implement such an algorithm?
Thank you for your time.
see more
H
I think the lock service and the storage have to be combined as one single service to resolve this tricky race condition.
Another workaround is the client sends back the hash of the whole original file before its modifications to the storage system. Then the storage system compares the hash with the current file hash to determine if the file has been modified between the client reading the file and sending a modified version to the storage.
see more
- - [−](https://disqus.com/embed/comments/?base=default&f=martinkl&t_i=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_u=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_d=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&t_t=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&s_o=default# "Collapse")
- [+](https://disqus.com/embed/comments/?base=default&f=martinkl&t_i=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_u=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_d=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&t_t=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&s_o=default# "Expand")
- [Flag as inappropriate](https://disqus.com/embed/comments/?base=default&f=martinkl&t_i=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_u=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_d=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&t_t=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&s_o=default# "Flag as inappropriate")
A
> Another workaround is the client sends back the hash of the whole original file before its modifications to the storage system.
Or just simply a (row)version - like in the optimistic concurrency check.
see more
- - [−](https://disqus.com/embed/comments/?base=default&f=martinkl&t_i=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_u=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_d=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&t_t=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&s_o=default# "Collapse")
- [+](https://disqus.com/embed/comments/?base=default&f=martinkl&t_i=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_u=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_d=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&t_t=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&s_o=default# "Expand")
- [Flag as inappropriate](https://disqus.com/embed/comments/?base=default&f=martinkl&t_i=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_u=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_d=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&t_t=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&s_o=default# "Flag as inappropriate")
K
the storage service is already validating the token somehow with lock service, so on top of that, additionally checking the TTL of current lock shouldn't be too hard. The latency is already there.
see more
![]()
Hi Neil, good question. An obvious thing you could try is that the storage service can make a request to the lock service to check whether a token is still current, before accepting a write, although that will add latency. Other than that I don't really have a good answer. It might be a fundamental limitation because you have a race condition between the two clients, and in the absence of further information (such as a request to the lock service, or assumptions about clock synchronisation) the storage service cannot know whether a lease is still valid.
see more
- - [−](https://disqus.com/embed/comments/?base=default&f=martinkl&t_i=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_u=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_d=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&t_t=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&s_o=default# "Collapse")
- [+](https://disqus.com/embed/comments/?base=default&f=martinkl&t_i=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_u=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_d=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&t_t=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&s_o=default# "Expand")
- [Flag as inappropriate](https://disqus.com/embed/comments/?base=default&f=martinkl&t_i=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_u=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_d=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&t_t=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&s_o=default# "Flag as inappropriate")
K
actually, validating the access token before the storage accepts the request is still not 100% reliable. Consider the case where the lock happen to expire and then be acquired by another client just after the validation passes. due to timing issue, the storage can then take concurrent writes from two clients at the same time.
This is because "validate the token" and "take the write" are not atomic operations and doesn't happen instantaneously.
see more
N
I want to answer the two questions I asked earlier now that I have a better understanding of distributed systems.
"I'm interested in enforcing that the client with the expired lease is not allowed to write to storage at all even if it writes before the client holding the lease.
1. Are there any algorithms to enforce this?"
If you need to enforce this for correctness then you should use your database to do it. What you want is to transactionally acquire a lock and write to the database and then release the lock. Databases already have this functionality built in! Putting the lock in another service makes everything harder because of the problems outlined in this blog post.
I would ask why do you have this requirement? Is it really about the client or about some other property you want to enforce. It is tempting to lock more than you need to enforce a constraint, but what is the minimum you need to lock? Can you structure your database transaction and table(s) such that the database enforces this in an efficient manner? The answer is almost always yes.
"2. Are there any libraries or storage solutions that implement such an algorithm?"
I would recommend you read about the locking strategy your database uses. The database already acquires locks by itself when you use it. It does a good job locking only what is needed
Postgres h ttps://www.postgresql.org/docs/7.1/locking-tables.html
MySQL https://dev.mysql.com/doc/refman/8.4/en/innodb-locking.html
see more
- - [−](https://disqus.com/embed/comments/?base=default&f=martinkl&t_i=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_u=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_d=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&t_t=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&s_o=default# "Collapse")
- [+](https://disqus.com/embed/comments/?base=default&f=martinkl&t_i=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_u=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_d=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&t_t=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&s_o=default# "Expand")
- [Flag as inappropriate](https://disqus.com/embed/comments/?base=default&f=martinkl&t_i=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_u=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_d=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&t_t=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&s_o=default# "Flag as inappropriate")
V
This is a low level algorithm to actually implement transactions properly by databases/filesystems. Application developers should not be using distributed lock/lease directly.
see more
V
Lock/lease expiration in this particular case is a red-herring non-problem that does not cause data correctness issues. To implement transactions with strict linearizability with a WAL (append only write ahead log), you just need to make sure all writes are linearized without conflict, so the fencing token at a WAL subsystem would suffice.
see more
S
What if the Lock Service also responds with lock start time and expiration time to client 1. Client 1 in turn send it to storage system. Storage system just validates expiry time in future or not. Assumping absense of clock synchronization issue.
see more
- - [−](https://disqus.com/embed/comments/?base=default&f=martinkl&t_i=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_u=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_d=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&t_t=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&s_o=default# "Collapse")
- [+](https://disqus.com/embed/comments/?base=default&f=martinkl&t_i=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_u=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_d=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&t_t=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&s_o=default# "Expand")
- [Flag as inappropriate](https://disqus.com/embed/comments/?base=default&f=martinkl&t_i=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_u=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_d=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&t_t=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&s_o=default# "Flag as inappropriate")
![]()
That only works if the lock service and the storage system have perfectly synchronised clocks. However, clock sync typically runs over NTP, which is subject to the same network delays as all other packets, and that delay causes the clocks to be slightly out of sync with each other. Even if the clock drift is on the order of milliseconds, that's not close enough for this to be safe.
see more
G
Hello.
But wouldn't this sort of fencing token scheme imply that the storage has a critical section where it can compare the current token with the last one it received from the clients? This critical section alone seems enough to avoid conflicts. Also, with this approach, the storage could simply store object versions, and reject operations that didn't specify the most recent version.
What I'm saying is, it sure looks like if it's possible to implement fenced updates in the storage at all, then you don't need a separate lock service anymore. What am I missing?
see more
![]()
Well observed. I didn't want to go into too much detail in this already long post. There are situations in which a fencing token is still useful. For example, if there are several replicas of the storage service, then each replica can independently check the freshness of the fencing token without coordinating with the other replicas. If a quorum (e.g. 2 out of 3) of replicas accept a fencing token as valid, then a write is reported as successful to the client.
This situation is actually very similar to the second phase of Paxos or Raft (if you replace "fencing token" with "term number"), where a leader tries to replicate a value to a quorum of nodes. The first phase of Paxos/Raft is the leader election, which corresponds closely to acquiring a distributed lock. Consensus algorithms ensure that there is at most one leader for a given term number, and that term numbers monotonically increase over time. Being the leader in a particular term number corresponds to acquiring the lock with a particular fencing token.
see more
- - [−](https://disqus.com/embed/comments/?base=default&f=martinkl&t_i=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_u=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_d=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&t_t=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&s_o=default# "Collapse")
- [+](https://disqus.com/embed/comments/?base=default&f=martinkl&t_i=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_u=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_d=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&t_t=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&s_o=default# "Expand")
- [Flag as inappropriate](https://disqus.com/embed/comments/?base=default&f=martinkl&t_i=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_u=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_d=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&t_t=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&s_o=default# "Flag as inappropriate")
![]()
Thanks, Martin! Makes sense, but it also means that while trying to come up with a distributed lock managing algorithm we ended up having nothing but another consensus algorithm:)
see more
![]()
Yes, this is exactly what I expressed in the comments here about 2 years ago. Well, I am glad I am not the only one who sees it this way :)
see more
![]()
"Instead, please use a proper consensus system such as ZooKeeper, probably."
I literally skimmed the entire article for this part. :-)
That said, it's a good writeup, and very informative. Thanks for sharing!
see more
A
After reading the article one conclusion can be derived i.e. it is specific to Redlock which has been implemented in Redisson library. But I would like to implement a lock from scratch using incrementing atomic counters in redis without setting the timeout .
Implementation would look like:
getLock(key) : increment the key by 1, if the value i get back is 1 then lock acquired otherwise some one has already taken the lock.
releaseLock(ey): delete the key.
is this implementation safe enough for correctness?
see more
![]()
Not having a lock timeout means that any process that grabs the lock and becomes unavailable, makes the whole LM (lock manager) unusable. So if you do not have a timeout, then there is no sence in DLM (distributed LM) at all (though, there is no sense in DLM anyway...). Thus if you simply need an LM, then using PostgreSQL advisory session locks makes more sense than doing it with Redis, because PostgreSQL session locks are at least auto-unlocked if the holding session disconnects, thus an unavailable process holding the lock will not render a PostgreSQL-based LM unusable.
see more
![]()
Hi Martin, now that Kafka 0.11 has support for transactions could we use it in some way shape or form for distributed locks? Perhaps using an event sourced pattern where an aggregate tracks who is locking it?
I have a particularly complex use-case where multiple resources are to be locked together for the purpose of a group-related operation, so I'm looking at many different options. Postgres seems like a safe bet, but if I can take advantage of Kafka transactions that would be great.
Thanks for the article!
see more
![]()
Not knowing all the details of your requirements it's hard to give an unequivocal answer, but my guess is that Kafka transactions are irrelevant here. The main feature of Kafka transactions is to allow you to atomically publish messages to several topic-partitions. I don't see how that would help you implement a distributed lock.
I am not sure Kafka is really the right tool for this particular purpose, although if you really want to use it, then I would suggest routing all lock requests through a single topic-partition. The monotonically increasing offset that Kafka assigns to messages in a partition could then be used for fencing.
see more
- - [−](https://disqus.com/embed/comments/?base=default&f=martinkl&t_i=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_u=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_d=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&t_t=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&s_o=default# "Collapse")
- [+](https://disqus.com/embed/comments/?base=default&f=martinkl&t_i=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_u=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_d=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&t_t=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&s_o=default# "Expand")
- [Flag as inappropriate](https://disqus.com/embed/comments/?base=default&f=martinkl&t_i=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_u=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_d=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&t_t=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&s_o=default# "Flag as inappropriate")
![]()
Thanks Martin! The last part you mentioned is along the lines of what I was thinking: Aggregates would each map to a single topic-partition. A client may then attempt to lock (for correctness) one or more aggregates in a transaction. The assumption is that having "exactly-once" guarantees would help with updating these aggregates together, for example updating their "lockedBy", "lockedAt" properties.
see more
![]()
It's a long dated post, nonetheless it is particularly relevant today with the wide distribution of blockchain systems.
As it is stated, due to the fundamental asynchronicity in communication among peers, a cascade of problems arise when we have the introduction of lock leases that expire.
However, if we remove the expiring function of lock leases we run into the problem of a peer that successfully acquired a lock (by majority voting) never releasing the lease and holding onto the lock forever (whether by malicious intent, network problems, etc).
On the other hand, we may solve the problem of a lock being held indefinitely by making the majority vote for a forced release of the lock.
If the majority already voted for the forced release of a lock, and in the case that the peer that held the lock comes online and tries to write or release the lock afterwards the forced release, the operation will simply be rejected by the majority and the node that used to own the lock will revert any operations in the context of its acquired lock.
In that sense, the lock manager itself is distributed among the peers that by majority voting perform three distinct operations: grant lock, commit + lock release, revert + forced release.
see more
![]()
Hi Martin ( Martin Kleppmann ), it seems to me that in order for the storage system which is protected by a distributed lock manager (DLM) to be able to participate in fencing, it has to be able to resolve concurrent requests on its own, because it must support the following operation (yes, the whole thing is a single operation, essentially compare-and-set (CAS)):
read the last-seen-fencing-number & compare it with the incoming-fencing-number, store the incoming-fencing-number if it is greater than the last-seen-fencing-number or reject the request with the incoming-fencing-number otherwise.
If a storage system can’t do such an operation, then the system can’t benefit from fencing. But if the system can do this, then such a system is able to resolve concurrent requests on its own because it has CAS, and the DLM becomes completely redundant for correctness purposes for such a system. Am I crazy, or am I correct? :)
It would be very cool if we could clarify this.
see more
![]()
Sorry for the bump.
Main use-case I can envisage: "best-effort transaction" support, over subsets of linearizable stores (plural!) that have CAS but not transactions.
During client x's lease, x could issue (potentially multiple) writes to (potentially multiple) linearizable stores participating in the arrangement. x could even extend their lease for a long-running "pseudo-transaction", if the network plays nice.
Atomicity is the kicker, though: x might need to live with only its first i of n writes executing with serial isolation, before its lease expires (say it's unable to renew it due to packet delay) - unless the next client y is forced to begin with a compensating transaction...
see more
- - [−](https://disqus.com/embed/comments/?base=default&f=martinkl&t_i=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_u=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_d=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&t_t=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&s_o=default# "Collapse")
- [+](https://disqus.com/embed/comments/?base=default&f=martinkl&t_i=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_u=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_d=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&t_t=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&s_o=default# "Expand")
- [Flag as inappropriate](https://disqus.com/embed/comments/?base=default&f=martinkl&t_i=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_u=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_d=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&t_t=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&s_o=default# "Flag as inappropriate")
![]()
Note that another commenter brought the same exact point I did (search for "What I'm saying is, it sure looks like if it's possible to implement fenced updates in the storage at all, then you don't need a separate lock service anymore."), and even got a reply from Martin.
see more
![]()
Great post! But I think a clarification has to be made regarding fencing tokens: What happens in the (unlikely) case that client 2 also suffers a STW pause and they write with increasing successive tokens?
see more
![]()
The storage system simply maintains the ‘ratchet’ that the token can only stay the same or increase, but not decrease. Thus, if client 2 pauses and client 3 acquires the lock, client 2 will have a lesser token. If client 3 has already made a request to the storage service, client 1 and 2 will both be blocked.
see more
- - [−](https://disqus.com/embed/comments/?base=default&f=martinkl&t_i=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_u=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_d=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&t_t=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&s_o=default# "Collapse")
- [+](https://disqus.com/embed/comments/?base=default&f=martinkl&t_i=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_u=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_d=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&t_t=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&s_o=default# "Expand")
- [Flag as inappropriate](https://disqus.com/embed/comments/?base=default&f=martinkl&t_i=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_u=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_d=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&t_t=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&s_o=default# "Flag as inappropriate")
![]()
what I mean is:
client 1 acquires lock
client 1 stops
client 1 lock expires
client 2 acquires lock
client 2 stops
client 1 resumes and writes to storage
client 2 resumes and writes to storage
This is more unlikely than the scenario you presented, but still possible, and breaks the desired "correctness" since both writes are accepted.
see more
- - [−](https://disqus.com/embed/comments/?base=default&f=martinkl&t_i=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_u=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_d=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&t_t=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&s_o=default# "Collapse")
- [+](https://disqus.com/embed/comments/?base=default&f=martinkl&t_i=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_u=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_d=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&t_t=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&s_o=default# "Expand")
- [Flag as inappropriate](https://disqus.com/embed/comments/?base=default&f=martinkl&t_i=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_u=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_d=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&t_t=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&s_o=default# "Flag as inappropriate")
![]()
Oh, I see what you mean now. Yes, you have a good point — that scenario does look risky. To reason about it properly, I think we would need to make some assumptions about the semantics of the write that occurs under the lock; I think it would probably turn out to be safe in some cases, but unsafe in others. I will think about it some more.
In consensus algorithms that use epochs/ballot numbers/round numbers (which have a similar function to the fencing token), the algorithm works because the type of write is constrained. Thus, Paxos for example can maintain the invariant that if one node decides x, no other node will decide a value other than x, even in the presence of arbitrary pauses. If unconstrained writes were allowed, the safety property could be violated.
Perhaps it would be useful to regard a storage system with fencing token support as participating in an approximation of a consensus algorithm, but further protocol constraints would be required to make it safe in all circumstances?
see more
- - [−](https://disqus.com/embed/comments/?base=default&f=martinkl&t_i=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_u=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_d=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&t_t=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&s_o=default# "Collapse")
- [+](https://disqus.com/embed/comments/?base=default&f=martinkl&t_i=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_u=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_d=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&t_t=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&s_o=default# "Expand")
- [Flag as inappropriate](https://disqus.com/embed/comments/?base=default&f=martinkl&t_i=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_u=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_d=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&t_t=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&s_o=default# "Flag as inappropriate")
![]()
After some thinking, this is how I would do it:
In lock manager:
- If since last released lock, no other (later) has been expired, then next returned token is an "ordinary token" (incrementing the previous one)
- Otherwise, the next returned token is a "paired token" containing major/minor information, being major: the current token numbering, and minor: the numbering of the first token not released at this time
In lock-aware resources:
- Keep record of the highest accepted token
- If the current token is ordinary then behave as usual (rejecting when it's not greater than the highest)
- If the current token is paired (granted after some others expiration) then accept only if its minor number is greater than highest known
This would be consistent with my previous example:
-client 1 acquires lock (token: "1")
-client 1 stops
-client 1 lock expires
-client 2 acquires lock (token "2:1", meaning "lock 2 given after 1 expired")
-client 2 stops
-client 1 resumes and writes to storage
-storage accepts token and sets the highest known to "1"
-client 2 resumes and writes to storage
-storage rejects token "2:1" since "1" is not greater than highest ("1")
what do you think?
see more
- - [−](https://disqus.com/embed/comments/?base=default&f=martinkl&t_i=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_u=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_d=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&t_t=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&s_o=default# "Collapse")
- [+](https://disqus.com/embed/comments/?base=default&f=martinkl&t_i=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_u=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_d=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&t_t=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&s_o=default# "Expand")
- [Flag as inappropriate](https://disqus.com/embed/comments/?base=default&f=martinkl&t_i=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_u=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_d=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&t_t=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&s_o=default# "Flag as inappropriate")
![]()
I think your method does not work:
-client 1 acquires lock (token: "1")
-client 1 stops
-client 1 lock expires
-client 2 acquires lock (token "2:1", meaning "lock 2 given after 1 expired")
-client 2 do not write storage
-client 2 released lock
-client 3 acquires lock (token: "3")
-client 3 stops
-client 3 lock expires
-client 1 resumes and writes to storage
-storage accepts token and sets the highest known to "1"
-client 3 resumes and writes to storage
-storage accepts token
see more
- - [−](https://disqus.com/embed/comments/?base=default&f=martinkl&t_i=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_u=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_d=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&t_t=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&s_o=default# "Collapse")
- [+](https://disqus.com/embed/comments/?base=default&f=martinkl&t_i=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_u=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_d=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&t_t=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&s_o=default# "Expand")
- [Flag as inappropriate](https://disqus.com/embed/comments/?base=default&f=martinkl&t_i=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_u=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_d=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&t_t=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&s_o=default# "Flag as inappropriate")
![]()
This approach means that the storage would reject all requests after the expiration of lock (token: 1) except for the client 1 requests until client 1 explicitly goes and releases the expired lock (token: 1). And this would eliminate the lock expiration idea: despite a lock can expire, the whole system (lock management + the storage) still must wait for it to be released if the owner of the lock made at least one request to the storage.
see more
- - [−](https://disqus.com/embed/comments/?base=default&f=martinkl&t_i=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_u=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_d=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&t_t=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&s_o=default# "Collapse")
- [+](https://disqus.com/embed/comments/?base=default&f=martinkl&t_i=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_u=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_d=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&t_t=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&s_o=default# "Expand")
- [Flag as inappropriate](https://disqus.com/embed/comments/?base=default&f=martinkl&t_i=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_u=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_d=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&t_t=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&s_o=default# "Flag as inappropriate")
![]()
No, this approach doesn't mean that.
First, let's clarify this idea was made in the context of fencing tokens (locking without making any timing assumptions, accepting a token only based on its value and the value of the previous accepted one).
Once the expiration of the first lock would occur, the locking system would give the lock to client 2 (token 2), but whatever comes first at the write operation would succeed. If 2 comes first, 1 is disabled automatically by the ordinary case (monotonically increasing verification), but if 1 comes first, 2 would be disabled by the new "paired" case
see more
- - [−](https://disqus.com/embed/comments/?base=default&f=martinkl&t_i=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_u=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_d=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&t_t=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&s_o=default# "Collapse")
- [+](https://disqus.com/embed/comments/?base=default&f=martinkl&t_i=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_u=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_d=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&t_t=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&s_o=default# "Expand")
- [Flag as inappropriate](https://disqus.com/embed/comments/?base=default&f=martinkl&t_i=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_u=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_d=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&t_t=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&s_o=default# "Flag as inappropriate")
![]()
Interesting idea. Seems plausible at first glance, but it's the kind of subtle protocol that would benefit from a formal proof of correctness. In these distributed systems topics it's terribly easy to accidentally overlook some edge-case.
see more
- - [−](https://disqus.com/embed/comments/?base=default&f=martinkl&t_i=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_u=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_d=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&t_t=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&s_o=default# "Collapse")
- [+](https://disqus.com/embed/comments/?base=default&f=martinkl&t_i=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_u=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_d=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&t_t=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&s_o=default# "Expand")
- [Flag as inappropriate](https://disqus.com/embed/comments/?base=default&f=martinkl&t_i=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_u=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_d=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&t_t=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&s_o=default# "Flag as inappropriate")
J
This actually looks a lot like the algorithm proposed in the paper "On Optimistic Methods for Concurrency Control" by Kung and Robinson, 1981 ( http://www.eecs.harvard.edu...
I believe that this paper addresses the exact issues that idelvall mentioned and also includes a formal proof. Additionally, in the case proposed as long as client 1 and client 2 have no conflicts in what they are writing then both would still be permitted, however in this case if there was a conflict then client 2 would be rejected prior to starting its write. Would be interested in hearing your thoughts on this though.
see more
- - [−](https://disqus.com/embed/comments/?base=default&f=martinkl&t_i=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_u=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_d=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&t_t=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&s_o=default# "Collapse")
- [+](https://disqus.com/embed/comments/?base=default&f=martinkl&t_i=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_u=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_d=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&t_t=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&s_o=default# "Expand")
- [Flag as inappropriate](https://disqus.com/embed/comments/?base=default&f=martinkl&t_i=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_u=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_d=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&t_t=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&s_o=default# "Flag as inappropriate")
A
http://www.eecs.harvard.edu...
see more
- - [−](https://disqus.com/embed/comments/?base=default&f=martinkl&t_i=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_u=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_d=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&t_t=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&s_o=default# "Collapse")
- [+](https://disqus.com/embed/comments/?base=default&f=martinkl&t_i=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_u=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_d=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&t_t=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&s_o=default# "Expand")
- [Flag as inappropriate](https://disqus.com/embed/comments/?base=default&f=martinkl&t_i=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_u=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_d=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&t_t=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&s_o=default# "Flag as inappropriate")
A
hi, the link is broken
see more
- - [−](https://disqus.com/embed/comments/?base=default&f=martinkl&t_i=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_u=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_d=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&t_t=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&s_o=default# "Collapse")
- [+](https://disqus.com/embed/comments/?base=default&f=martinkl&t_i=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_u=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_d=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&t_t=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&s_o=default# "Expand")
- [Flag as inappropriate](https://disqus.com/embed/comments/?base=default&f=martinkl&t_i=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_u=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_d=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&t_t=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&s_o=default# "Flag as inappropriate")
![]()
Agreed, just an idea. I'll take a look into Chubby's paper and see how they handle this
see more
- - [−](https://disqus.com/embed/comments/?base=default&f=martinkl&t_i=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_u=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_d=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&t_t=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&s_o=default# "Collapse")
- [+](https://disqus.com/embed/comments/?base=default&f=martinkl&t_i=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_u=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_d=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&t_t=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&s_o=default# "Expand")
- [Flag as inappropriate](https://disqus.com/embed/comments/?base=default&f=martinkl&t_i=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_u=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_d=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&t_t=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&s_o=default# "Flag as inappropriate")
![]()
Chubby solved this problem by introuducing sequencer.
Quote from 2.4:
At any time, a lock holder may request a sequencer , an opaque byte-string that describes the state of the lock immediately after acquisition. It contains the name of the lock, the mode in which it was acquired(exclusive or shared), and the lock generation number. The client passes the sequencer to servers(such as file servers) if it expectes the operation to be protected by the lock. The recipient server is expected to test whether the sequencer is still valid and has the appropriate mode; if not, it should reject the request. The validity of a sequencer can be cheched against the server's Chubby cache or, if the server does not wish to maintain a session with Chubby, against the most recent sequencer that the server has observed.
In my opinion, it is only safe to check with Chubby cache, or it suffers the same problem.
And Chubby paper mentions another imperfect solution, lock-delay, which says, if a lock becomes free because the holder has failed or become inaccessible, the lock server will prevent other clients from claiming the lock for a lock-delay(1 minute) period.
see more
![]()
Note for the readers: that there is an error in the way the Redlock algorithm is used in the blog post: the final step after the majority is acquired, is to check if the total time elapsed is already over the lock TTL, and in such a case the client does not consider the lock as valid. This makes Redlock immune from client <-> lock-server delays in the messages, and makes every other delay *after* the lock validity is tested as any other GC pause during the processing of the locked resource. This is also equivalent to what happens, when using a remote lock server, if the "OK, you have the lock" reply from the server remains in the kernel buffers since the socket pauses before reading it. So where in this blog post its assumed that network delays or GC pauses during the lock acquisition stage are a problem, there is an error.
see more
![]()
This is correct, I had overlooked that additional clock check after messages are received. However, I believe that the additional check does not substantially alter the properties of the algorithm:
- Large network delay between the application and the shared resource (the thing that is protected by the lock) can still cause the resource to receive a message from an application process that no longer holds the lock, so fencing is still required.
- A GC pause between the final clock check and the resource access will not be caught by the clock check. As I point out in the article: "Remember that GC can pause a running thread at any point, including the point that is maximally inconvenient for you".
- All the dependencies on accuracy of clock measurement still hold.
see more
- - [−](https://disqus.com/embed/comments/?base=default&f=martinkl&t_i=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_u=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_d=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&t_t=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&s_o=default# "Collapse")
- [+](https://disqus.com/embed/comments/?base=default&f=martinkl&t_i=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_u=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_d=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&t_t=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&s_o=default# "Expand")
- [Flag as inappropriate](https://disqus.com/embed/comments/?base=default&f=martinkl&t_i=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_u=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_d=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&t_t=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&s_o=default# "Flag as inappropriate")
![]()
Hello Martin, thanks for your reply. Network delays between the app and the shared resource, and a GC pause *after* the check, but before doing the actual work, are all conceptually exactly the same as the "point 1" of your argument, that is, GC pauses (or other pauses) make the algo require an incremental token. So, regarding the safety of the algorithm itself, the only remaining thing would be the dependency on clock drifts, that can be argued depending on point of view. So I'm sorry to have to say that IMHO the current version of the article, by showing the wrong implementation of the algorithm, and not citing the equivalence of GC pauses processing the shared resource, with GC pauses immediately after the token is acquired, does not provide a fair picture.
see more
- - [−](https://disqus.com/embed/comments/?base=default&f=martinkl&t_i=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_u=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_d=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&t_t=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&s_o=default# "Collapse")
- [+](https://disqus.com/embed/comments/?base=default&f=martinkl&t_i=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_u=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_d=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&t_t=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&s_o=default# "Expand")
- [Flag as inappropriate](https://disqus.com/embed/comments/?base=default&f=martinkl&t_i=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_u=http%3A%2F%2Fmartin.kleppmann.com%2F2016%2F02%2F08%2Fhow-to-do-distributed-locking.html&t_d=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&t_t=How%20to%20do%20distributed%20locking%20%E2%80%94%20Martin%20Kleppmann%E2%80%99s%20blog&s_o=default# "Flag as inappropriate")
M
Bottom line is for an application programmer like me this implementation looks doubtful enough not to use it in production systems. If this does not work perfectly then can introduce bugs which will be impossible to fix.
see more
李
Today, I learn and think about the "redlock:. I am very agree you!
see more
M
So basically you're proposing to use optimistic locking with a token instead of pessimistic locking.
The obvious question is then: why use a lock at all? It seems you could just use a token generator service + a storage service check and have same guarantees.
see more
한
There is a typo in 'The man page for gettimeofday explicitly says that the time '
man -> main
see more
H
manual
see more
H
It is man page actually. https://man7.org/linux/man-pages/man2/gettimeofday.2.html
see more
G
This is an interesting topic, incidentally, I just released an article about achieving distributed locking using Redis regular IO (read & set).
Please check it out here for details: https://www.linkedin.com/pulse/master-less-cluster-wide-resource-locking-gerardo-recinto-g4mec/?trackingId=%2FRvVnuf4TwirkXACC4TYCw%3D%3D
Enjoy! :)
see more
![]()
Hi, Martin!
Thanks for the article. The fencing token approach is very helpful in case of correctness, but the main problem (as for me) a storage should support tokens validation. For example, if I have S3, and I want to store a file only from leader, how vanilla S3 implementation could check the correctness of a token. There's no a conditional put in API. Seems we need some kind of proxy on top of S3 which would handle put requests with tokens?
see more
live.rezync.com
live.rezync.com is blocked
This page has been blocked by an extension
- Try disabling your extensions.
ERR_BLOCKED_BY_CLIENT
Reload
This page has been blocked by an extension
pippio.com
pippio.com is blocked
This page has been blocked by an extension
- Try disabling your extensions.
ERR_BLOCKED_BY_CLIENT
Reload
This page has been blocked by an extension
tempest.services.disqus.com
tempest.services.disqus.com is blocked
This page has been blocked by an extension
- Try disabling your extensions.
ERR_BLOCKED_BY_CLIENT
Reload
This page has been blocked by an extension

