一、zookeeper的数据模型二、zk的监听机制三、获取分布式锁(公平锁)的总体思路四、java代码示例

一、zookeeper的数据模型

ZooKeeper的视图数据结构,很像Unix文件系统,也是树状的,这样可以确定每个路径都是唯一的。zookeeper的节点统一叫做

znode,它是可以通过路径来标识,结构图如下:

znode的四种类型

1.持久节点(PERSISTENT):这类节点被创建后,就会一直存在于Zk服务器上。直到手动删除。

2.持久顺序节点(PERSISTENT_SEQUENTIAL):它的基本特性同持久节点,不同在于增加了顺序性。父节点会维护一个自增整性数

字,用于子节点的创建的先后顺序。

3.临时节点(EPHEMERAL):临时节点的生命周期与客户端的会话绑定,一旦客户端会话失效(非TCP连接断开),那么这个节点

就会被自动清理掉。zk规定临时节点只能作为叶子节点。

4.临时顺序节点(EPHEMERAL_SEQUENTIAL):基本特性同临时节点,添加了顺序的特性。(分布式锁的实现就依靠与它

二、zk的监听机制

Zookeeper 允许客户端向服务端的某个Znode注册一个Watcher监听,当服务端的一些指定事件触发了这个Watcher,服务端会向指定客

户端发送一个事件通知来实现分布式的通知功能,然后客户端根据 Watcher通知状态和事件类型做出业务上的改变。

/*
可以把Watcher理解成客户端注册在某个Znode上的触发器,当这个Znode节点发生变化时(增删改查),就会触发Znode对应的注册事件,注册的 
客户端就会收到异步通知,然后做出业务的改变。


zk可监控的类型:
            EVENT.EventType     事件类型:(znode节点相关的)
                        NodeCreated(1):创建节点
                        NodeDeleted(2):删除节点
                        NodeDataChanged(3):节点数据变更
                        NodeChildrenChanged(4):子节点下的数据变更
            EVENT.KEEPSTATE     状态类型:(是跟客户端实例相关的)
                        @Deprecated
                        Unknown(-1):
                        Disconnected(0):连接失败
                        @Deprecated
                        NoSyncConnected(1):
                        SyncConnected(3):连接成功
                        AuthFailed(4):认证失败
                        ConnectedReadOnly(5):
                        SaslAuthenticated(6):
                        Expired(-112):会话过期
 */
//监听连接示例
zkClient = new ZooKeeper("localhost:2182", 10000, new Watcher() {
            @Override
            public void process(WatchedEvent event) {
                if (event.getState() == Event.KeeperState.Disconnected) {
                    System.out.println("失去连接");

                }
            }
        });
        
 //监听节点(LOCK_ROOT_PATH + "/" + preLockPath)示例    
  Stat stat = zkClient.exists(LOCK_ROOT_PATH + "/" + preLockPath, new Watcher() {
                @Override
                public void process(WatchedEvent event) {
                  //监听到(LOCK_ROOT_PATH + "/" + preLockPath)节点的删除节点事件
                    if(Event.EventType.NodeDeleted == event.getType()){
                        System.out.println(event.getPath() + " 前锁释放");
                        synchronized (this) {
                            notifyAll();
                        }
                    }
                }
            });

三、获取分布式锁(公平锁)的总体思路

在获取分布式锁的时候在locker(自己定义)节点下创建临时顺序节点,释放锁的时候删除该临时节点。客户端调用createNode方法

在locker下创建临时顺序节点,然后调用getChildren(“locker”)来获取locker下面的所有子节点,注意此时不用设置任何Watcher。客户端

获取到所有的子节点path之后,如果发现自己在之前创建的子节点序号最小,那么就认为该客户端获取到了锁。如果发现自己创建的节点

并非locker所有子节点中最小的,说明自己还没有获取到锁,此时客户端需要找到比自己小的那个节点,然后对其调用exist()方法,同时

对其注册事件监听器。之后,让这个被关注的节点删除,则客户端的Watcher会收到相应通知,此时再次判断自己创建的节点是否是

locker子节点中序号最小的,如果是则获取到了锁,如果不是则重复以上步骤继续获取到比自己小的一个节点并注册监听。当前这个过程

中还需要许多的逻辑判断。

非公平锁:全部客户端创建临时节点(不是顺序节点),谁先创建就谁获得锁

四、java代码示例

创建锁删除锁监听节点:

import org.apache.zookeeper.*;
import org.apache.zookeeper.data.Stat;

import java.io.IOException;
import java.util.Collections;
import java.util.List;

public class ZkClient {
    ZooKeeper zkClient;
    final String LOCK_ROOT_PATH = "/Locks";
    final String LOCK_NODE_NAME = "Lock_";
    String lockPath;

    public ZkClient() throws IOException {
        zkClient = new ZooKeeper("localhost:2182", 10000, new Watcher() {
            @Override
            public void process(WatchedEvent event) {
                if (event.getState() == Event.KeeperState.Disconnected) {
                    System.out.println("失去连接");

                }
            }
        });
    }

    public void releaseLock() throws KeeperException, InterruptedException {
        zkClient.delete(lockPath, -1);
//        zkClient.close();
        System.out.println(" 锁释放:" + lockPath);
    }

    public  void acquireLock() throws InterruptedException, KeeperException {
        //创建锁节点
        createLock();
        //尝试获取锁
        attemptLock();
    }

    private void createLock() throws KeeperException, InterruptedException {
        //如果根节点不存在,则创建根节点
        Stat stat = zkClient.exists(LOCK_ROOT_PATH, false);
        if (stat == null) {
            zkClient.create(LOCK_ROOT_PATH, new byte[0], ZooDefs.Ids.OPEN_ACL_UNSAFE, CreateMode.PERSISTENT);
        }
        // 创建EPHEMERAL_SEQUENTIAL类型节点
        String lockPath = zkClient.create(LOCK_ROOT_PATH + "/" + LOCK_NODE_NAME,
                Thread.currentThread().getName().getBytes(), ZooDefs.Ids.OPEN_ACL_UNSAFE,
                CreateMode.EPHEMERAL_SEQUENTIAL);
        System.out.println(Thread.currentThread().getName() + " 锁创建: " + lockPath);
        this.lockPath=lockPath;
    }

    private Watcher watcher = new Watcher() {
        @Override
        public void process(WatchedEvent event) {
            System.out.println(event.getPath() + " 前锁释放");
            synchronized (this) {
                notifyAll();
            }
        }
    };

    private void attemptLock() throws KeeperException, InterruptedException {
        // 获取Lock所有子节点,按照节点序号排序
        List<String> lockPaths = null;

        lockPaths = zkClient.getChildren(LOCK_ROOT_PATH, false);

        Collections.sort(lockPaths);
        int index = lockPaths.indexOf(lockPath.substring(LOCK_ROOT_PATH.length() + 1));

        // 如果lockPath是序号最小的节点,则获取锁
        if (index == 0) {
            System.out.println(Thread.currentThread().getName() + " 锁获得, lockPath: " + lockPath);
            return ;
        } else {
            // lockPath不是序号最小的节点,监听前一个节点
            String preLockPath = lockPaths.get(index - 1);

            Stat stat = zkClient.exists(LOCK_ROOT_PATH + "/" + preLockPath, new Watcher() {
                @Override
                public void process(WatchedEvent event) {
                    if(Event.EventType.NodeDeleted == event.getType()){
                        System.out.println(event.getPath() + " 前锁释放");
                        synchronized (this) {
                            notifyAll();
                        }
                    }
                }
            });

            // 假如前一个节点不存在了,比如说执行完毕,或者执行节点掉线,重新获取锁
            if (stat == null) {
                attemptLock();
            } else { // 阻塞当前进程,直到preLockPath释放锁,被watcher观察到,notifyAll后,重新acquireLock
                System.out.println(" 等待前锁释放,prelocakPath:"+preLockPath);
                synchronized (watcher) {
                    watcher.wait();
                }
                attemptLock();
            }
        }
    }
}

模拟业务:

//售票员1
public class Seller1 {
    public static void main(String[] args) throws Exception {
        ZkClient lock = new ZkClient();
        for (int i = 0; i < 100; i++) {
            lock.acquireLock();
            System.out.println("售票开始");
            // 线程随机休眠数毫秒,模拟现实中的费时操作
            int sleepMillis = (int) (Math.random() * 20000);
            try {
                //代表复杂逻辑执行了一段时间
                Thread.sleep(sleepMillis);
                System.out.println("售票结束");
                lock.releaseLock();
                } catch (InterruptedException e) {
                  e.printStackTrace();
                }
         }
    }
}

//售票员2
public class Seller2 {
    public static void main(String[] args) throws Exception {
        ZkClient lock = new ZkClient();
        for (int i = 0; i < 100; i++) {
            lock.acquireLock();
            System.out.println("售票开始");
            // 线程随机休眠数毫秒,模拟现实中的费时操作
            int sleepMillis = (int) (Math.random() * 20000);
            try {
                //代表复杂逻辑执行了一段时间
                Thread.sleep(sleepMillis);
                System.out.println("售票结束");
                lock.releaseLock();
                } catch (InterruptedException e) {
                  e.printStackTrace();
                }
         }
    }
}

Logo

汇聚全球AI编程工具,助力开发者即刻编程。

更多推荐