行业资讯
Zookeeper--08---zk实现分布式锁、案例
提示文章写完后目录可以自动生成如何生成可参考右边的帮助文档文章目录zk实现分布式锁1.zk中锁的种类2.zk如何上读锁3.zk如何上写锁4.⽺群效应可以调整成链式监听。解决这个问题。5.curator实现读写锁分布式锁案例案例分析依赖分布式锁原理--序号节点持久序号节点临时序号节点分布式锁实现测试对比单体模式下---ReentrantLockCurator框架实现分布式锁案例依赖获取客户端连接测试案例Redis分布式锁的实现zk实现分布式锁1.zk中锁的种类读锁⼤家都可以读要想上读锁的前提之前的锁没有写锁写锁只有得到写锁的才能写。要想上写锁的前提是之前没有任何锁。2.zk如何上读锁1.创建⼀个临时序号节点节点的数据是read表示是读锁2.获取当前zk中序号⽐⾃⼰⼩的所有节点3.判断最⼩节点是否是读锁如果不是读锁的话则上锁失败为最⼩节点设置监听。阻塞等待zk的watch机制会当最⼩节点发⽣变化时通知当前节点于是再执⾏第⼆步的流程如果是读锁的话则上锁成功3.zk如何上写锁1.创建⼀个临时序号节点节点的数据是write表示是 写锁2.获取zk中所有的⼦节点3.判断⾃⼰是否是最⼩的节点如果是则上写锁成功如果不是说明前⾯还有锁则上锁失败监听最⼩的节点如果最⼩节点有变化 则回到第⼆步。4.⽺群效应如果⽤上述的上锁⽅式只要有节点发⽣变化就会触发其他节点的监听事件这样的话对zk的压⼒⾮常⼤——⽺群效应。可以调整成链式监听。解决这个问题。5.curator实现读写锁importorg.apache.curator.framework.CuratorFramework;importorg.apache.curator.framework.recipes.locks.InterProcessLock;importorg.apache.curator.framework.recipes.locks.InterProcessReadWriteLock;importorg.junit.jupiter.api.Test;importorg.springframework.beans.factory.annotation.Autowired;importorg.springframework.boot.test.context.SpringBootTest;SpringBootTestpublicclassTestReadWriteLock{AutowiredprivateCuratorFrameworkclient;TestvoidtestGetReadLock()throwsException{// 读写锁InterProcessReadWriteLockinterProcessReadWriteLocknewInterProcessReadWriteLock(client,/lock1);// 获取读锁对象InterProcessLockinterProcessLockinterProcessReadWriteLock.readLock();System.out.println(等待获取读锁对象!);// 获取锁interProcessLock.acquire();for(inti1;i100;i){Thread.sleep(3000);System.out.println(i);}// 释放锁interProcessLock.release();System.out.println(等待释放锁!);}TestvoidtestGetWriteLock()throwsException{// 读写锁InterProcessReadWriteLockinterProcessReadWriteLocknewInterProcessReadWriteLock(client,/lock1);// 获取写锁对象InterProcessLockinterProcessLockinterProcessReadWriteLock.writeLock();System.out.println(等待获取写锁对象!);// 获取锁interProcessLock.acquire();for(inti1;i100;i){Thread.sleep(3000);System.out.println(i);}// 释放锁interProcessLock.release();System.out.println(等待释放锁!);}}分布式锁案例案例分析比如说 进程 1在使用该资源的时候会先去获得锁在使用该资源的时候保持独占这样其他进程就无法访问该资源进程1用完该资源以后就将锁释放掉让其他进程来获得锁那么通过这个锁机制我们就能保证了分布式系统中多个进程能够有序的访问该临界资源。那么我们把这个分布式环境下的这个锁叫作分布式锁。接收到请求后在/locks节点下创建一个临时顺序节点判断自己是不是当前节点下最小的节点是获取到锁不是对前一个节点进行监听获取到锁处理完业务后delete节点释放锁然后下面的节点将收到通知重复第二步判断依赖dependencygroupIdorg.apache.zookeeper/groupIdartifactIdzookeeper/artifactIdversion3.5.7/version/dependency分布式锁原理–序号节点持久序号节点临时序号节点分布式锁实现获取连接对zk加锁对zk解锁importorg.apache.zookeeper.*;importorg.apache.zookeeper.data.Stat;importjava.io.IOException;importjava.util.Collections;importjava.util.List;importjava.util.concurrent.CountDownLatch;publicclassDistributedLock{privatefinalStringconnectStringhadoop102:2181,hadoop103:2181,hadoop104:2181;privatefinalintsessionTimeout2000;privatefinalZooKeeperzk;privateCountDownLatchconnectLatchnewCountDownLatch(1);privateCountDownLatchwaitLatchnewCountDownLatch(1);privateStringwaitPath;privateStringcurrentMode;publicDistributedLock()throwsIOException,InterruptedException,KeeperException{// 获取连接zknewZooKeeper(connectString,sessionTimeout,newWatcher(){Overridepublicvoidprocess(WatchedEventwatchedEvent){// connectLatch 如果连接上zk 可以释放if(watchedEvent.getState()Event.KeeperState.SyncConnected){connectLatch.countDown();}// waitLatch 需要释放if(watchedEvent.getType()Event.EventType.NodeDeletedwatchedEvent.getPath().equals(waitPath)){waitLatch.countDown();}}});// 等待zk正常连接后往下走程序connectLatch.await();// 判断根节点/locks是否存在Statstatzk.exists(/locks,false);if(statnull){// 创建一下根节点zk.create(/locks,locks.getBytes(),ZooDefs.Ids.OPEN_ACL_UNSAFE,CreateMode.PERSISTENT);}}// 对zk加锁publicvoidzklock(){// 创建对应的临时带序号节点try{currentModezk.create(/locks/seq-,null,ZooDefs.Ids.OPEN_ACL_UNSAFE,CreateMode.EPHEMERAL_SEQUENTIAL);// wait一小会, 让结果更清晰一些Thread.sleep(10);// 判断创建的节点是否是最小的序号节点如果是获取到锁如果不是监听他序号前一个节点ListStringchildrenzk.getChildren(/locks,false);// 如果children 只有一个值那就直接获取锁 如果有多个节点需要判断谁最小if(children.size()1){return;}else{Collections.sort(children);// 获取节点名称 seq-00000000StringthisNodecurrentMode.substring(/locks/.length());// 通过seq-00000000获取该节点在children集合的位置intindexchildren.indexOf(thisNode);// 判断if(index-1){System.out.println(数据异常);}elseif(index0){// 就一个节点可以获取锁了return;}else{// 需要监听 他前一个节点变化waitPath/locks/children.get(index-1);zk.getData(waitPath,true,newStat());// 等待监听waitLatch.await();return;}}}catch(KeeperExceptione){e.printStackTrace();}catch(InterruptedExceptione){e.printStackTrace();}}// 解锁publicvoidunZkLock(){// 删除节点try{zk.delete(this.currentMode,-1);}catch(InterruptedExceptione){e.printStackTrace();}catch(KeeperExceptione){e.printStackTrace();}}}测试importorg.apache.zookeeper.KeeperException;importjava.io.IOException;publicclassDistributedLockTest{publicstaticvoidmain(String[]args)throwsInterruptedException,IOException,KeeperException{finalDistributedLocklock1newDistributedLock();finalDistributedLocklock2newDistributedLock();newThread(newRunnable(){Overridepublicvoidrun(){try{lock1.zklock();System.out.println(线程1 启动获取到锁);Thread.sleep(5*1000);lock1.unZkLock();System.out.println(线程1 释放锁);}catch(InterruptedExceptione){e.printStackTrace();}}}).start();newThread(newRunnable(){Overridepublicvoidrun(){try{lock2.zklock();System.out.println(线程2 启动获取到锁);Thread.sleep(5*1000);lock2.unZkLock();System.out.println(线程2 释放锁);}catch(InterruptedExceptione){e.printStackTrace();}}}).start();}}对比单体模式下—ReentrantLockCurator框架实现分布式锁案例依赖dependencygroupIdorg.apache.curator/groupIdartifactIdcurator-framework/artifactIdversion4.3.0/version/dependencydependencygroupIdorg.apache.curator/groupIdartifactIdcurator-recipes/artifactIdversion4.3.0/version/dependencydependencygroupIdorg.apache.curator/groupIdartifactIdcurator-client/artifactIdversion4.3.0/version/dependency获取客户端连接privatestaticCuratorFrameworkgetCuratorFramework(){ExponentialBackoffRetrypolicynewExponentialBackoffRetry(3000,3);CuratorFrameworkclientCuratorFrameworkFactory.builder().connectString(hadoop102:2181,hadoop103:2181,hadoop104:2181).connectionTimeoutMs(2000).sessionTimeoutMs(2000).retryPolicy(policy).build();// 启动客户端client.start();System.out.println(zookeeper 启动成功);returnclient;}测试案例importorg.apache.curator.framework.CuratorFramework;importorg.apache.curator.framework.CuratorFrameworkFactory;importorg.apache.curator.framework.recipes.locks.InterProcessMutex;importorg.apache.curator.retry.ExponentialBackoffRetry;publicclassCuratorLockTest{publicstaticvoidmain(String[]args){// 创建分布式锁1InterProcessMutexlock1newInterProcessMutex(getCuratorFramework(),/locks);// 创建分布式锁2InterProcessMutexlock2newInterProcessMutex(getCuratorFramework(),/locks);newThread(newRunnable(){Overridepublicvoidrun(){try{lock1.acquire();System.out.println(线程1 获取到锁);lock1.acquire();System.out.println(线程1 再次获取到锁);Thread.sleep(5*1000);lock1.release();System.out.println(线程1 释放锁);lock1.release();System.out.println(线程1 再次释放锁);}catch(Exceptione){e.printStackTrace();}}}).start();newThread(newRunnable(){Overridepublicvoidrun(){try{lock2.acquire();System.out.println(线程2 获取到锁);lock2.acquire();System.out.println(线程2 再次获取到锁);Thread.sleep(5*1000);lock2.release();System.out.println(线程2 释放锁);lock2.release();System.out.println(线程2 再次释放锁);}catch(Exceptione){e.printStackTrace();}}}).start();}privatestaticCuratorFrameworkgetCuratorFramework(){ExponentialBackoffRetrypolicynewExponentialBackoffRetry(3000,3);CuratorFrameworkclientCuratorFrameworkFactory.builder().connectString(hadoop102:2181,hadoop103:2181,hadoop104:2181).connectionTimeoutMs(2000).sessionTimeoutMs(2000).retryPolicy(policy).build();// 启动客户端client.start();System.out.println(zookeeper 启动成功);returnclient;}}Redis分布式锁的实现Redis–12–Redis分布式锁的实现
郑州网站建设
网页设计
企业官网