八、ZooKeeper 分布式锁案例
·
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) {
//偷懒
}
}
更多推荐



所有评论(0)