zookeeper分布式锁原理及使用 curator 实现分布式锁
本文为博主原创,未经允许不得转载:
1. zookeeper 分布式锁应用场景及特点分析
2. zookeeper 分布式原理
3. curator 实现分布式锁
1. zookeeper 分布式锁:
(1)优点:ZooKeeper分布式锁(如InterProcessMutex),能有效的解决分布式问题,不可重入问题,使用起来也较为简单。
(2)缺点:ZooKeeper实现的分布式锁,性能并不太高。为啥呢?
因为每次在创建锁和释放锁的过程中,都要动态创建、销毁瞬时节点来实现锁功能。大家知道,ZK中创建和删除节点只能通过Leader服务器来执行,
然后Leader服务器还需要将数据同不到所有的Follower机器上,这样频繁的网络通信,性能的短板是非常突出的。
总之,在高性能,高并发的场景下,不建议使用ZooKeeper的分布式锁。而由于ZooKeeper的高可用特性,所以在并发量不是太高的场景,推荐
使用ZooKeeper的分布式锁。
对比 redis 分布式锁:
(1)基于ZooKeeper的分布式锁,适用于高可靠(高可用)而并发量不是太大的场景;
(2)基于Redis的分布式锁,适用于并发量很大、性能要求很高的、而可靠性问题可以通过其他方案去弥补的场景。
2. zookeeper 分布式锁原理:
非公共锁方式实现流程:
如上实现方式在并发问题比较严重的情况下,性能会下降的比较厉害,主要原因是,所有的连接都在对同一个节点进行监听,当服务器检测到删除事件时, 要通知所有的连接,所有的连接同时收到事件,再次并发竞争,这就是羊群效应。这种加锁方式是非公平锁的具体实现
zookeeper 公平锁实现:
流程:
1. 请求进来,直接在 /lock 节点下穿件一个临时顺序节点
2. 判断自己是不是 lock 节点下最小的节点: 如果是最小的,则获取锁,反之对前面的节点进行监听
3. 获得锁的请求,处理完释放锁,即delete 节点,然后后继第一个节点收到通知,并重复第二步判断。
借助于临时顺序节点,可以避免同时多个节点的并发竞争锁,缓解了服务端压力。这种实现方式所有加锁请求都进行排队加锁,是公平锁的具体实现 3. 使用curator 封装分布式锁 Curator 是一套由netflix 公司开源的,Java 语言编程的 ZooKeeper 客户端框架,Curator项目是现在ZooKeeper 客户端中使用最多,Curator 把我们平时 常用的很多 ZooKeeper 服务开发功能做了封装,例如 Leader 选举、分布式计数器、分布式锁。这就减少了技术人员在使用 ZooKeeper 时的大部分底层细节 开发工作。在会话重新连接、Watch 反复注册、多种异常处理等使用场景中,用原生的 ZooKeeper处理比较复杂。而在使用 Curator 时,由于其对这些功能都 做了高度的封装,使用起来更加简单,不但减少了开发时间,而且增强了程序的可靠性。 引入curator 相关的pom 依赖: 将 Curaotr 框架引用到工程项目里,在配置文件中分别引用了两个 Curator 相关的包,第一个是 curator-framework 包,该包是对 ZooKeeper 底层 API 的一些封装。 另一个是 curator-recipes 包,该包封装了一些 ZooKeeper 服务的高级特性,如:Cache 事件监听、选举、分布式锁、分布式 Barrier。org.apache.curator curator‐recipes 5.0.0 org.apache.zookeeper zookeeper org.apache.zookeeper zookeeper 3.5.8
封装的工具类:
import org.apache.commons.lang3.exception.ExceptionUtils; import org.apache.commons.logging.Log; import org.apache.commons.logging.LogFactory; import org.apache.curator.RetryPolicy; import org.apache.curator.framework.CuratorFramework; import org.apache.curator.framework.CuratorFrameworkFactory; import org.apache.curator.framework.recipes.leader.LeaderLatch; import org.apache.curator.retry.ExponentialBackoffRetry; import org.springframework.beans.factory.annotation.Value; import org.springframework.stereotype.Component; import java.util.concurrent.TimeUnit; @Component public class BaseJob { Log log = LogFactory.getLog(this.getClass()); @Value("${zookeeper.connectstring}") private String zookeeperHost; @Value("${zookeeper.parent.path}") private String zookeeperPath = "/dlocks/zk/demo"; @Value("${curator.leadership.wait}") private int leaderShipWait = 60000; @Value("${curator.leadership.release}") private int leaderShipRelease = 360000; protected static String Zookeeper_KEY = "zk/lag/count";
public boolean getZkLock(String tokenKey) { return getZkLock(tokenKey, leaderShipRelease); } public boolean getZkLock(String tokenKey, final int lShipRelease) { CuratorFramework client = null; LeaderLatch ll = null; try { RetryPolicy policy = new ExponentialBackoffRetry(4000, 3); client = CuratorFrameworkFactory.newClient(zookeeperHost, policy); client.start(); ll = new LeaderLatch(client, String.format("%s/ppcloud/"+tokenKey+"/lock", zookeeperPath)); ll.start(); boolean get = ll.await(leaderShipWait, TimeUnit.MILLISECONDS); if (get) { final CuratorFramework client2 = client; final LeaderLatch ll2 = ll; new Thread(new Runnable() { public void run() { try { Thread.sleep(lShipRelease); } catch (Throwable e1) { } try { ll2.close(); } catch (Throwable t) { } try { client2.close(); } catch (Throwable e) { } } }).start(); return true; } else { ll.close(); client.close(); log.error(String.format("Zookeeper return:%s", get)); //hack code temp //return true; } } catch (Throwable e) { log.error(String.format("Zookeeper failed:%s", ExceptionUtils.getStackTrace(e))); try { if (ll != null) { ll.close(); } } catch (Exception e2) { } try { if (client != null) { client.close(); } } catch (Exception e2) { } } return false; } }