A distributed lock is a method for controlling synchronized access to shared resources between distributed systems.

The following introduces how Zookeeper implements distributed locks, explaining two types of distributed locks: exclusive locks and shared locks.

Exclusive Lock

Exclusive Locks, also known as write locks or exclusive locks, if transaction T1 adds an exclusive lock to data object O1, then during the entire locking period, only transaction T1 is allowed to read and update O1; no other transaction can read or write.

Define the lock:

/exclusive_lock/lock

Implementation:

Utilizing the uniqueness feature of Zookeeper's sibling nodes, when an exclusive lock needs to be acquired, all clients attempt to call the create() interface to create/exclusive_locka temporary child node under the node/exclusive_lock/lock, ultimately only one client can create successfully, and that client obtains the distributed lock. At the same time, all clients that did not acquire the lock can register/exclusive_locka watcher listening event for child node changes on the node, in order to re-contend for the lock.

Shared Lock

Shared Locks, also known as read locks. If transaction T1 adds a shared lock to data object O1, then the current transaction can only perform read operations on O1, and other transactions can only add shared locks to this data object until all shared locks on the data object are released.

Define the lock:

/shared_lock/[hostname]-请求类型W/R-序号

Implementation:

1. The client calls the create method to create a temporary sequential node similar to the lock definition method.

2. The client calls the getChildren interface to obtain a list of all created child nodes.

3. Determine whether the lock is obtained. For read requests, if all child nodes smaller than itself are read requests, or there are no child nodes with a smaller sequence number than itself, it indicates that the shared lock has been successfully acquired and the read logic starts executing. For write requests, if it is not the child node with the smallest sequence number, it enters a waiting state.

4. If the shared lock is not acquired, read requests register a watcher on the last write request node with a smaller sequence number than itself, and write requests register a watcher on the last node with a smaller sequence number than itself.

In actual development, we can use the APIs encapsulated in the Curator toolkit to help us implement distributed locks.

<dependency> <groupId>org.apache.curator</groupId> <artifactId>curator-recipes</artifactId> <version>x.x.x</version> </dependency>

Curator's several lock schemes:

  • 1、InterProcessMutex: Distributed reentrant exclusive lock
  • 2、InterProcessSemaphoreMutex: Distributed exclusive lock
  • 3、InterProcessReadWriteLock: Distributed read-write lock

The following example simulates 50 threads using the reentrant exclusive lock InterProcessMutex to contend for the lock simultaneously:

Example

public class InterprocessLock {
    public static void main(String[] args)  {
        CuratorFramework zkClient = getZkClient();
        String lockPath = "/lock";
        InterProcessMutex lock = new InterProcessMutex(zkClient, lockPath);
        //Simulate 50 threads contending for the lock
        for (int i = 0; i < 50; i++) {
            new Thread(new TestThread(i, lock)).start();
        }
    }


    static class TestThread implements Runnable {
        private Integer threadFlag;
        private InterProcessMutex lock;

        public TestThread(Integer threadFlag, InterProcessMutex lock) {
            this.threadFlag = threadFlag;
            this.lock = lock;
        }

        @Override
        public void run() {
            try {
                lock.acquire();
                System.out.println(The+threadFlag+thread acquired the lock);
                //Wait 1 second then release the lock
                Thread.sleep(1000);
            } catch (Exception e) {
                e.printStackTrace();
            }finally {
                try {
                    lock.release();
                } catch (Exception e) {
                    e.printStackTrace();
                }
            }
        }
    }

    private static CuratorFramework getZkClient() {
        String zkServerAddress = "192.168.3.39:2181";
        ExponentialBackoffRetry retryPolicy = new ExponentialBackoffRetry(1000, 3, 5000);
        CuratorFramework zkClient = CuratorFrameworkFactory.builder()
                .connectString(zkServerAddress)
                .sessionTimeoutMs(5000)
                .connectionTimeoutMs(5000)
                .retryPolicy(retryPolicy)
                .build();
        zkClient.start();
        return zkClient;
    }
}

The console outputs a record every second:

Source Code Download for This Chapter

Download