ZooKeeper 分布式锁案例

什么叫做分布式锁呢?

比如说 “进程1” 在使用该资源的时候,会先去获得锁,"进程 1"获得锁以后会对该资源保持独占,这样其他进程就无法访问该资源,"进程 1"用完该资源以后就将锁释放掉,让其他进程来获得锁,那么通过这个锁机制,我们就能保证了分布式系统中多个进程能够有序的访问该临界资源。

那么我们把这个分布式环境下的这个锁叫作分布式锁。

分布式锁案例分析:

客户端访问集群,客户端访问,常见临时的带序号的节点,序号最小的拿到资源(加锁),走业务逻辑,解锁就是删除改节点。

所以每个客户端创建的节点都会判断自己是不是序号最小的节点,如果不是就监听其前一个节点。
在这里插入图片描述

1、原生 Zookeeper 实现分布式锁案例

1.1、分布式锁实现
测试类,主要模拟十个线程进行抢占锁资源

import com.msb.zookeeper.config.ZKUtils;
import org.apache.zookeeper.ZooKeeper;
import org.junit.After;
import org.junit.Before;
import org.junit.Test;

public class TestLock {


    ZooKeeper zk ;

    @Before
    public void conn (){
        zk  = ZKUtils.getZK();//获取zookeeper的客户端连接
    }

    @After
    public void close (){
        try {
            zk.close();//关闭zookeeper连接
        } catch (InterruptedException e) {
            e.printStackTrace();
        }
    }

    @Test
    public void lock(){

    	//模拟10个线程来抢占锁资源
        for (int i = 0; i < 10; i++) {
            new Thread(){
                @Override
                public void run() {//在每一个线程内部new WatchCallBack对象,
                	//这样每个线程内部独享CountDownLatch等资源,如果十个线程共享一个CountDownLatch,一个减一其余全部都要通过了
                    WatchCallBack watchCallBack = new WatchCallBack();
                    watchCallBack.setZk(zk);
                    String threadName = Thread.currentThread().getName();
                    watchCallBack.setThreadName(threadName);
                    //每一个线程:
                    //抢锁
                    watchCallBack.tryLock();
                    //干活
                    System.out.println(threadName+" working...");
//                    try {
//                        Thread.sleep(1000);
//                    } catch (InterruptedException e) {
//                        e.printStackTrace();
//                    }
                    //释放锁
                    watchCallBack.unLock();


                }
            }.start();
        }
        while(true){

        }


    }




}

进行锁资源的抢占和释放,同时在一个目录下建立节点,有序号,序号最小表示得到了锁资源,没有抢到锁的监听前一个节点,前一个得到锁被释放之后,只有后一个节点会得到响应通知,进行回调判断

package com.msb.zookeeper.lock;

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

import java.util.Collections;
import java.util.List;
import java.util.concurrent.CountDownLatch;

public class WatchCallBack   implements Watcher, AsyncCallback.StringCallback ,AsyncCallback.Children2Callback ,AsyncCallback.StatCallback {

    ZooKeeper zk ;
    String threadName;
    CountDownLatch cc = new CountDownLatch(1);
    String pathName;

    public String getPathName() {
        return pathName;
    }

    public void setPathName(String pathName) {
        this.pathName = pathName;
    }

    public String getThreadName() {
        return threadName;
    }

    public void setThreadName(String threadName) {
        this.threadName = threadName;
    }

    public ZooKeeper getZk() {
        return zk;
    }

    public void setZk(ZooKeeper zk) {
        this.zk = zk;
    }

    public void tryLock(){
        try {

            System.out.println(threadName + "  create....");
//            if(zk.getData("/"))
            zk.create("/lock",threadName.getBytes(), ZooDefs.Ids.OPEN_ACL_UNSAFE, CreateMode.EPHEMERAL_SEQUENTIAL,this,"abc");

            cc.await();
        } catch (InterruptedException e) {
            e.printStackTrace();
        }
    }

    public void unLock(){
        try {
            zk.delete(pathName,-1);
            System.out.println(threadName + " over work....");
        } catch (InterruptedException e) {
            e.printStackTrace();
        } catch (KeeperException e) {
            e.printStackTrace();
        }
    }


    //没有抢占到锁的节点监听前一个节点,当前一个节点有变会调用这个方法,根据前一个节点的变化类型走对应方法  监控节点
    @Override
    public void process(WatchedEvent event) {


        //如果第一个哥们,那个锁释放了,其实只有第二个收到了回调事件!!
        //如果,不是第一个哥们,某一个,挂了,也能造成他后边的收到这个通知,从而让他后边那个跟去watch挂掉这个哥们前边的。。。
        switch (event.getType()) {
            case None:
                break;
            case NodeCreated:
                break;
            case NodeDeleted:
                zk.getChildren("/",false,this ,"sdf");//节点被删除了,后一个节点收到通知调用Children2Callback回调方法,没有监听
                break;
            case NodeDataChanged:
                break;
            case NodeChildrenChanged:
                break;
        }

    }

    //创建完结点以后的异步回掉方法,获取创建的节点名称
    @Override
    public void processResult(int rc, String path, Object ctx, String name) {
        if(name != null ){
            System.out.println(threadName  +"  create node : " +  name );
            pathName =  name ;
            zk.getChildren("/",false,this ,"sdf");//获取根目录下的子节点,不监视,有异步回调Children2Callback获取子节点
        }

    }

    //getChildren  call back  获取子节点的异步回调方法 获取的是当前创建节点这一时刻的前面的子节点
    //当前节点的前一个节点被删除了下一个节点收到通知,也会调用这个节点
    @Override
    public void processResult(int rc, String path, Object ctx, List<String> children, Stat stat) {

        //一定能看到自己前边的。。

//        System.out.println(threadName+"look locks.....");
//        for (String child : children) {
//            System.out.println(child);
//        }

        Collections.sort(children);//对子节点排序
        int i = children.indexOf(pathName.substring(1));//找到节点所在的下标索引顺序


        //是不是第一个
        if(i == 0){//如果是第一个就说明最靠前,所以说抢到了锁资源
            //yes
            System.out.println(threadName +" i am first....");
            try {
                zk.setData("/",threadName.getBytes(),-1);//当前节点放数据,忽略版本号
                cc.countDown();// -1 执行 cc.await(); 程序继续向下执行直到释放锁资源

            } catch (KeeperException e) {
                e.printStackTrace();
            } catch (InterruptedException e) {
                e.printStackTrace();
            }
        }else{//不是第一个 说明没有获取到锁资源
            //no
            zk.exists("/"+children.get(i-1),this,this,"sdf");//监控前一个节点 process(WatchedEvent event)
        }

    }

    //StatCallback的回调方法
    @Override
    public void processResult(int rc, String path, Object ctx, Stat stat) {
        //偷懒
    }
}

更多推荐