NOTE

Distributed Locks with ZooKeeper

ZooKeeper distributed-lock implementations using Apache Curator and native ZooKeeper, the lock principle, herd-effect avoidance, and known issues.

ZooKeeperCreated Updated 2 min readhistorical

This is a historical learning note and may contain outdated or incomplete understanding.

1. ZooKeeper Distributed-Lock Implementation

1.1. Based on Apache Curator

  • Maven
<dependency>
    <groupId>org.apache.curator</groupId>
    <artifactId>curator-recipes</artifactId>
    <version>4.0.0</version>
</dependency>
  • Example
public static void main(String[] args) throws Exception {
    // Create a ZooKeeper client
    RetryPolicy retryPolicy = new ExponentialBackoffRetry(1000, 3);
    CuratorFramework client = CuratorFrameworkFactory.newClient("<PRIVATE_IP>:2181,<PRIVATE_IP>:2181,<PRIVATE_IP>:2181", retryPolicy);
    client.start();

    // Create a distributed lock. The root node path of the lock namespace is /curator/lock
    InterProcessMutex mutex = new InterProcessMutex(client, "/curator/lock");
    mutex.acquire();
    // The lock has been acquired; execute the business process
    System.out.println("Enter mutex");
    // The business process is complete; release the lock
    mutex.release();

    // Close the client
    client.close();
}

1.1.1. Notes

Only the same InterProcessSemaphoreMutex instance in the same thread can acquire and release the lock. The details are as follows.

  1. Create two instances in the same thread and acquire the lock with both.
InterProcessMutex mutex = new InterProcessMutex(client, "/curator/lock");
mutex.acquire();

// It can be created successfully, and lock acquisition blocks
// This is because sequential child node 2 is created and waits for child node 1 to be released

InterProcessMutex mutex2 = new InterProcessMutex(client, "/curator/lock");
mutex2.acquire();
  1. Create two instances in the same thread, use one to acquire the lock and the other to release it.
InterProcessMutex mutex = new InterProcessMutex(client, "/curator/lock");
mutex.acquire();

InterProcessMutex mutex2 = new InterProcessMutex(client, "/curator/lock");
// Throws an exception: this is not the thread holding the lock. (The reason is that the node currently holding the lock is 1 rather than 2.)
mutex2.release();
  1. Start a new thread to release the lock.
InterProcessMutex mutex = new InterProcessMutex(client, "/curator/lock");
mutex.acquire();

new Thread(()->{
    mutex.release();// Throws an exception: this is not the thread holding the lock. (The reason is that this thread has no associated lock data.)
}).start();
  1. Reentrant
InterProcessMutex mutex = new InterProcessMutex(client, "/curator/lock");
mutex.acquire();// Can be acquired repeatedly
mutex.acquire();

// Release it as many times as it was acquired
mutex.release();
mutex.release();

1.2. Based on Native ZooKeeper

1.2.1. pom.xml

<?xml version="1.0" encoding="UTF-8"?>
<project xmlns="http://maven.apache.org/POM/4.0.0"
         xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
         xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
    <modelVersion>4.0.0</modelVersion>

    <groupId>com.zsk</groupId>
    <artifactId>test_zk</artifactId>
    <version>1.0-SNAPSHOT</version>
    <parent>
        <groupId>org.springframework.boot</groupId>
        <artifactId>spring-boot-starter-parent</artifactId>
        <version>1.5.20.RELEASE</version>
    </parent>

    <dependencies>
        <!--zookeeper-->
        <dependency>
            <groupId>org.apache.zookeeper</groupId>
            <artifactId>zookeeper</artifactId>
            <version>3.5.6</version>
        </dependency>
        <!--test-->
        <dependency>
            <groupId>org.springframework.boot</groupId>
            <artifactId>spring-boot-starter-test</artifactId>
            <scope>test</scope>
        </dependency>
    </dependencies>

    <build>
        <plugins>
            <plugin>
                <groupId>org.springframework.boot</groupId>
                <artifactId>spring-boot-maven-plugin</artifactId>
                <version>1.5.20.RELEASE</version>
                <configuration>
                    <executable>true</executable>
                </configuration>
            </plugin>
        </plugins>
    </build>

</project>

1.2.2. Code

public class DistributedLock
{
    private static final String LOCK_PATH = "/lock/";
    private String machineName;

    public DistributedLock(String machineName)
    {
        this.machineName = machineName;
    }

    public static void main(String[] args) throws Exception
    {
        String orderId = "1";
        IntStream.rangeClosed(1, 5)//number sequence
                .mapToObj(index -> "Machine" + index)//transform: add "Machine" in front
                .map(DistributedLock::new)//transform: DistributedLock constructor
                .map(lock -> (Runnable) () -> {//transform: create Runnable

                    ZooKeeper zookeeper = null;
                    try
                    {
                        zookeeper = lock.connect();

                        lock.lock(zookeeper, LOCK_PATH + orderId);
                        TimeUnit.SECONDS.sleep(3);//simulate a business operation
                    }
                    catch (Exception e)
                    {
                        e.printStackTrace();
                    }
                    finally
                    {
                        lock.unlock(zookeeper, LOCK_PATH + orderId);
                    }
                }).map(Thread::new)//transform: create Thread
                .forEach(Thread::start);//iterate and start

        TimeUnit.SECONDS.sleep(1000);
    }

    public ZooKeeper connect() throws Exception
    {
        // Execute asynchronously. Use CountDownLatch for synchronization and wait until the connection is created before continuing.
        CountDownLatch countDownLatch = new CountDownLatch(1);
        ZooKeeper zooKeeper = new ZooKeeper("127.0.0.1:2181", 5000, new Watcher()
        {
            @Override
            public void process(WatchedEvent watchedEvent)
            {
                countDownLatch.countDown();
            }
        });
        countDownLatch.await();
        System.out.println(machineName + " connected to ZooKeeper successfully");
        return zooKeeper;
    }

    public void lock(ZooKeeper zooKeeper, String lock)
    {
        zooKeeper.create(lock, "".getBytes(), ZooDefs.Ids.OPEN_ACL_UNSAFE, CreateMode.EPHEMERAL_SEQUENTIAL, new AsyncCallback.StringCallback()
        {
            @Override
            public void processResult(int rc, String path, Object ctx, String name)
            {
                if (rc == KeeperException.Code.OK.intValue())
                {
                    System.out.println(machineName + " acquired the lock successfully");
                    // Check whether this is the smallest node. If so, the lock is acquired successfully.
                    // Otherwise, watch the immediately preceding node, and when data changes, check again whether this is the smallest node.
                }
            }
        }, "ctx_data");
    }

    public void unlock(ZooKeeper zooKeeper, String lock)
    {
        // Delete the node
    }

}

2. ZooKeeper Distributed-Lock Principle

Create an ephemeral sequential node + watch the deletion event of the node immediately before your own.

ZooKeeper Use Case - Distributed Lock

  1. The client connects to ZooKeeper and creates an ephemeral, sequential child node under /lock. The first client’s child node is /lock/lock-0000000000, the second is /lock/lock-0000000001, and so on.
  2. The client gets the child-node list under /lock and determines whether the child node it created has the smallest sequence number in the current child-node list. If so, it is considered to have acquired the lock. Otherwise, it watches for child-node changes under /lock. After receiving a child-node change notification, it repeats this step until it acquires the lock.
  3. Execute the business code.
  4. After completing the business process, delete the corresponding child node to release the lock.

2.1. Why Are Ephemeral Nodes Needed?

To prevent a lock from being impossible to release after a crash.

2.2. Why Are Sequential Nodes Needed?

The smallest node acquires the lock.

2.3. How to Prevent the Herd Effect

When a lock is released, all clients are awakened. In fact, only the client whose sequence is immediately after the released node needs to be awakened.

  1. The client connects to ZooKeeper and creates an ephemeral, sequential child node under /lock. The first client’s child node is /lock/lock-0000000000, the second is /lock/lock-0000000001, and so on.
  2. The client gets the child-node list under /lock and determines whether the child node it created has the smallest sequence number in the current child-node list. If so, it is considered to have acquired the lock. Otherwise, it watches the deletion event of the child node immediately before its own. After receiving a child-node change notification, it repeats this step until it acquires the lock.
  3. Execute the business code.
  4. After completing the business process, delete the corresponding child node to release the lock.

3. Problems with ZooKeeper Distributed Locks

3.1. Performance Problem

  • ZooKeeper QPS is not high enough for high-concurrency scenarios.

3.2. Full GC Problem

  • ZooKeeper is implemented in Java. If a Full GC causes the heartbeat between ZooKeeper and the client to stop long enough for the persistent connection to disconnect, the lock is released.

4. References

Discussion

Sign in with GitHub to comment. Discussions are stored as GitHub Issues.View on GitHub