ZooKeeper 实战(五) Curator实现分布式锁

文章目录

  • ZooKeeper 实战(五) Curator实现分布式锁
    • 1.简介
      • 1.1.分布式锁概念
      • 1.2.Curator 分布式锁的实现方式
      • 1.3.分布式锁接口
    • 2.准备工作
    • 3.分布式可重入锁
      • 3.1.锁对象
      • 3.2.非重入式抢占锁
        • 测试代码
        • 输出日志
      • 3.3.重入式抢占锁
        • 测试代码
        • 输出日志
    • 4.分布式非可重入锁
      • 4.1.锁对象
      • 4.2.重入式抢占锁
        • 测试代码
        • 输出日志
    • 5.分布式可重入读写锁
      • 5.1.锁对象
      • 5.2.读锁和写锁的竞争
        • 测试代码
        • 输出日志
    • 6.共享信号量
      • 6.1.锁对象
      • 6.2.信号量抢占
        • 测试代码
        • 输出日志
    • 7.多共享锁
      • 7.1.锁对象
      • 7.2.获取共享锁
        • 测试代码
        • 输出日志

ZooKeeper 实战(五) Curator实现分布式锁

1.简介

1.1.分布式锁概念

分布式锁是一种用于实现分布式系统中的同步机制的技术。它允许在多个进程或线程之间实现互斥访问共享资源,以避免并发访问时的数据不一致问题。分布式锁的主要目的是在分布式系统中提供类似于全局锁的效果,以确保在任何时刻只有一个进程或线程可以访问特定的资源。

zookeeper基于临时有序节点实现分布式锁。每个客户端对某个临界资源加锁时,在zookeeper上的与该临界资源对应的指定节点的目录下,生成一个唯一的临时有序节点。 判断是否获取锁的方式很简单,只需要判断临时有序节点中序号最小的那个是否由自身创建。 当释放锁的时候,只需将这个临时有序节点删除即可。同时,其可以避免服务宕机导致的锁无法释放,而产生的死锁问题。

1.2.Curator 分布式锁的实现方式

curator-recipes中实现的锁有五种:

  • Shared Reentrant Lock 分布式可重入锁
  • Shared Lock 分布式非可重入锁
  • Shared Reentrant Read Write Lock 可重入读写锁
  • Shared Semaphore 共享信号量
  • Multi Shared Lock 多共享锁

1.3.分布式锁接口

org.apache.curator.framework.recipes.locks.InterProcessLock 该接口为分布式锁的行为规范,定义了获取锁和释放锁方法。

public interface InterProcessLock
{/*** 阻塞获取锁(如果未获取到锁,则一直阻塞)*/public void acquire() throws Exception;/*** 非阻塞获取锁** @param time 超时时间* @param unit 时间单位* @return 如果在超时时间内获取到锁则返回true,否则返回false* @throws Exception ZK errors, connection interruptions*/public boolean acquire(long time, TimeUnit unit) throws Exception;/*** 释放锁。每次获取锁之后必须调用该方法释放锁(acquire和release成对出现)*/public void release() throws Exception;/*** 判断该线程是否获取到锁* @return true/false*/boolean isAcquiredInThisProcess();
}

2.准备工作

首先,需要读者掌握 Ideal 同一项目启动多个Service的能力,详细教程可参考博主另一篇博客,博主创建了两个启动实例,一个端口号为8888,另一个9981,此处随意只要不与其他服务的端口号冲突即可。

在这里插入图片描述

另外,在启动项目之前,请根据先前所写的教程启动zookeeper的单机服务器 ,参考ZooKeeper 实战(一) 超详细的单机与集群部署教程,并创建一个存储数据节点:路径/ahao/data,数据内容40。(可参照如下操作,由于博主已经创建过,所以重新设置了一遍)

在这里插入图片描述

3.分布式可重入锁

可重入锁是一种自我保护的锁,允许同一进程或线程多次获得相同的锁,而不会造成死锁。

3.1.锁对象

类路径:org.apache.curator.framework.recipes.locks.InterProcessMutex

公开构造方法如下:

     /*** @param client 当前客户端实例* @param path   锁节点路径*/public InterProcessMutex(CuratorFramework client, String path)/*** @param client 当前客户端实例* @param path   锁节点路径* @param driver 锁驱动实例(工具类,只要提供两个方法:#createsTheLock 创建锁节点,#getsTheLock 获取当前锁)*/public InterProcessMutex(CuratorFramework client, String path, LockInternalsDriver driver)

3.2.非重入式抢占锁

由于非重入式抢占锁的场景对于5种分布式锁实现方式均适用,并且测试场景均一样,所以本节介绍完非重入式抢占锁的测试场景之后,对应的另4种实现方式的该测试场景将不在赘述。

测试场景:有两台服务实例 C1,C2。C1和C2个有两个线程(线程池调度)并发执行,对同一共享资源(/ahao/data中的数据)进行减1 操作。

测试代码
/*** @Name: CuratorDemoApplication* @Description:* @Author: ahao* @Date: 2024/1/10 3:29 PM*/
@Slf4j
@SpringBootApplication
public class CuratorDemoApplication implements ApplicationRunner{@Autowiredprivate CuratorFramework client;public static void main(String[] args) {SpringApplication.run(CuratorDemoApplication.class,args);}// 存储数据节点的路径final String dataPath = "/ahao/data";@Overridepublic void run(ApplicationArguments args) throws Exception {log.info("。。。。。。。。。。。。。容器初始化完毕。。。。。。。。。。。。。。");// 锁节点路径String lockPath = "/ahao/lock";TimeUnit.SECONDS.sleep(3);// 创建可重入锁InterProcessMutex mutex = new InterProcessMutex(client,lockPath);// 创建一个线程池ExecutorService executorService = Executors.newFixedThreadPool(2);// 保证并发执行,当前时间的秒针部分为0则结束循环int seconds = LocalDateTime.now().getSecond();while (seconds != 0){seconds = LocalDateTime.now().getSecond();}// 提交两个任务for (int i = 0; i < 2; i++) {executorService.submit(() -> {while (share(mutex)) {try {// 睡眠0.5秒TimeUnit.MILLISECONDS.sleep(500);} catch (InterruptedException e) {throw new RuntimeException(e);}}});}}/*** 用来模拟临界资源的方法*/public boolean share(final InterProcessMutex mutex){boolean b = true;try {// 获取锁if (mutex.acquire(3, TimeUnit.SECONDS)) {// 获取数据节点中的值byte[] bytes = client.getData().forPath(dataPath);String s = new String(bytes);Integer integer = Integer.valueOf(s);if(integer > 0){// 设置新值client.setData().forPath(dataPath,String.valueOf(integer-1).getBytes(StandardCharsets.UTF_8));log.info("当前值:{}",integer);b = true;}else {log.info("任务已完成。。。。");b = false;}}} catch (Exception e) {throw new RuntimeException(e);} finally {try {// 释放锁mutex.release();return b;} catch (Exception e) {log.error("释放锁失败");throw new RuntimeException(e);}}}
}
输出日志

从下方日志可见,没有出现重复的数值,保证了分布式系统中实现互斥访问共享资源,避免并发访问时的数据不一致问题。

实例c1:

2024-01-15 INFO 65348 --- [pool-3-thread-2] com.ahao.demo.CuratorDemoApplication     : 当前值:40
2024-01-15 INFO 65348 --- [pool-3-thread-1] com.ahao.demo.CuratorDemoApplication     : 当前值:39
2024-01-15 INFO 65348 --- [pool-3-thread-2] com.ahao.demo.CuratorDemoApplication     : 当前值:36
2024-01-15 INFO 65348 --- [pool-3-thread-1] com.ahao.demo.CuratorDemoApplication     : 当前值:35
2024-01-15 INFO 65348 --- [pool-3-thread-2] com.ahao.demo.CuratorDemoApplication     : 当前值:32
2024-01-15 INFO 65348 --- [pool-3-thread-1] com.ahao.demo.CuratorDemoApplication     : 当前值:31
2024-01-15 INFO 65348 --- [pool-3-thread-2] com.ahao.demo.CuratorDemoApplication     : 当前值:28
2024-01-15 INFO 65348 --- [pool-3-thread-1] com.ahao.demo.CuratorDemoApplication     : 当前值:27
2024-01-15 INFO 65348 --- [pool-3-thread-2] com.ahao.demo.CuratorDemoApplication     : 当前值:24
2024-01-15 INFO 65348 --- [pool-3-thread-1] com.ahao.demo.CuratorDemoApplication     : 当前值:23
2024-01-15 INFO 65348 --- [pool-3-thread-2] com.ahao.demo.CuratorDemoApplication     : 当前值:20
2024-01-15 INFO 65348 --- [pool-3-thread-1] com.ahao.demo.CuratorDemoApplication     : 当前值:19
2024-01-15 INFO 65348 --- [pool-3-thread-2] com.ahao.demo.CuratorDemoApplication     : 当前值:16
2024-01-15 INFO 65348 --- [pool-3-thread-1] com.ahao.demo.CuratorDemoApplication     : 当前值:15
2024-01-15 INFO 65348 --- [pool-3-thread-2] com.ahao.demo.CuratorDemoApplication     : 当前值:12
2024-01-15 INFO 65348 --- [pool-3-thread-1] com.ahao.demo.CuratorDemoApplication     : 当前值:11
2024-01-15 INFO 65348 --- [pool-3-thread-2] com.ahao.demo.CuratorDemoApplication     : 当前值:8
2024-01-15 INFO 65348 --- [pool-3-thread-1] com.ahao.demo.CuratorDemoApplication     : 当前值:7
2024-01-15 INFO 65348 --- [pool-3-thread-2] com.ahao.demo.CuratorDemoApplication     : 当前值:4
2024-01-15 INFO 65348 --- [pool-3-thread-1] com.ahao.demo.CuratorDemoApplication     : 当前值:3
2024-01-15 INFO 65348 --- [pool-3-thread-2] com.ahao.demo.CuratorDemoApplication     : 任务已完成。。。。
2024-01-15 INFO 65348 --- [pool-3-thread-1] com.ahao.demo.CuratorDemoApplication     : 任务已完成。。。。

实例c2:

2024-01-15 INFO 65387 --- [pool-3-thread-2] com.ahao.demo.CuratorDemoApplication     : 当前值:38
2024-01-15 INFO 65387 --- [pool-3-thread-1] com.ahao.demo.CuratorDemoApplication     : 当前值:37
2024-01-15 INFO 65387 --- [pool-3-thread-2] com.ahao.demo.CuratorDemoApplication     : 当前值:34
2024-01-15 INFO 65387 --- [pool-3-thread-1] com.ahao.demo.CuratorDemoApplication     : 当前值:33
2024-01-15 INFO 65387 --- [pool-3-thread-2] com.ahao.demo.CuratorDemoApplication     : 当前值:30
2024-01-15 INFO 65387 --- [pool-3-thread-1] com.ahao.demo.CuratorDemoApplication     : 当前值:29
2024-01-15 INFO 65387 --- [pool-3-thread-2] com.ahao.demo.CuratorDemoApplication     : 当前值:26
2024-01-15 INFO 65387 --- [pool-3-thread-1] com.ahao.demo.CuratorDemoApplication     : 当前值:25
2024-01-15 INFO 65387 --- [pool-3-thread-2] com.ahao.demo.CuratorDemoApplication     : 当前值:22
2024-01-15 INFO 65387 --- [pool-3-thread-1] com.ahao.demo.CuratorDemoApplication     : 当前值:21
2024-01-15 INFO 65387 --- [pool-3-thread-2] com.ahao.demo.CuratorDemoApplication     : 当前值:18
2024-01-15 INFO 65387 --- [pool-3-thread-1] com.ahao.demo.CuratorDemoApplication     : 当前值:17
2024-01-15 INFO 65387 --- [pool-3-thread-2] com.ahao.demo.CuratorDemoApplication     : 当前值:14
2024-01-15 INFO 65387 --- [pool-3-thread-1] com.ahao.demo.CuratorDemoApplication     : 当前值:13
2024-01-15 INFO 65387 --- [pool-3-thread-2] com.ahao.demo.CuratorDemoApplication     : 当前值:10
2024-01-15 INFO 65387 --- [pool-3-thread-1] com.ahao.demo.CuratorDemoApplication     : 当前值:9
2024-01-15 INFO 65387 --- [pool-3-thread-2] com.ahao.demo.CuratorDemoApplication     : 当前值:6
2024-01-15 INFO 65387 --- [pool-3-thread-1] com.ahao.demo.CuratorDemoApplication     : 当前值:5
2024-01-15 INFO 65387 --- [pool-3-thread-2] com.ahao.demo.CuratorDemoApplication     : 当前值:2
2024-01-15 INFO 65387 --- [pool-3-thread-1] com.ahao.demo.CuratorDemoApplication     : 当前值:1
2024-01-15 INFO 65387 --- [pool-3-thread-2] com.ahao.demo.CuratorDemoApplication     : 任务已完成。。。。
2024-01-15 INFO 65387 --- [pool-3-thread-1] com.ahao.demo.CuratorDemoApplication     : 任务已完成。。。。

3.3.重入式抢占锁

测试场景:有两台服务实例 C1,C2。C1和C2个有两个线程(线程池调度)并发执行,对同一共享资源(/ahao/data中的数据)进行连续两次的减1 操作,并且第一次操作成功后才能进行第二次操作。

测试代码

具体的减1操作抽离成方法,以便接下来的测试。

		@Overridepublic void run(ApplicationArguments args) throws Exception {log.info("。。。。。。。。。。。。。容器初始化完毕。。。。。。。。。。。。。。");String lockPath = "/ahao/lock";TimeUnit.SECONDS.sleep(3);// 创建可重入锁InterProcessMutex mutex = new InterProcessMutex(client,lockPath);// 创建一个线程池ExecutorService executorService = Executors.newFixedThreadPool(2);// 保证并发执行,当前时间的秒针部分为0则结束循环int seconds = LocalDateTime.now().getSecond();while (seconds != 0){seconds = LocalDateTime.now().getSecond();}// 提交两个任务for (int i = 0; i < 2; i++) {executorService.submit(() -> {// 循环执行while (share(mutex,1)) {try {// 睡眠0.5秒TimeUnit.MILLISECONDS.sleep(500);} catch (InterruptedException e) {throw new RuntimeException(e);}}});}}/*** 用来模拟临界资源的方法*/public boolean share(final InterProcessLock mutex, int n){boolean b = true;try {// 获取锁if (mutex.acquire(3, TimeUnit.SECONDS)) {// 减1操作b = doLock(n);// 最多执行三次if (b && n < 3){b = share(mutex,n+1);}}} catch (Exception e) {throw new RuntimeException(e);} finally {try {mutex.release();return b;} catch (Exception e) {log.error("释放锁失败");throw new RuntimeException(e);}}}// 减1操作private boolean doLock(int n) throws Exception {// 获取数据节点中的值byte[] bytes = client.getData().forPath(dataPath);String s = new String(bytes);Integer integer = Integer.valueOf(s);// 判断是否为0if(integer > 0){// 设置新值client.setData().forPath(dataPath,String.valueOf(integer-1).getBytes(StandardCharsets.UTF_8));log.info("第{}次加锁当前值:{}",n,integer);return true;}else {log.info("任务已完成。。。。");return false;}}
输出日志

由此可见,通过日志可以发现,每次获取到锁的线程都是连续执行三次并且重复获取锁并释放。

实例c1:

2024-01-16 INFO 91120 --- [pool-3-thread-2] com.ahao.demo.CuratorDemoApplication     : 第1次加锁当前值:34
2024-01-16 INFO 91120 --- [pool-3-thread-2] com.ahao.demo.CuratorDemoApplication     : 第2次加锁当前值:33
2024-01-16 INFO 91120 --- [pool-3-thread-2] com.ahao.demo.CuratorDemoApplication     : 第3次加锁当前值:32
2024-01-16 INFO 91120 --- [pool-3-thread-1] com.ahao.demo.CuratorDemoApplication     : 第1次加锁当前值:31
2024-01-16 INFO 91120 --- [pool-3-thread-1] com.ahao.demo.CuratorDemoApplication     : 第2次加锁当前值:30
2024-01-16 INFO 91120 --- [pool-3-thread-1] com.ahao.demo.CuratorDemoApplication     : 第3次加锁当前值:29
2024-01-16 INFO 91120 --- [pool-3-thread-2] com.ahao.demo.CuratorDemoApplication     : 第1次加锁当前值:22
2024-01-16 INFO 91120 --- [pool-3-thread-2] com.ahao.demo.CuratorDemoApplication     : 第2次加锁当前值:21
2024-01-16 INFO 91120 --- [pool-3-thread-2] com.ahao.demo.CuratorDemoApplication     : 第3次加锁当前值:20
2024-01-16 INFO 91120 --- [pool-3-thread-1] com.ahao.demo.CuratorDemoApplication     : 第1次加锁当前值:19
2024-01-16 INFO 91120 --- [pool-3-thread-1] com.ahao.demo.CuratorDemoApplication     : 第2次加锁当前值:18
2024-01-16 INFO 91120 --- [pool-3-thread-1] com.ahao.demo.CuratorDemoApplication     : 第3次加锁当前值:17
2024-01-16 INFO 91120 --- [pool-3-thread-2] com.ahao.demo.CuratorDemoApplication     : 第1次加锁当前值:10
2024-01-16 INFO 91120 --- [pool-3-thread-2] com.ahao.demo.CuratorDemoApplication     : 第2次加锁当前值:9
2024-01-16 INFO 91120 --- [pool-3-thread-2] com.ahao.demo.CuratorDemoApplication     : 第3次加锁当前值:8
2024-01-16 INFO 91120 --- [pool-3-thread-1] com.ahao.demo.CuratorDemoApplication     : 第1次加锁当前值:7
2024-01-16 INFO 91120 --- [pool-3-thread-1] com.ahao.demo.CuratorDemoApplication     : 第2次加锁当前值:6
2024-01-16 INFO 91120 --- [pool-3-thread-1] com.ahao.demo.CuratorDemoApplication     : 第3次加锁当前值:5
2024-01-16 INFO 91120 --- [pool-3-thread-2] com.ahao.demo.CuratorDemoApplication     : 任务已完成。。。。
2024-01-16 INFO 91120 --- [pool-3-thread-1] com.ahao.demo.CuratorDemoApplication     : 任务已完成。。。。

实例c2:

2024-01-16 INFO 91115 --- [pool-3-thread-1] com.ahao.demo.CuratorDemoApplication     : 第1次加锁当前值:40
2024-01-16 INFO 91115 --- [pool-3-thread-1] com.ahao.demo.CuratorDemoApplication     : 第2次加锁当前值:39
2024-01-16 INFO 91115 --- [pool-3-thread-1] com.ahao.demo.CuratorDemoApplication     : 第3次加锁当前值:38
2024-01-16 INFO 91115 --- [pool-3-thread-2] com.ahao.demo.CuratorDemoApplication     : 第1次加锁当前值:37
2024-01-16 INFO 91115 --- [pool-3-thread-2] com.ahao.demo.CuratorDemoApplication     : 第2次加锁当前值:36
2024-01-16 INFO 91115 --- [pool-3-thread-2] com.ahao.demo.CuratorDemoApplication     : 第3次加锁当前值:35
2024-01-16 INFO 91115 --- [pool-3-thread-1] com.ahao.demo.CuratorDemoApplication     : 第1次加锁当前值:28
2024-01-16 INFO 91115 --- [pool-3-thread-1] com.ahao.demo.CuratorDemoApplication     : 第2次加锁当前值:27
2024-01-16 INFO 91115 --- [pool-3-thread-1] com.ahao.demo.CuratorDemoApplication     : 第3次加锁当前值:26
2024-01-16 INFO 91115 --- [pool-3-thread-2] com.ahao.demo.CuratorDemoApplication     : 第1次加锁当前值:25
2024-01-16 INFO 91115 --- [pool-3-thread-2] com.ahao.demo.CuratorDemoApplication     : 第2次加锁当前值:24
2024-01-16 INFO 91115 --- [pool-3-thread-2] com.ahao.demo.CuratorDemoApplication     : 第3次加锁当前值:23
2024-01-16 INFO 91115 --- [pool-3-thread-1] com.ahao.demo.CuratorDemoApplication     : 第1次加锁当前值:16
2024-01-16 INFO 91115 --- [pool-3-thread-1] com.ahao.demo.CuratorDemoApplication     : 第2次加锁当前值:15
2024-01-16 INFO 91115 --- [pool-3-thread-1] com.ahao.demo.CuratorDemoApplication     : 第3次加锁当前值:14
2024-01-16 INFO 91115 --- [pool-3-thread-2] com.ahao.demo.CuratorDemoApplication     : 第1次加锁当前值:13
2024-01-16 INFO 91115 --- [pool-3-thread-2] com.ahao.demo.CuratorDemoApplication     : 第2次加锁当前值:12
2024-01-16 INFO 91115 --- [pool-3-thread-2] com.ahao.demo.CuratorDemoApplication     : 第3次加锁当前值:11
2024-01-16 INFO 91115 --- [pool-3-thread-1] com.ahao.demo.CuratorDemoApplication     : 第1次加锁当前值:4
2024-01-16 INFO 91115 --- [pool-3-thread-1] com.ahao.demo.CuratorDemoApplication     : 第2次加锁当前值:3
2024-01-16 INFO 91115 --- [pool-3-thread-1] com.ahao.demo.CuratorDemoApplication     : 第3次加锁当前值:2
2024-01-16 INFO 91115 --- [pool-3-thread-2] com.ahao.demo.CuratorDemoApplication     : 第1次加锁当前值:1
2024-01-16 INFO 91115 --- [pool-3-thread-2] com.ahao.demo.CuratorDemoApplication     : 任务已完成。。。。
2024-01-16 INFO 91115 --- [pool-3-thread-1] com.ahao.demo.CuratorDemoApplication     : 任务已完成。。。。

4.分布式非可重入锁

org.apache.curator.framework.recipes.locks.InterProcessMutex 类似,只是不可重入。

4.1.锁对象

类路径:InterProcessSemaphoreMutex

公开构造方法如下:

 		/*** @param client 当前客户端实例* @param path   锁节点路径*/		public InterProcessSemaphoreMutex(CuratorFramework client, String path);

4.2.重入式抢占锁

测试场景:有一台服务实例 C1。C1中的main线程进行连续两次的加锁操作。

测试代码
/*** @Name: CuratorDemoApplication* @Description:* @Author: ahao* @Date: 2024/1/10 3:29 PM*/
@Slf4j
@SpringBootApplication
public class CuratorDemoApplication implements ApplicationRunner{@Autowiredprivate CuratorFramework client;public static void main(String[] args) {SpringApplication.run(CuratorDemoApplication.class,args);}String dataPath = "/ahao/data";@Overridepublic void run(ApplicationArguments args) throws Exception {log.info("。。。。。。。。。。。。。容器初始化完毕。。。。。。。。。。。。。。");String lockPath = "/ahao/lock";TimeUnit.SECONDS.sleep(3);// 创建非可重入锁InterProcessSemaphoreMutex mutex = new InterProcessSemaphoreMutex(client, lockPath);// 调用测试方法share(mutex);}/*** 用来模拟临界资源的方法*/public void share(final InterProcessLock mutex){try {// 获取锁if (mutex.acquire(5, TimeUnit.SECONDS)) {log.info("第一次加锁成功");}// 再次获取锁if (!mutex.acquire(5, TimeUnit.SECONDS)) {log.info("第二次加锁失败了");}} catch (Exception e) {throw new RuntimeException(e);} finally {try {// 获取多少次锁就要释放多少次mutex.release();mutex.release();} catch (Exception e) {log.error("释放锁失败:{}",e.getStackTrace());throw new RuntimeException(e);}}}}
输出日志

本次测试仅需要启动一个实例c1即可,用于测试分布式非可重入锁多次加锁的场景。根据以下输出日志可知,第一次加锁成功,在第二次加锁时超时失败了,导致之后在第二次释放锁时出现异常。

2024-01-16 INFO 45025 --- [           main] com.ahao.demo.CuratorDemoApplication     : 第一次加锁成功
2024-01-16 INFO 45025 --- [           main] com.ahao.demo.CuratorDemoApplication     : 第二次加锁失败了
2024-01-16 ERROR 45025 --- [           main] com.ahao.demo.CuratorDemoApplication     : 释放锁失败:org.apache.curator.shaded.com.google.common.base.Preconditions.checkState(Preconditions.java:444)
2024-01-16 INFO 45025 --- [           main] ConditionEvaluationReportLoggingListener : Error starting ApplicationContext. To display the conditions report re-run your application with 'debug' enabled.
2024-01-16 ERROR 45025 --- [           main] o.s.boot.SpringApplication               : Application run failedjava.lang.IllegalStateException: Failed to execute ApplicationRunnerat org.springframework.boot.SpringApplication.callRunner(SpringApplication.java:762)at org.springframework.boot.SpringApplication.callRunners(SpringApplication.java:749)at org.springframework.boot.SpringApplication.run(SpringApplication.java:314)at org.springframework.boot.SpringApplication.run(SpringApplication.java:1303)at org.springframework.boot.SpringApplication.run(SpringApplication.java:1292)at com.ahao.demo.CuratorDemoApplication.main(CuratorDemoApplication.java:41)
Caused by: java.lang.RuntimeException: java.lang.IllegalStateException: Not acquiredat com.ahao.demo.CuratorDemoApplication.share(CuratorDemoApplication.java:83)at com.ahao.demo.CuratorDemoApplication.run(CuratorDemoApplication.java:56)at org.springframework.boot.SpringApplication.callRunner(SpringApplication.java:759)... 5 common frames omitted
Caused by: java.lang.IllegalStateException: Not acquiredat org.apache.curator.shaded.com.google.common.base.Preconditions.checkState(Preconditions.java:444)at org.apache.curator.framework.recipes.locks.InterProcessSemaphoreMutex.release(InterProcessSemaphoreMutex.java:68)at com.ahao.demo.CuratorDemoApplication.share(CuratorDemoApplication.java:80)... 7 common frames omitted

5.分布式可重入读写锁

读写锁顾名思义,包含两把锁:读锁和写锁。当写锁未生效(未被获取)时,读锁能够被多个线程获取使用。但是写锁只能被一个线程获取持有。 只有当写锁释放时,读锁才能被持有。可重入表示一个拥有写锁的线程可重入读锁,但是读锁却不能进入写锁,读锁可以重入读锁。 这也意味着写锁可以降级成读锁, 比如请求写锁 —>读锁 —->释放写锁。 从读锁升级成写锁是不行的。可重入读写锁是“公平的”,每个实例将按请求的顺序获取锁。

5.1.锁对象

类路径:org.apache.curator.framework.recipes.locks.InterProcessReadWriteLock

公开构造方法如下:

    /*** @param client 当前客户端实例* @param path   锁节点路径*/		public InterProcessReadWriteLock(CuratorFramework client, String basePath)/***  @param client 当前客户端实例* @param path   锁节点路径* @param lockData 存储在锁节点的数据内容*/public InterProcessReadWriteLock(CuratorFramework client, String basePath, byte[] lockData)

5.2.读锁和写锁的竞争

测试场景:有两台服务实例 C1,C2。C1和C2个有3个线程(2个读线程和1个写线程)并发执行,其中写线程进行加1操作并重复执行5次,每次都加写锁,而读线程进行查询操作并重复执行10次每次都加读锁。

测试代码
@Slf4j
@SpringBootApplication
public class CuratorDemoApplication implements ApplicationRunner {@Autowiredprivate CuratorFramework client;public static void main(String[] args) {SpringApplication.run(CuratorDemoApplication.class, args);}String dataPath = "/ahao/data";@Overridepublic void run(ApplicationArguments args) throws Exception {log.info("。。。。。。。。。。。。。容器初始化完毕。。。。。。。。。。。。。。");String lockPath = "/ahao/lock";TimeUnit.SECONDS.sleep(3);InterProcessReadWriteLock readWriteLock = new InterProcessReadWriteLock(client, lockPath);// 获取读锁InterProcessMutex readLock = readWriteLock.readLock();// 获取写锁InterProcessMutex writeLock = readWriteLock.writeLock();// 保证并发执行,当前时间的秒针部分为30的整数倍则结束循环int seconds = LocalDateTime.now().getSecond();while (seconds/30 != 0){seconds = LocalDateTime.now().getSecond();}for (int j = 0; j < 2; j++) {// 读线程new Thread(() -> {for (int i = 0; i < 10; i++) {try {// 加锁if (readLock.acquire(3, TimeUnit.SECONDS)) {doLock(false);}} catch (Exception e) {throw new RuntimeException(e);} finally {// 释放锁try {readLock.release();} catch (Exception e) {throw new RuntimeException(e);}}}}, "读线程"+j).start();}// 写线程new Thread(() -> {for (int i = 0; i < 5; i++) {try {// 加锁if (writeLock.acquire(3, TimeUnit.SECONDS)) {doLock(true);}} catch (Exception e) {throw new RuntimeException(e);} finally {// 释放锁try {writeLock.release();} catch (Exception e) {throw new RuntimeException(e);}}}}, "写线程").start();}/*** 加操作* @param isAdd true 表示加1,false 表示不加,查询数据* @return* @throws Exception*/public void doLock(boolean isAdd) throws Exception {// 获取数据节点中的值byte[] bytes = client.getData().forPath(dataPath);Integer integer = Integer.valueOf(new String(bytes));if (isAdd) {// 设置新值client.setData().forPath(dataPath, String.valueOf(integer + 1).getBytes(StandardCharsets.UTF_8));log.info("加1操作后:{}", integer);} else {log.info("查询数据:{}", integer);}}   
}
输出日志

可以观察到,读线程所查询的数据存在重复数据,说明了在同一时刻可以加多个读锁,而写线程不会出现重复数据,只能有一个线程可以获取到写锁。

实例c1:

2024-01-16  INFO 64462 --- [  写线程] com.ahao.demo.CuratorDemoApplication     : 加1操作后:40
2024-01-16  INFO 64462 --- [ 读线程0] com.ahao.demo.CuratorDemoApplication     : 查询数据:41
2024-01-16  INFO 64462 --- [ 读线程1] com.ahao.demo.CuratorDemoApplication     : 查询数据:42
2024-01-16  INFO 64462 --- [  写线程] com.ahao.demo.CuratorDemoApplication     : 加1操作后:42
2024-01-16  INFO 64462 --- [ 读线程0] com.ahao.demo.CuratorDemoApplication     : 查询数据:43
2024-01-16  INFO 64462 --- [ 读线程1] com.ahao.demo.CuratorDemoApplication     : 查询数据:44
2024-01-16  INFO 64462 --- [  写线程] com.ahao.demo.CuratorDemoApplication     : 加1操作后:44
2024-01-16  INFO 64462 --- [ 读线程0] com.ahao.demo.CuratorDemoApplication     : 查询数据:45
2024-01-16  INFO 64462 --- [ 读线程1] com.ahao.demo.CuratorDemoApplication     : 查询数据:46
2024-01-16  INFO 64462 --- [  写线程] com.ahao.demo.CuratorDemoApplication     : 加1操作后:46
2024-01-16  INFO 64462 --- [ 读线程0] com.ahao.demo.CuratorDemoApplication     : 查询数据:47
2024-01-16  INFO 64462 --- [ 读线程1] com.ahao.demo.CuratorDemoApplication     : 查询数据:48
2024-01-16  INFO 64462 --- [  写线程] com.ahao.demo.CuratorDemoApplication     : 加1操作后:48
2024-01-16  INFO 64462 --- [ 读线程0] com.ahao.demo.CuratorDemoApplication     : 查询数据:49
2024-01-16  INFO 64462 --- [ 读线程0] com.ahao.demo.CuratorDemoApplication     : 查询数据:50
2024-01-16  INFO 64462 --- [ 读线程1] com.ahao.demo.CuratorDemoApplication     : 查询数据:50
2024-01-16  INFO 64462 --- [ 读线程0] com.ahao.demo.CuratorDemoApplication     : 查询数据:50
2024-01-16  INFO 64462 --- [ 读线程1] com.ahao.demo.CuratorDemoApplication     : 查询数据:50
2024-01-16  INFO 64462 --- [ 读线程0] com.ahao.demo.CuratorDemoApplication     : 查询数据:50
2024-01-16  INFO 64462 --- [ 读线程1] com.ahao.demo.CuratorDemoApplication     : 查询数据:50
2024-01-16  INFO 64462 --- [ 读线程1] com.ahao.demo.CuratorDemoApplication     : 查询数据:50
2024-01-16  INFO 64462 --- [ 读线程0] com.ahao.demo.CuratorDemoApplication     : 查询数据:50
2024-01-16  INFO 64462 --- [ 读线程1] com.ahao.demo.CuratorDemoApplication     : 查询数据:50
2024-01-16  INFO 64462 --- [ 读线程0] com.ahao.demo.CuratorDemoApplication     : 查询数据:50
2024-01-16  INFO 64462 --- [ 读线程1] com.ahao.demo.CuratorDemoApplication     : 查询数据:50

实例c2:

2024-01-16  INFO 64453 --- [ 读线程1] com.ahao.demo.CuratorDemoApplication     : 查询数据:40
2024-01-16  INFO 64453 --- [  写线程] com.ahao.demo.CuratorDemoApplication     : 加1操作后:41
2024-01-16  INFO 64453 --- [ 读线程0] com.ahao.demo.CuratorDemoApplication     : 查询数据:42
2024-01-16  INFO 64453 --- [ 读线程1] com.ahao.demo.CuratorDemoApplication     : 查询数据:42
2024-01-16  INFO 64453 --- [  写线程] com.ahao.demo.CuratorDemoApplication     : 加1操作后:43
2024-01-16  INFO 64453 --- [ 读线程1] com.ahao.demo.CuratorDemoApplication     : 查询数据:44
2024-01-16  INFO 64453 --- [ 读线程0] com.ahao.demo.CuratorDemoApplication     : 查询数据:44
2024-01-16  INFO 64453 --- [  写线程] com.ahao.demo.CuratorDemoApplication     : 加1操作后:45
2024-01-16  INFO 64453 --- [ 读线程1] com.ahao.demo.CuratorDemoApplication     : 查询数据:46
2024-01-16  INFO 64453 --- [ 读线程0] com.ahao.demo.CuratorDemoApplication     : 查询数据:46
2024-01-16  INFO 64453 --- [  写线程] com.ahao.demo.CuratorDemoApplication     : 加1操作后:47
2024-01-16  INFO 64453 --- [ 读线程1] com.ahao.demo.CuratorDemoApplication     : 查询数据:48
2024-01-16  INFO 64453 --- [ 读线程0] com.ahao.demo.CuratorDemoApplication     : 查询数据:48
2024-01-16  INFO 64453 --- [  写线程] com.ahao.demo.CuratorDemoApplication     : 加1操作后:49
2024-01-16  INFO 64453 --- [ 读线程1] com.ahao.demo.CuratorDemoApplication     : 查询数据:50
2024-01-16  INFO 64453 --- [ 读线程0] com.ahao.demo.CuratorDemoApplication     : 查询数据:50
2024-01-16  INFO 64453 --- [ 读线程0] com.ahao.demo.CuratorDemoApplication     : 查询数据:50
2024-01-16  INFO 64453 --- [ 读线程1] com.ahao.demo.CuratorDemoApplication     : 查询数据:50
2024-01-16  INFO 64453 --- [ 读线程0] com.ahao.demo.CuratorDemoApplication     : 查询数据:50
2024-01-16  INFO 64453 --- [ 读线程1] com.ahao.demo.CuratorDemoApplication     : 查询数据:50
2024-01-16  INFO 64453 --- [ 读线程0] com.ahao.demo.CuratorDemoApplication     : 查询数据:50
2024-01-16  INFO 64453 --- [ 读线程1] com.ahao.demo.CuratorDemoApplication     : 查询数据:50
2024-01-16  INFO 64453 --- [ 读线程0] com.ahao.demo.CuratorDemoApplication     : 查询数据:50
2024-01-16  INFO 64453 --- [ 读线程1] com.ahao.demo.CuratorDemoApplication     : 查询数据:50
2024-01-16  INFO 64453 --- [ 读线程0] com.ahao.demo.CuratorDemoApplication     : 查询数据:50

6.共享信号量

对 JUC 熟悉的读者应该了解Semaphore(信号量)。Semaphore是一种用于实现同步的对象,通常用于限制对共享资源的并发访问。Semaphore可以用来控制进入临界区的线程数量,通过使用Semaphore,可以防止过多线程同时访问共享资源,从而避免出现资源竞争和死锁等问题。

而Curator中的Semaphore和JUC中的Semaphore,如出一辙,只是一个应用于分布式场景,一个应用于进程(服务器)内部。

6.1.锁对象

类路径:org.apache.curator.framework.recipes.locks.InterProcessSemaphoreV2

公开构造方法如下:

    /*** @param client 当前客户端实例* @param path   锁节点路径* @param maxLeases 允许线程进入临界区的最大数量(许可证数量)*/public InterProcessSemaphoreV2(CuratorFramework client, String path, int maxLeases);/*** @param client 当前客户端实例* @param path   锁节点路径* @param count  用于监听许可证数量*/public InterProcessSemaphoreV2(CuratorFramework client, String path, SharedCountReader count);

6.2.信号量抢占

测试场景:有两台服务实例 C1,C2。C1和C2个有3个线程并发执行。但是有2个许可证,每个线程一次只能获取一个许可证,获取成功则执行4s耗时操作(睡眠4s),然后释放锁。

测试代码
    @Overridepublic void run(ApplicationArguments args) throws Exception {log.info("。。。。。。。。。。。。。容器初始化完毕。。。。。。。。。。。。。。");String lockPath = "/ahao/lock";TimeUnit.SECONDS.sleep(3);// 创建共享信号量InterProcessSemaphoreV2 semaphoreV2 = new InterProcessSemaphoreV2(client,lockPath,2);// 保证并发执行,当前时间的秒针部分为30的整数倍则结束循环int seconds = LocalDateTime.now().getSecond();while (seconds/30 != 0){seconds = LocalDateTime.now().getSecond();}// 创建线程争夺许可证for (int i = 0; i < 3; i++) {new Thread(() -> {Lease acquire = null;try {acquire = semaphoreV2.acquire(5, TimeUnit.SECONDS);if (acquire != null){log.info("抢到许可证,参与竞争的节点:{}",semaphoreV2.getParticipantNodes());// 睡眠4秒TimeUnit.SECONDS.sleep(4);}else {log.info("抢到许可证失败");}} catch (Exception e) {throw new RuntimeException(e);}finally {if (acquire != null){semaphoreV2.returnLease(acquire);}}}, "线程"+i).start();}}
输出日志

实例c1:

2024-01-16  INFO 4319 --- [ 线程2] com.ahao.demo.CuratorDemoApplication     : 抢到许可证
2024-01-16  INFO 4319 --- [ 线程1] com.ahao.demo.CuratorDemoApplication     : 抢到许可证
2024-01-16  INFO 4319 --- [ 线程0] com.ahao.demo.CuratorDemoApplication     : 抢到许可证失败

实例c2:

2024-01-16  INFO 4307 --- [ 线程0] com.ahao.demo.CuratorDemoApplication     : 抢到许可证
2024-01-16  INFO 4307 --- [ 线程2] com.ahao.demo.CuratorDemoApplication     : 抢到许可证
2024-01-16  INFO 4307 --- [ 线程1] com.ahao.demo.CuratorDemoApplication     : 抢到许可证失败

7.多共享锁

表示将多个锁合并为一个锁。在获取多共享锁时,必须获取其内部所有的锁,才算获取成功,否则释放所有已获取的锁。同样调用释放锁方法时,会释放所有的锁。

7.1.锁对象

类路径:org.apache.curator.framework.recipes.locks.InterProcessMultiLock

公开构造方法如下:

		/*** @param client 当前客户端实例* @param path   多个锁节点路径*/public InterProcessMultiLock(CuratorFramework client, List<String> paths);/*** @param locks 多个锁对象*/public InterProcessMultiLock(List<InterProcessLock> locks);

7.2.获取共享锁

测试场景:有一台服务实例 C1,启动3个线程并发执行,抢占同一个共享锁。

测试代码
@Overridepublic void run(ApplicationArguments args) throws Exception {log.info("。。。。。。。。。。。。。容器初始化完毕。。。。。。。。。。。。。。");String lockPath = "/ahao/lock";String lockPath2 = "/ahao/lock2";TimeUnit.SECONDS.sleep(3);// 创建锁1InterProcessMutex mutex = new InterProcessMutex(client,lockPath);// 创建锁2InterProcessMutex mutex2 = new InterProcessMutex(client,lockPath2);// 创建共享锁InterProcessMultiLock multiLock = new InterProcessMultiLock(List.of(mutex,mutex2));for (int i = 0; i < 3; i++) {new Thread(()->{try {if (multiLock.acquire(5, TimeUnit.SECONDS)) {log.info("获取到锁");TimeUnit.SECONDS.sleep(3);}else {log.info("获取失败");}} catch (Exception e) {throw new RuntimeException(e);} finally {try {multiLock.release();} catch (Exception e) {throw new RuntimeException(e);}}},"线程"+i).start();}}
输出日志

可见在第3个线程获取锁时,由于没有获取到/ahao/lock2对应的锁对象导致的超时。

2024-01-16  INFO 14352 --- [ 线程1] com.ahao.demo.CuratorDemoApplication     : 获取到锁
2024-01-16  INFO 14352 --- [ 线程2] com.ahao.demo.CuratorDemoApplication     : 获取到锁
2024-01-16  INFO 14352 --- [ 线程0] com.ahao.demo.CuratorDemoApplication     : 获取失败
Exception in thread "线程0" java.lang.RuntimeException: java.lang.Exception: java.lang.IllegalMonitorStateException: You do not own the lock: /ahao/lock2at com.ahao.demo.CuratorDemoApplication.lambda$run$0(CuratorDemoApplication.java:70)at java.base/java.lang.Thread.run(Thread.java:834)
Caused by: java.lang.Exception: java.lang.IllegalMonitorStateException: You do not own the lock: /ahao/lock2at org.apache.curator.framework.recipes.locks.InterProcessMultiLock.release(InterProcessMultiLock.java:169)at com.ahao.demo.CuratorDemoApplication.lambda$run$0(CuratorDemoApplication.java:68)... 1 more
Caused by: java.lang.IllegalMonitorStateException: You do not own the lock: /ahao/lock2at org.apache.curator.framework.recipes.locks.InterProcessMutex.release(InterProcessMutex.java:140)at org.apache.curator.framework.recipes.locks.InterProcessMultiLock.release(InterProcessMultiLock.java:158)... 2 more

本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若转载,请注明出处:http://www.hqwc.cn/news/409828.html

如若内容造成侵权/违法违规/事实不符,请联系编程知识网进行投诉反馈email:809451989@qq.com,一经查实,立即删除!

相关文章

R语言【paleobioDB】——pbdb_orig_ext():绘制随着时间变化而出现的新类群

Package paleobioDB version 0.7.0 paleobioDB 包在2020年已经停止更新&#xff0c;该包依赖PBDB v1 API。 可以选择在Index of /src/contrib/Archive/paleobioDB (r-project.org)下载安装包后&#xff0c;执行本地安装。 Usage pbdb_orig_ext (data, rank, temporal_extent…

Spark---累加器和广播变量

文章目录 1.累加器实现原理2.自定义累加器3.广播变量 1.累加器实现原理 累加器用来把 Executor 端变量信息聚合到 Driver 端。在 Driver 程序中定义的变量&#xff0c;在Executor 端的每个 Task 都会得到这个变量的一份新的副本&#xff0c;每个 task 更新这些副本的值后&…

【windows】右键添加git bash here菜单

在vs 里安装了git for windows 后&#xff0c;之前git-bash 右键菜单消失了。难道是git for windows 覆盖了原来自己安装的git &#xff1f;大神给出解决方案 手动添加Git Bash Here到右键菜单&#xff08;超详细&#xff09; 安装路径&#xff1a;我老的 &#xff1f; vs的gi…

Spring5深入浅出篇:Spring工厂设计模式拓展应用

Spring5深入浅出篇:Spring工厂设计模式拓展应用 简单工厂实现 这里直接上代码举例子 UserService.java public interface UserService {public void register(User user);public void login(String name, String password); }UserServiceImpl.java public class UserService…

Netty-Netty源码分析

Netty线程模型图 Netty线程模型源码剖析图 Netty高并发高性能架构设计精髓 主从Reactor线程模型NIO多路复用非阻塞无锁串行化设计思想支持高性能序列化协议零拷贝(直接内存的使用)ByteBuf内存池设计灵活的TCP参数配置能力并发优化 无锁串行化设计思想 在大多数场景下&#…

【计算机网络】网络层——详解IP协议

个人主页&#xff1a;兜里有颗棉花糖 欢迎 点赞&#x1f44d; 收藏✨ 留言✉ 加关注&#x1f493;本文由 兜里有颗棉花糖 原创 收录于专栏【网络编程】 本专栏旨在分享学习计算机网络的一点学习心得&#xff0c;欢迎大家在评论区交流讨论&#x1f48c; 目录 &#x1f431;一、I…

【MATLAB源码-第113期】基于matlab的孔雀优化算法(POA)机器人栅格路径规划,输出做短路径图和适应度曲线。

操作环境&#xff1a; MATLAB 2022a 1、算法描述 POA&#xff08;孔雀优化算法&#xff09;是一种基于孔雀羽毛开屏行为启发的优化算法。这种算法模仿孔雀通过展开其色彩斑斓的尾羽来吸引雌性的自然行为。在算法中&#xff0c;每个孔雀代表一个潜在的解决方案&#xff0c;而…

linux驱动(六):input(key)

本文主要探讨210的input子系统。 input子系统 input子系统包含:设备驱动层,输入核心层,事件驱动层 事件处理层&#xff1a;接收核心层上报事件选择对应struct input_handler处理,每个input_handler对象处理一类事件,同类事件的设备驱动共用同一handler …

TCP连接TIME_WAIT

TCP断开过程: TIME_WAIT的作用: TIME_WAIT状态存在的理由&#xff1a; 1&#xff09;可靠地实现TCP全双工连接的终止 在进行关闭连接四次挥手协议时&#xff0c;最后的ACK是由主动关闭端发出的&#xff0c;如果这个最终的ACK丢失&#xff0c;服务器将重发最终的FIN&#xf…

Unity与Android交互通信系列(4)

上篇文章我们实现了模块化调用&#xff0c;运用了模块化设计思想和简化了调用流程&#xff0c;本篇文章讲述UnityPlayerActivity类的继承和使用。 在一些深度交互场合&#xff0c;比如Activity切换、程序启动预处理等&#xff0c;这时可能会需要继承Application和UnityPlayerAc…

【目标检测】YOLOv7算法实现(一):模型搭建

本系列文章记录本人硕士阶段YOLO系列目标检测算法自学及其代码实现的过程。其中算法具体实现借鉴于ultralytics YOLO源码Github&#xff0c;删减了源码中部分内容&#xff0c;满足个人科研需求。   本篇文章在YOLOv5算法实现的基础上&#xff0c;进一步完成YOLOv7算法的实现。…

设置了uni.chooseLocation,小程序中打不开

设置了uni.chooseLocation&#xff0c;在小程序打不开&#xff0c;点击没反应&#xff0c;地图显现不出来&#xff1b; 解决方案&#xff1a; 1.Hbuilder——微信开发者工具路径没有配置 打开工具——>设置 2.微信小程序服务端口没有开 解决方法&#xff1a;打开微信开发…