欢迎光临
我们一直在努力

基于Curator实现分布式锁

更加完整详细内容可查看【免费版Java学习笔记】和【免费版Java面试题】

免费版Java学习笔记(28w字)链接:https://www.yuque.com/aoyouaoyou/sgcqr8
免费版Java面试题(20w字)链接:https://www.yuque.com/aoyouaoyou/wh3hto
完整版Java学习笔记200w字,附有代码实现,图解清楚,仅需9.9
完整版Java面试题,150w字,高频面试题,内容详细,仅需9.9
完整版链接:
https://www.xiaohongshu.com/user/profile/63c2d512000000002601232c
祝您新的一年事事马到成功,身体健康,阖家幸福,大展宏图!

代码位置:

一、Curator分布式锁概述

1.1 核心实现类

Curator提供了多种分布式锁实现:

  • InterProcessMutex:可重入互斥锁(最常用)
  • InterProcessSemaphoreMutex:不可重入互斥锁
  • InterProcessReadWriteLock:读写锁
  • InterProcessMultiLock:多重锁
  • InterProcessSemaphoreV2:信号量

1.2 实现原理

ZooKeeper节点结构:
/lock-root (持久节点)
├── /lock-0000000001 (临时顺序节点)
├── /lock-0000000002 (临时顺序节点)
├── /lock-0000000003 (临时顺序节点)
└── …

工作原理:
1. 每个客户端创建临时顺序节点
2. 客户端获取/lock-root下所有子节点
3. 按顺序排序,检查自己是否是最小节点
4. 如果是则获得锁,否则监听前一个节点
5. 前一个节点删除时唤醒等待的客户端

Curator的InterProcessMutex是基于ZooKeeper实现的分布式可重入互斥锁,逻辑:

  • 加锁:客户端在指定路径下创建临时顺序节点,通过ZooKeeper的原子性操作保证节点创建的唯一性;若为最小序号节点则获取锁,否则监听前序节点的删除事件,等待锁释放。
  • 可重入:通过线程本地计数(ThreadLocal)记录当前线程的锁持有次数,同一线程多次调用acquire()时仅增加计数,release()时减少计数,计数为0时才真正释放锁。
  • 释放锁:计数为0时删除持有的临时顺序节点,ZooKeeper会在客户端会话过期/断开时自动删除节点,避免死锁。
  • 二、完整实现代码

    2.1 依赖引入

    <!– ZooKeeper核心依赖 –>
    <dependency>
    <groupId>org.apache.zookeeper</groupId>
    <artifactId>zookeeper</artifactId>
    <version>3.8.0</version>
    <exclusions>
    <exclusion>
    <groupId>org.slf4j</groupId>
    <artifactId>slf4j-log4j12</artifactId>
    </exclusion>
    </exclusions>
    </dependency>
    <!– Curator Framework –>
    <dependency>
    <groupId>org.apache.curator</groupId>
    <artifactId>curator-framework</artifactId>
    <version>5.1.0</version>
    <exclusions>
    <exclusion>
    <groupId>org.apache.zookeeper</groupId>
    <artifactId>zookeeper</artifactId>
    </exclusion>
    </exclusions>
    </dependency>
    <!– Curator Recipes(包含分布式队列实现) –>
    <dependency>
    <groupId>org.apache.curator</groupId>
    <artifactId>curator-recipes</artifactId>
    <version>5.1.0</version>
    <exclusions>
    <exclusion>
    <groupId>org.apache.zookeeper</groupId>
    <artifactId>zookeeper</artifactId>
    </exclusion>
    </exclusions>
    </dependency>
    <dependency>
    <groupId>org.apache.curator</groupId>
    <artifactId>curator-test</artifactId>
    <version>5.1.0</version>
    <scope>test</scope>
    </dependency>
    <!– SLF4J –>
    <dependency>
    <groupId>org.slf4j</groupId>
    <artifactId>slf4j-api</artifactId>
    <version>1.7.36</version>
    </dependency>
    <dependency>
    <groupId>org.slf4j</groupId>
    <artifactId>slf4j-simple</artifactId>
    <version>1.7.36</version>
    <scope>test</scope>
    </dependency>

    2.2 基础分布式锁实现

    package com.ao.c_zk_lock.lock;

    import org.apache.curator.framework.CuratorFramework;
    import org.apache.curator.framework.CuratorFrameworkFactory;
    import org.apache.curator.framework.recipes.locks.InterProcessLock;
    import org.apache.curator.framework.recipes.locks.InterProcessMutex;
    import org.apache.curator.framework.recipes.locks.InterProcessReadWriteLock;
    import org.apache.curator.framework.recipes.locks.InterProcessSemaphoreMutex;
    import org.apache.curator.retry.ExponentialBackoffRetry;
    import org.slf4j.Logger;
    import org.slf4j.LoggerFactory;

    import java.util.concurrent.TimeUnit;
    import java.util.concurrent.atomic.AtomicInteger;

    /**
    * Curator分布式锁基础实现
    * 支持多种锁类型和锁策略
    */
    public class CuratorDistributedLock {

    private static final Logger logger = LoggerFactory.getLogger(CuratorDistributedLock.class);

    // 默认配置
    private static final String DEFAULT_ZK_ADDRESS = "192.168.6.100:2181";
    private static final String DEFAULT_LOCK_ROOT = "/aoyou/locks";
    private static final int DEFAULT_RETRY_COUNT = 3;
    private static final int DEFAULT_BASE_SLEEP_TIME = 1000;

    // 锁类型枚举
    public enum LockType {
    MUTEX, // 可重入互斥锁
    SEMAPHORE_MUTEX, // 不可重入互斥锁
    READ_WRITE, // 读写锁
    READ, // 读锁
    WRITE // 写锁
    }

    /**
    * 锁配置类
    */
    public static class LockConfig {
    private String zkAddress = DEFAULT_ZK_ADDRESS;
    private String lockPath;
    private LockType lockType = LockType.MUTEX;
    private int retryCount = DEFAULT_RETRY_COUNT;
    private int baseSleepTime = DEFAULT_BASE_SLEEP_TIME;
    private int sessionTimeoutMs = 15000;
    private int connectionTimeoutMs = 10000;
    private String namespace = "aoyou";

    public LockConfig(String lockPath) {
    this.lockPath = DEFAULT_LOCK_ROOT + "/" + lockPath;
    }

    // 链式配置方法
    public LockConfig withZkAddress(String zkAddress) {
    this.zkAddress = zkAddress;
    return this;
    }

    public LockConfig withLockType(LockType lockType) {
    this.lockType = lockType;
    return this;
    }

    public LockConfig withRetryCount(int retryCount) {
    this.retryCount = retryCount;
    return this;
    }

    public LockConfig withBaseSleepTime(int baseSleepTime) {
    this.baseSleepTime = baseSleepTime;
    return this;
    }

    public LockConfig withSessionTimeout(int timeoutMs) {
    this.sessionTimeoutMs = timeoutMs;
    return this;
    }

    public LockConfig withConnectionTimeout(int timeoutMs) {
    this.connectionTimeoutMs = timeoutMs;
    return this;
    }

    public LockConfig withNamespace(String namespace) {
    this.namespace = namespace;
    return this;
    }
    }

    // 锁实例
    private final String lockPath;
    private final CuratorFramework client;
    private final InterProcessLock lock;
    private final LockType lockType;

    // 锁状态
    private final AtomicInteger lockCount = new AtomicInteger(0);
    private volatile boolean isLocked = false;
    private volatile Thread lockOwner;

    /**
    * 构造函数 – 默认配置
    */
    public CuratorDistributedLock(String lockName) {
    this(new LockConfig(lockName));
    }

    /**
    * 构造函数 – 自定义配置
    */
    public CuratorDistributedLock(LockConfig config) {
    this.lockPath = config.lockPath;
    this.lockType = config.lockType;

    // 创建Curator客户端
    ExponentialBackoffRetry retryPolicy = new ExponentialBackoffRetry(
    config.baseSleepTime, config.retryCount);

    this.client = CuratorFrameworkFactory.builder()
    .connectString(config.zkAddress)
    .namespace(config.namespace)
    .retryPolicy(retryPolicy)
    .sessionTimeoutMs(config.sessionTimeoutMs)
    .connectionTimeoutMs(config.connectionTimeoutMs)
    .build();

    // 创建锁实例
    this.lock = createLockInstance(config);

    // 启动客户端
    client.start();

    logger.info("初始化分布式锁: {},类型: {}", lockPath, lockType);
    }

    /**
    * 创建锁实例
    */
    private InterProcessLock createLockInstance(LockConfig config) {
    switch (config.lockType) {
    case MUTEX:
    return new InterProcessMutex(client, config.lockPath);

    case SEMAPHORE_MUTEX:
    return new InterProcessSemaphoreMutex(client, config.lockPath);

    case READ_WRITE:
    // 读写锁需要特殊处理
    InterProcessReadWriteLock readWriteLock =
    new InterProcessReadWriteLock(client, config.lockPath);
    return readWriteLock.writeLock(); // 默认返回写锁

    case READ:
    InterProcessReadWriteLock readLock =
    new InterProcessReadWriteLock(client, config.lockPath);
    return readLock.readLock();

    case WRITE:
    InterProcessReadWriteLock writeLock =
    new InterProcessReadWriteLock(client, config.lockPath);
    return writeLock.writeLock();

    default:
    throw new IllegalArgumentException("不支持的锁类型: " + config.lockType);
    }
    }

    /**
    * 获取锁(阻塞)
    * @return 是否成功获取锁
    */
    public boolean lock() throws Exception {
    return lock(0, TimeUnit.MILLISECONDS);
    }

    /**
    * 获取锁(带超时)
    * @param timeout 超时时间
    * @param unit 时间单位
    * @return 是否成功获取锁
    */
    public boolean lock(long timeout, TimeUnit unit) throws Exception {
    long startTime = System.currentTimeMillis();
    long timeoutMs = unit.toMillis(timeout);

    while (true) {
    try {
    boolean acquired;

    if (timeout > 0) {
    acquired = lock.acquire(timeout, TimeUnit.MILLISECONDS);
    } else {
    acquired = lock.acquire(timeout, TimeUnit.MILLISECONDS);
    // 对于阻塞锁,这里需要特殊处理
    if (!acquired) {
    // 如果支持带超时的acquire方法
    acquired = lock.acquire(timeoutMs, TimeUnit.MILLISECONDS);
    }
    }

    if (acquired) {
    isLocked = true;
    lockOwner = Thread.currentThread();
    lockCount.incrementAndGet();

    logger.debug("获取锁成功: {},线程: {}", lockPath, Thread.currentThread().getName());
    return true;
    }

    // 检查超时
    if (timeout > 0 && System.currentTimeMillis() – startTime > timeoutMs) {
    logger.warn("获取锁超时: {},超时时间: {}ms", lockPath, timeoutMs);
    return false;
    }

    // 等待后重试
    Thread.sleep(100);

    } catch (Exception e) {
    logger.error("获取锁异常: {}", lockPath, e);

    // 检查超时
    if (timeout > 0 && System.currentTimeMillis() – startTime > timeoutMs) {
    return false;
    }

    // 等待后重试
    Thread.sleep(100);
    }
    }
    }

    /**
    * 尝试获取锁(非阻塞)
    * @return 是否成功获取锁
    */
    public boolean tryLock() throws Exception {
    return tryLock(100, TimeUnit.MILLISECONDS);
    }

    /**
    * 尝试获取锁(带超时)
    * @param timeout 超时时间
    * @param unit 时间单位
    * @return 是否成功获取锁
    */
    public boolean tryLock(long timeout, TimeUnit unit) throws Exception {
    return lock(timeout, unit);
    }

    /**
    * 释放锁
    */
    public void unlock() throws Exception {
    if (!isLocked || lockOwner != Thread.currentThread()) {
    throw new IllegalMonitorStateException("当前线程不持有锁: " + lockPath);
    }

    int currentCount = lockCount.decrementAndGet();

    if (currentCount == 0) {
    lock.release();
    isLocked = false;
    lockOwner = null;

    logger.debug("释放锁成功: {},线程: {}", lockPath, Thread.currentThread().getName());
    } else {
    logger.debug("锁重入计数减少: {},当前计数: {}", lockPath, currentCount);
    }
    }

    /**
    * 强制释放锁(谨慎使用)
    */
    public void forceUnlock() throws Exception {
    try {
    if (isLocked) {
    lock.release();
    isLocked = false;
    lockOwner = null;
    lockCount.set(0);

    logger.warn("强制释放锁: {},原持有线程: {}", lockPath,
    lockOwner != null ? lockOwner.getName() : "未知");
    }
    } catch (Exception e) {
    logger.error("强制释放锁失败: {}", lockPath, e);
    }
    }

    /**
    * 检查是否持有锁
    */
    public boolean isHeldByCurrentThread() {
    return isLocked && lockOwner == Thread.currentThread();
    }

    /**
    * 检查锁是否被任何线程持有
    */
    public boolean isLocked() {
    return isLocked;
    }

    /**
    * 获取锁的重入计数
    */
    public int getHoldCount() {
    return lockCount.get();
    }

    /**
    * 获取锁的持有线程
    */
    public Thread getLockOwner() {
    return lockOwner;
    }

    /**
    * 获取锁路径
    */
    public String getLockPath() {
    return lockPath;
    }

    /**
    * 获取锁类型
    */
    public LockType getLockType() {
    return lockType;
    }

    /**
    * 获取锁状态信息
    */
    public LockStatus getStatus() {
    LockStatus status = new LockStatus();
    status.lockPath = lockPath;
    status.lockType = lockType;
    status.isLocked = isLocked;
    status.holdCount = lockCount.get();
    status.lockOwner = lockOwner != null ? lockOwner.getName() : "无";
    status.isHeldByCurrentThread = isHeldByCurrentThread();

    return status;
    }

    /**
    * 关闭锁
    */
    public void close() {
    try {
    if (isLocked) {
    forceUnlock();
    }

    if (client != null) {
    client.close();
    }

    logger.info("关闭分布式锁: {}", lockPath);

    } catch (Exception e) {
    logger.error("关闭锁失败: {}", lockPath, e);
    }
    }

    /**
    * 锁状态类
    */
    public static class LockStatus {
    public String lockPath;
    public LockType lockType;
    public boolean isLocked;
    public int holdCount;
    public String lockOwner;
    public boolean isHeldByCurrentThread;

    @Override
    public String toString() {
    return String.format(
    "LockStatus{路径='%s', 类型=%s, 已锁定=%s, 重入计数=%d, 持有者='%s', 当前线程持有=%s}",
    lockPath, lockType, isLocked ? "是" : "否", holdCount,
    lockOwner, isHeldByCurrentThread ? "是" : "否"
    );
    }
    }

    /**
    * 锁执行器 – 简化锁的使用
    */
    public <T> T executeWithLock(LockOperation<T> operation) throws Exception {
    return executeWithLock(operation, 0, TimeUnit.MILLISECONDS);
    }

    /**
    * 带超时的锁执行器
    */
    public <T> T executeWithLock(LockOperation<T> operation, long timeout, TimeUnit unit) throws Exception {
    boolean locked = false;

    try {
    locked = lock(timeout, unit);
    if (!locked) {
    throw new LockTimeoutException("获取锁超时: " + lockPath);
    }

    return operation.execute();

    } finally {
    if (locked && isHeldByCurrentThread()) {
    unlock();
    }
    }
    }

    /**
    * 锁操作接口
    */
    public interface LockOperation<T> {
    T execute() throws Exception;
    }

    /**
    * 锁超时异常
    */
    public static class LockTimeoutException extends RuntimeException {
    public LockTimeoutException(String message) {
    super(message);
    }
    }
    }

    2.3 高级锁实现(读写锁、重入锁)

    package com.ao.c_zk_lock.lock;

    import org.apache.curator.framework.CuratorFramework;
    import org.apache.curator.framework.recipes.locks.*;
    import org.slf4j.Logger;
    import org.slf4j.LoggerFactory;

    import java.util.concurrent.TimeUnit;
    import java.util.concurrent.locks.Condition;
    import java.util.concurrent.locks.Lock;

    /**
    * 高级分布式锁实现
    * 支持读写锁、可重入锁等高级特性
    */
    public class AdvancedDistributedLock implements Lock {

    private static final Logger logger = LoggerFactory.getLogger(AdvancedDistributedLock.class);

    // 锁实例
    private final InterProcessLock internalLock;
    private final String lockPath;
    private final LockType lockType;

    // 锁状态
    private volatile Thread lockOwner;
    private volatile int holdCount = 0;

    /**
    * 锁类型
    */
    public enum LockType {
    MUTEX, // 可重入互斥锁
    READ, // 读锁
    WRITE, // 写锁
    SEMAPHORE_MUTEX // 信号量互斥锁
    }

    /**
    * 锁统计信息
    */
    public static class LockStats {
    private final String lockPath;
    private final LockType lockType;
    private final long lockCount;
    private final long waitCount;
    private final long timeoutCount;
    private final long averageWaitTime;

    public LockStats(String lockPath, LockType lockType, long lockCount,
    long waitCount, long timeoutCount, long averageWaitTime) {
    this.lockPath = lockPath;
    this.lockType = lockType;
    this.lockCount = lockCount;
    this.waitCount = waitCount;
    this.timeoutCount = timeoutCount;
    this.averageWaitTime = averageWaitTime;
    }

    // Getters
    public String getLockPath() { return lockPath; }
    public LockType getLockType() { return lockType; }
    public long getLockCount() { return lockCount; }
    public long getWaitCount() { return waitCount; }
    public long getTimeoutCount() { return timeoutCount; }
    public long getAverageWaitTime() { return averageWaitTime; }

    /**
    * 获取锁成功率
    */
    public double getSuccessRate() {
    if (lockCount == 0) return 0.0;
    return (double) (lockCount – timeoutCount) / lockCount;
    }

    @Override
    public String toString() {
    return String.format(
    "LockStats{路径='%s', 类型=%s, 加锁次数=%d, 等待次数=%d, 超时次数=%d, 平均等待时间=%dms, 成功率=%.2f%%}",
    lockPath, lockType, lockCount, waitCount, timeoutCount,
    averageWaitTime, getSuccessRate() * 100
    );
    }
    }

    /**
    * 构造函数
    */
    public AdvancedDistributedLock(CuratorFramework client, String lockPath, LockType lockType) {
    this.lockPath = lockPath;
    this.lockType = lockType;

    switch (lockType) {
    case MUTEX:
    this.internalLock = new InterProcessMutex(client, lockPath);
    break;

    case READ:
    InterProcessReadWriteLock readWriteLock = new InterProcessReadWriteLock(client, lockPath);
    this.internalLock = readWriteLock.readLock();
    break;

    case WRITE:
    InterProcessReadWriteLock writeLock = new InterProcessReadWriteLock(client, lockPath);
    this.internalLock = writeLock.writeLock();
    break;

    case SEMAPHORE_MUTEX:
    this.internalLock = new InterProcessSemaphoreMutex(client, lockPath);
    break;

    default:
    throw new IllegalArgumentException("不支持的锁类型: " + lockType);
    }

    logger.info("创建高级分布式锁: {},类型: {}", lockPath, lockType);
    }

    @Override
    public void lock() {
    try {
    internalLock.acquire();
    lockOwner = Thread.currentThread();
    holdCount++;

    logger.debug("获取锁成功: {},线程: {},重入计数: {}",
    lockPath, Thread.currentThread().getName(), holdCount);

    } catch (Exception e) {
    throw new LockException("获取锁失败: " + lockPath, e);
    }
    }

    @Override
    public void lockInterruptibly() throws InterruptedException {
    try {
    internalLock.acquire();
    lockOwner = Thread.currentThread();
    holdCount++;

    logger.debug("可中断获取锁成功: {},线程: {}",
    lockPath, Thread.currentThread().getName());

    } catch (Exception e) {
    Thread.currentThread().interrupt();
    throw new LockException("可中断获取锁失败: " + lockPath, e);
    }
    }

    @Override
    public boolean tryLock() {
    return tryLock(0, TimeUnit.MILLISECONDS);
    }

    @Override
    public boolean tryLock(long time, TimeUnit unit) {
    long startTime = System.currentTimeMillis();

    try {
    boolean acquired = internalLock.acquire(time, unit);

    if (acquired) {
    lockOwner = Thread.currentThread();
    holdCount++;

    long waitTime = System.currentTimeMillis() – startTime;
    logger.debug("尝试获取锁成功: {},线程: {},等待时间: {}ms",
    lockPath, Thread.currentThread().getName(), waitTime);
    }

    return acquired;

    } catch (Exception e) {
    logger.error("尝试获取锁失败: {}", lockPath, e);
    return false;
    }
    }

    @Override
    public void unlock() {
    if (lockOwner != Thread.currentThread()) {
    throw new IllegalMonitorStateException(
    "当前线程不是锁的持有者: " + Thread.currentThread().getName());
    }

    try {
    holdCount–;

    if (holdCount == 0) {
    internalLock.release();
    lockOwner = null;

    logger.debug("释放锁成功: {},线程: {}",
    lockPath, Thread.currentThread().getName());
    } else {
    logger.debug("锁重入计数减少: {},当前计数: {},线程: {}",
    lockPath, holdCount, Thread.currentThread().getName());
    }

    } catch (Exception e) {
    throw new LockException("释放锁失败: " + lockPath, e);
    }
    }

    @Override
    public Condition newCondition() {
    throw new UnsupportedOperationException("分布式锁不支持Condition");
    }

    /**
    * 强制释放锁(谨慎使用)
    */
    public void forceUnlock() {
    try {
    if (lockOwner != null) {
    internalLock.release();
    lockOwner = null;
    holdCount = 0;

    logger.warn("强制释放锁: {},原持有线程: {}",
    lockPath, lockOwner != null ? lockOwner.getName() : "未知");
    }
    } catch (Exception e) {
    logger.error("强制释放锁失败: {}", lockPath, e);
    }
    }

    /**
    * 检查是否持有锁
    */
    public boolean isHeldByCurrentThread() {
    return lockOwner == Thread.currentThread();
    }

    /**
    * 检查锁是否被持有
    */
    public boolean isLocked() {
    return lockOwner != null;
    }

    /**
    * 获取重入计数
    */
    public int getHoldCount() {
    return holdCount;
    }

    /**
    * 获取锁持有者
    */
    public Thread getLockOwner() {
    return lockOwner;
    }

    /**
    * 获取锁路径
    */
    public String getLockPath() {
    return lockPath;
    }

    /**
    * 获取锁类型
    */
    public LockType getLockType() {
    return lockType;
    }

    /**
    * 执行带锁的操作
    */
    public <T> T executeWithLock(LockableOperation<T> operation) throws Exception {
    return executeWithLock(operation, 0, TimeUnit.MILLISECONDS);
    }

    /**
    * 执行带锁的操作(带超时)
    */
    public <T> T executeWithLock(LockableOperation<T> operation, long timeout, TimeUnit unit) throws Exception {
    boolean locked = false;

    try {
    locked = tryLock(timeout, unit);
    if (!locked) {
    throw new LockTimeoutException("获取锁超时: " + lockPath);
    }

    return operation.execute();

    } finally {
    if (locked && isHeldByCurrentThread()) {
    unlock();
    }
    }
    }

    /**
    * 锁操作接口
    */
    public interface LockableOperation<T> {
    T execute() throws Exception;
    }

    /**
    * 锁异常
    */
    public static class LockException extends RuntimeException {
    public LockException(String message) {
    super(message);
    }

    public LockException(String message, Throwable cause) {
    super(message, cause);
    }
    }

    /**
    * 锁超时异常
    */
    public static class LockTimeoutException extends RuntimeException {
    public LockTimeoutException(String message) {
    super(message);
    }
    }
    }

    2.4 锁管理工厂

    package com.ao.c_zk_lock.factory;

    import com.ao.c_zk_lock.lock.AdvancedDistributedLock;
    import com.ao.c_zk_lock.lock.CuratorDistributedLock;
    import org.apache.curator.framework.CuratorFramework;
    import org.apache.curator.framework.CuratorFrameworkFactory;
    import org.apache.curator.retry.ExponentialBackoffRetry;
    import org.slf4j.Logger;
    import org.slf4j.LoggerFactory;

    import java.util.Map;
    import java.util.concurrent.ConcurrentHashMap;
    import java.util.concurrent.atomic.AtomicLong;

    /**
    * 分布式锁管理工厂
    * 统一管理锁实例和ZooKeeper连接
    */
    public class DistributedLockFactory {

    private static final Logger logger = LoggerFactory.getLogger(DistributedLockFactory.class);

    // ZooKeeper客户端缓存
    private static CuratorFramework defaultClient;

    // 锁实例缓存
    private static final Map<String, CuratorDistributedLock> lockCache =
    new ConcurrentHashMap<>();
    private static final Map<String, AdvancedDistributedLock> advancedLockCache =
    new ConcurrentHashMap<>();

    // 锁统计信息
    private static final Map<String, LockStats> lockStatsMap =
    new ConcurrentHashMap<>();

    // 默认配置
    private static final String DEFAULT_ZK_ADDRESS = "192.168.6.100:2181";
    private static final String DEFAULT_NAMESPACE = "aoyou";

    private DistributedLockFactory() {
    // 私有构造函数
    }

    /**
    * 获取默认Curator客户端
    */
    public static synchronized CuratorFramework getDefaultClient() {
    if (defaultClient == null) {
    ExponentialBackoffRetry retryPolicy = new ExponentialBackoffRetry(1000, 3);

    defaultClient = CuratorFrameworkFactory.builder()
    .connectString(DEFAULT_ZK_ADDRESS)
    .namespace(DEFAULT_NAMESPACE)
    .retryPolicy(retryPolicy)
    .sessionTimeoutMs(15000)
    .connectionTimeoutMs(10000)
    .build();

    defaultClient.start();

    logger.info("创建默认Curator客户端: {}", DEFAULT_ZK_ADDRESS);

    // 添加关闭钩子
    Runtime.getRuntime().addShutdownHook(new Thread(() -> {
    shutdownAll();
    System.out.println("分布式锁工厂已关闭");
    }));
    }

    return defaultClient;
    }

    /**
    * 获取或创建基础分布式锁
    */
    public static CuratorDistributedLock getLock(String lockName) {
    return getLock(lockName, CuratorDistributedLock.LockType.MUTEX);
    }

    /**
    * 获取或创建指定类型的分布式锁
    */
    public static CuratorDistributedLock getLock(String lockName, CuratorDistributedLock.LockType lockType) {
    String cacheKey = lockName + ":" + lockType;

    return lockCache.computeIfAbsent(cacheKey, key -> {
    CuratorDistributedLock.LockConfig config =
    new CuratorDistributedLock.LockConfig(lockName)
    .withLockType(lockType);

    CuratorDistributedLock lock = new CuratorDistributedLock(config);

    logger.info("创建分布式锁: {},类型: {}", lockName, lockType);
    return lock;
    });
    }

    /**
    * 获取或创建高级分布式锁
    */
    public static AdvancedDistributedLock getAdvancedLock(String lockName,
    AdvancedDistributedLock.LockType lockType) {
    String cacheKey = lockName + ":" + lockType;

    return advancedLockCache.computeIfAbsent(cacheKey, key -> {
    CuratorFramework client = getDefaultClient();
    AdvancedDistributedLock lock =
    new AdvancedDistributedLock(client, "/locks/" + lockName, lockType);

    logger.info("创建高级分布式锁: {},类型: {}", lockName, lockType);
    return lock;
    });
    }

    /**
    * 获取读写锁
    */
    public static ReadWriteLockPair getReadWriteLock(String lockName) {
    AdvancedDistributedLock readLock = getAdvancedLock(lockName,
    AdvancedDistributedLock.LockType.READ);
    AdvancedDistributedLock writeLock = getAdvancedLock(lockName,
    AdvancedDistributedLock.LockType.WRITE);

    return new ReadWriteLockPair(readLock, writeLock);
    }

    /**
    * 记录锁统计信息
    */
    public static void recordLockStats(String lockPath, long waitTime, boolean success) {
    LockStats stats = lockStatsMap.computeIfAbsent(lockPath,
    key -> new LockStats(lockPath,
    AdvancedDistributedLock.LockType.MUTEX, 0, 0, 0, 0));

    // 这里需要重新设计LockStats类为可变的,或者创建新的统计类
    // 为了简化,这里只记录日志
    logger.debug("锁操作统计 – 路径: {},等待时间: {}ms,成功: {}",
    lockPath, waitTime, success);
    }

    /**
    * 获取所有锁的状态信息
    */
    public static String getAllLockStatus() {
    StringBuilder status = new StringBuilder();
    status.append("========== 分布式锁状态 ==========\\n");

    status.append("基础锁 (").append(lockCache.size()).append("个):\\n");
    for (Map.Entry<String, CuratorDistributedLock> entry : lockCache.entrySet()) {
    try {
    CuratorDistributedLock.LockStatus lockStatus = entry.getValue().getStatus();
    status.append(" ").append(entry.getKey()).append(": ").append(lockStatus).append("\\n");
    } catch (Exception e) {
    status.append(" ").append(entry.getKey()).append(": 获取状态失败 – ")
    .append(e.getMessage()).append("\\n");
    }
    }

    status.append("\\n高级锁 (").append(advancedLockCache.size()).append("个):\\n");
    for (Map.Entry<String, AdvancedDistributedLock> entry : advancedLockCache.entrySet()) {
    status.append(" ").append(entry.getKey())
    .append(": 持有者=")
    .append(entry.getValue().getLockOwner() != null ?
    entry.getValue().getLockOwner().getName() : "无")
    .append(",重入计数=").append(entry.getValue().getHoldCount())
    .append("\\n");
    }

    return status.toString();
    }

    /**
    * 关闭所有锁和连接
    */
    public static void shutdownAll() {
    logger.info("关闭所有分布式锁…");

    // 关闭基础锁
    for (CuratorDistributedLock lock : lockCache.values()) {
    try {
    lock.close();
    } catch (Exception e) {
    logger.error("关闭锁失败", e);
    }
    }
    lockCache.clear();

    // 关闭高级锁
    for (AdvancedDistributedLock lock : advancedLockCache.values()) {
    try {
    // 高级锁没有close方法,这里强制解锁
    if (lock.isLocked()) {
    lock.forceUnlock();
    }
    } catch (Exception e) {
    logger.error("关闭高级锁失败", e);
    }
    }
    advancedLockCache.clear();

    // 关闭客户端连接
    if (defaultClient != null) {
    defaultClient.close();
    defaultClient = null;
    }

    logger.info("所有分布式锁已关闭");
    }

    /**
    * 清理指定的锁
    */
    public static void cleanupLock(String lockName) {
    // 清理基础锁
    lockCache.keySet().removeIf(key -> key.startsWith(lockName + ":"));

    // 清理高级锁
    advancedLockCache.keySet().removeIf(key -> key.startsWith(lockName + ":"));

    logger.info("清理锁: {}", lockName);
    }

    /**
    * 获取锁实例数量
    */
    public static int getLockCount() {
    return lockCache.size() + advancedLockCache.size();
    }

    /**
    * 读写锁对
    */
    public static class ReadWriteLockPair {
    private final AdvancedDistributedLock readLock;
    private final AdvancedDistributedLock writeLock;

    public ReadWriteLockPair(AdvancedDistributedLock readLock, AdvancedDistributedLock writeLock) {
    this.readLock = readLock;
    this.writeLock = writeLock;
    }

    public AdvancedDistributedLock getReadLock() {
    return readLock;
    }

    public AdvancedDistributedLock getWriteLock() {
    return writeLock;
    }

    public <T> T executeWithReadLock(AdvancedDistributedLock.LockableOperation<T> operation) throws Exception {
    return readLock.executeWithLock(operation);
    }

    public <T> T executeWithWriteLock(AdvancedDistributedLock.LockableOperation<T> operation) throws Exception {
    return writeLock.executeWithLock(operation);
    }
    }

    /**
    * 锁统计信息(内部类)
    */
    private static class LockStats {
    private final String lockPath;
    private final AdvancedDistributedLock.LockType lockType;
    private final AtomicLong lockCount = new AtomicLong(0);
    private final AtomicLong waitCount = new AtomicLong(0);
    private final AtomicLong timeoutCount = new AtomicLong(0);
    private final AtomicLong totalWaitTime = new AtomicLong(0);

    public LockStats(String lockPath, AdvancedDistributedLock.LockType lockType,
    long initialCount, long initialWait, long initialTimeout, long initialTime) {
    this.lockPath = lockPath;
    this.lockType = lockType;
    this.lockCount.set(initialCount);
    this.waitCount.set(initialWait);
    this.timeoutCount.set(initialTimeout);
    this.totalWaitTime.set(initialTime);
    }

    public void recordOperation(long waitTime, boolean success) {
    lockCount.incrementAndGet();
    waitCount.incrementAndGet();
    totalWaitTime.addAndGet(waitTime);

    if (!success) {
    timeoutCount.incrementAndGet();
    }
    }

    public double getAverageWaitTime() {
    return waitCount.get() > 0 ?
    (double) totalWaitTime.get() / waitCount.get() : 0.0;
    }

    public double getSuccessRate() {
    return lockCount.get() > 0 ?
    (double) (lockCount.get() – timeoutCount.get()) / lockCount.get() : 0.0;
    }

    @Override
    public String toString() {
    return String.format(
    "LockStats{路径='%s', 加锁次数=%d, 等待次数=%d, 超时次数=%d, 平均等待时间=%.2fms, 成功率=%.2f%%}",
    lockPath, lockCount.get(), waitCount.get(), timeoutCount.get(),
    getAverageWaitTime(), getSuccessRate() * 100
    );
    }
    }
    }

    2.5 订单号生成器示例

    package com.ao.c_zk_lock.generator;

    import com.ao.c_zk_lock.factory.DistributedLockFactory;
    import com.ao.c_zk_lock.lock.CuratorDistributedLock;

    import java.text.SimpleDateFormat;
    import java.util.Date;
    import java.util.concurrent.atomic.AtomicLong;

    /**
    * 订单号生成器
    * 演示如何使用分布式锁保证订单号唯一
    */
    public class OrderCodeGenerator {

    private static final String ORDER_PREFIX = "ORD";
    private static final SimpleDateFormat DATE_FORMAT =
    new SimpleDateFormat("yyyyMMddHHmmss");

    private final AtomicLong sequence = new AtomicLong(0);
    private final String generatorId;

    public OrderCodeGenerator() {
    this.generatorId = "GEN-" + System.currentTimeMillis();
    }

    public OrderCodeGenerator(String generatorId) {
    this.generatorId = generatorId;
    }

    /**
    * 生成订单号(无锁,可能有重复)
    */
    public String generateOrderCode() {
    String timestamp = DATE_FORMAT.format(new Date());
    long seq = sequence.incrementAndGet() % 10000; // 4位序列号

    return String.format("%s-%s-%04d-%s",
    ORDER_PREFIX, timestamp, seq, generatorId);
    }

    /**
    * 生成唯一订单号(使用分布式锁)
    */
    public String generateUniqueOrderCode() throws Exception {
    // 获取分布式锁
    CuratorDistributedLock lock =
    DistributedLockFactory.getLock("order-code-generator");

    try {
    // 获取锁
    if (!lock.lock(5000, java.util.concurrent.TimeUnit.MILLISECONDS)) {
    throw new RuntimeException("获取订单号生成锁超时");
    }

    // 生成订单号
    return generateOrderCode();

    } finally {
    // 释放锁
    if (lock.isHeldByCurrentThread()) {
    lock.unlock();
    }
    }
    }

    /**
    * 批量生成订单号
    */
    public java.util.List<String> batchGenerateOrderCodes(int count) throws Exception {
    java.util.List<String> codes = new java.util.ArrayList<>();

    // 使用锁执行器简化代码
    CuratorDistributedLock lock =
    DistributedLockFactory.getLock("order-code-batch-generator");

    lock.executeWithLock(() -> {
    for (int i = 0; i < count; i++) {
    codes.add(generateOrderCode());
    }
    return null; // 这里不需要返回值
    }, 10000, java.util.concurrent.TimeUnit.MILLISECONDS);

    return codes;
    }

    /**
    * 获取生成器ID
    */
    public String getGeneratorId() {
    return generatorId;
    }

    /**
    * 重置序列号
    */
    public void resetSequence() {
    sequence.set(0);
    }

    /**
    * 获取当前序列号
    */
    public long getCurrentSequence() {
    return sequence.get();
    }
    }

    2.6 单元测试类

    package com.ao.c_zk_lock;

    import com.ao.c_zk_lock.factory.DistributedLockFactory;
    import com.ao.c_zk_lock.generator.OrderCodeGenerator;
    import com.ao.c_zk_lock.lock.AdvancedDistributedLock;
    import com.ao.c_zk_lock.lock.CuratorDistributedLock;
    import org.junit.jupiter.api.AfterAll;
    import org.junit.jupiter.api.BeforeEach;
    import org.junit.jupiter.api.Test;
    import org.junit.jupiter.api.Timeout;
    import org.springframework.boot.test.context.SpringBootTest;

    import java.util.HashSet;
    import java.util.Set;
    import java.util.concurrent.CountDownLatch;
    import java.util.concurrent.ExecutorService;
    import java.util.concurrent.Executors;
    import java.util.concurrent.TimeUnit;
    import java.util.concurrent.atomic.AtomicInteger;

    import static org.junit.jupiter.api.Assertions.assertSame;
    import static org.springframework.test.util.AssertionErrors.*;

    /**
    * Curator分布式锁单元测试
    */
    @SpringBootTest
    public class CuratorDistributedLockTest {

    private OrderCodeGenerator orderCodeGenerator;

    @BeforeEach
    public void setUp() {
    orderCodeGenerator = new OrderCodeGenerator("test-generator");
    }

    @AfterAll
    public static void tearDown() {
    DistributedLockFactory.shutdownAll();
    }

    @Test
    public void testBasicLockOperations() throws Exception {
    System.out.println("=== 测试基础锁操作 ===");

    // 获取锁实例
    CuratorDistributedLock lock = DistributedLockFactory.getLock("test-basic-lock");

    // 获取锁状态
    CuratorDistributedLock.LockStatus status = lock.getStatus();
    System.out.println("锁初始状态: " + status);

    assertFalse("初始状态应该未锁定", status.isLocked);
    assertEquals("重入计数应该为0", 0, status.holdCount);

    // 测试加锁和解锁
    boolean locked = lock.lock(1000, TimeUnit.MILLISECONDS);
    assertTrue("应该成功获取锁", locked);

    status = lock.getStatus();
    assertTrue("加锁后应该显示已锁定", status.isLocked);
    assertTrue("当前线程应该持有锁", status.isHeldByCurrentThread);
    assertEquals("重入计数应该为1", 1, status.holdCount);

    // 测试重入
    locked = lock.lock(1000, TimeUnit.MILLISECONDS);
    assertTrue("应该成功重入锁", locked);

    status = lock.getStatus();
    assertEquals("重入后计数应该为2", 2, status.holdCount);

    // 测试解锁
    lock.unlock();
    status = lock.getStatus();
    assertEquals("解锁后计数应该为1", 1, status.holdCount);

    lock.unlock();
    status = lock.getStatus();
    assertFalse("完全解锁后应该未锁定", status.isLocked);
    assertEquals("完全解锁后计数应该为0", 0, status.holdCount);

    System.out.println("基础锁操作测试通过");
    }

    @Test
    public void testOrderCodeGeneratorWithLock() throws Exception {
    System.out.println("=== 测试带锁的订单号生成 ===");

    Set<String> orderCodes = new HashSet<>();
    int testCount = 100;

    // 生成订单号(使用锁保证唯一性)
    for (int i = 0; i < testCount; i++) {
    String orderCode = orderCodeGenerator.generateUniqueOrderCode();

    // 检查唯一性
    if (orderCodes.contains(orderCode)) {
    fail("订单号重复: " + orderCode);
    }

    orderCodes.add(orderCode);

    if ((i + 1) % 20 == 0) {
    System.out.println("已生成 " + (i + 1) + " 个唯一订单号");
    }
    }

    assertEquals("应该生成指定数量的订单号", testCount, orderCodes.size());
    System.out.println("订单号生成测试通过,生成数量: " + orderCodes.size());
    }

    @Test
    public void testConcurrentOrderCodeGeneration() throws Exception {
    System.out.println("=== 测试并发订单号生成 ===");

    int threadCount = 20;
    int codesPerThread = 10;

    Set<String> allOrderCodes = new HashSet<>();
    CountDownLatch latch = new CountDownLatch(threadCount);
    AtomicInteger duplicateCount = new AtomicInteger(0);

    ExecutorService executor = Executors.newFixedThreadPool(threadCount);

    for (int i = 0; i < threadCount; i++) {
    final int threadId = i;
    executor.submit(() -> {
    try {
    for (int j = 0; j < codesPerThread; j++) {
    String orderCode = orderCodeGenerator.generateUniqueOrderCode();

    synchronized (allOrderCodes) {
    if (allOrderCodes.contains(orderCode)) {
    duplicateCount.incrementAndGet();
    System.err.println("发现重复订单号: " + orderCode);
    } else {
    allOrderCodes.add(orderCode);
    }
    }

    // 稍微延迟
    Thread.sleep(10);
    }

    System.out.println("线程 " + threadId + " 完成");

    } catch (Exception e) {
    e.printStackTrace();
    } finally {
    latch.countDown();
    }
    });
    }

    latch.await(30, TimeUnit.SECONDS);
    executor.shutdown();

    int expectedCount = threadCount * codesPerThread;
    int actualCount = allOrderCodes.size();
    int duplicates = duplicateCount.get();

    System.out.println("并发测试结果:");
    System.out.println(" 期望数量: " + expectedCount);
    System.out.println(" 实际数量: " + actualCount);
    System.out.println(" 重复数量: " + duplicates);
    System.out.println(" 成功率: " + (actualCount * 100.0 / expectedCount) + "%");

    assertEquals("不应该有重复订单号", 0, duplicates);
    assertEquals("应该生成所有订单号", expectedCount, actualCount);

    System.out.println("并发订单号生成测试通过");
    }

    @Test
    public void testLockTimeout() throws Exception {
    System.out.println("=== 测试锁超时 ===");

    CuratorDistributedLock lock1 = DistributedLockFactory.getLock("timeout-test-lock");
    CuratorDistributedLock lock2 = DistributedLockFactory.getLock("timeout-test-lock");

    // 线程1获取锁
    Thread thread1 = new Thread(() -> {
    try {
    boolean locked = lock1.lock();
    if (locked) {
    System.out.println("线程1获取锁成功,持有10秒");
    Thread.sleep(10000); // 持有10秒
    lock1.unlock();
    System.out.println("线程1释放锁");
    }
    } catch (Exception e) {
    e.printStackTrace();
    }
    });

    // 线程2尝试获取锁(应该超时)
    Thread thread2 = new Thread(() -> {
    try {
    System.out.println("线程2尝试获取锁(超时3秒)");
    long startTime = System.currentTimeMillis();
    boolean locked = lock2.lock(3000, TimeUnit.MILLISECONDS);
    long elapsedTime = System.currentTimeMillis() – startTime;

    if (locked) {
    System.out.println("线程2获取锁成功,耗时: " + elapsedTime + "ms");
    lock2.unlock();
    } else {
    System.out.println("线程2获取锁超时,耗时: " + elapsedTime + "ms");
    }

    assertFalse("线程2不应该获取到锁", locked);
    assertTrue("应该等待大约3秒", elapsedTime >= 2800 && elapsedTime <= 3500);

    } catch (Exception e) {
    e.printStackTrace();
    }
    });

    thread1.start();
    Thread.sleep(100); // 确保线程1先启动

    thread2.start();

    thread1.join();
    thread2.join();

    System.out.println("锁超时测试通过");
    }

    @Test
    public void testLockExecuteWithLock() throws Exception {
    System.out.println("=== 测试锁执行器 ===");

    CuratorDistributedLock lock = DistributedLockFactory.getLock("execute-test-lock");

    // 使用锁执行器执行任务
    String result = lock.executeWithLock(() -> {
    System.out.println("在锁保护下执行任务");
    return "任务结果";
    }, 5000, TimeUnit.MILLISECONDS);

    assertEquals("应该返回任务结果", "任务结果", result);

    // 验证锁已释放
    CuratorDistributedLock.LockStatus status = lock.getStatus();
    assertFalse("锁执行器应该自动释放锁", status.isLocked);

    System.out.println("锁执行器测试通过");
    }

    @Test
    public void testBatchOrderCodeGeneration() throws Exception {
    System.out.println("=== 测试批量订单号生成 ===");

    int batchSize = 50;

    // 批量生成订单号
    java.util.List<String> orderCodes =
    orderCodeGenerator.batchGenerateOrderCodes(batchSize);

    // 检查唯一性
    Set<String> uniqueCodes = new HashSet<>(orderCodes);

    assertEquals("批量生成的订单号应该唯一", batchSize, uniqueCodes.size());
    assertEquals("应该生成指定数量的订单号", batchSize, orderCodes.size());

    System.out.println("批量生成了 " + orderCodes.size() + " 个唯一订单号");
    System.out.println("前5个订单号:");
    for (int i = 0; i < Math.min(5, orderCodes.size()); i++) {
    System.out.println(" " + orderCodes.get(i));
    }

    System.out.println("批量订单号生成测试通过");
    }

    @Test
    public void testLockFactory() {
    System.out.println("=== 测试锁工厂 ===");

    // 获取多个锁实例
    CuratorDistributedLock lock1 = DistributedLockFactory.getLock("factory-test-1");
    CuratorDistributedLock lock2 = DistributedLockFactory.getLock("factory-test-2");
    AdvancedDistributedLock lock3 =
    DistributedLockFactory.getAdvancedLock("factory-test-3",
    AdvancedDistributedLock.LockType.MUTEX);

    // 验证锁实例
    assertNotNull("锁1不应该为null", lock1);
    assertNotNull("锁2不应该为null", lock2);
    assertNotNull("锁3不应该为null", lock3);

    // 验证单例模式
    CuratorDistributedLock sameLock1 = DistributedLockFactory.getLock("factory-test-1");
    assertSame(lock1, sameLock1, "相同名称的锁应该是同一个实例");

    // 获取所有锁状态
    String status = DistributedLockFactory.getAllLockStatus();
    assertNotNull("状态信息不应该为null", status);
    assertTrue("状态信息应该包含锁信息", status.contains("factory-test-1"));

    System.out.println("锁工厂状态:\\n" + status);
    System.out.println("锁总数: " + DistributedLockFactory.getLockCount());

    System.out.println("锁工厂测试通过");
    }

    @Test
    @Timeout(10)
    public void testLockPerformance() throws Exception {
    System.out.println("=== 测试锁性能 ===");

    CuratorDistributedLock lock = DistributedLockFactory.getLock("performance-test-lock");

    int testCount = 100;

    // 预热
    for (int i = 0; i < 10; i++) {
    lock.lock();
    lock.unlock();
    }

    // 性能测试
    long startTime = System.currentTimeMillis();

    for (int i = 0; i < testCount; i++) {
    lock.lock();
    lock.unlock();

    if (i % 20 == 0 && i > 0) {
    long currentTime = System.currentTimeMillis();
    double ops = i * 1000.0 / (currentTime – startTime);
    System.out.printf("已执行 %d 次加锁/解锁,当前OPS: %.2f\\n", i, ops);
    }
    }

    long endTime = System.currentTimeMillis();
    long duration = endTime – startTime;

    double ops = testCount * 1000.0 / duration;

    System.out.println("性能测试结果:");
    System.out.println(" 总操作数: " + testCount * 2 + " (加锁+解锁)");
    System.out.println(" 总耗时: " + duration + "ms");
    System.out.println(" 平均OPS: " + String.format("%.2f", ops));

    assertTrue("锁操作OPS应该较高", ops > 50);

    System.out.println("锁性能测试通过");
    }
    }

    运行结果

    运行截图

    三、使用示例

    package com.ao.c_zk_lock.example;

    import com.ao.c_zk_lock.factory.DistributedLockFactory;
    import com.ao.c_zk_lock.generator.OrderCodeGenerator;
    import com.ao.c_zk_lock.lock.CuratorDistributedLock;

    import java.util.concurrent.TimeUnit;

    /**
    * Curator分布式锁使用示例
    */
    public class CuratorLockExample {

    public static void main(String[] args) throws Exception {
    System.out.println("=== Curator分布式锁使用示例 ===\\n");

    try {
    // 示例1:基本锁使用
    basicLockExample();

    // 示例2:订单号生成示例
    orderCodeExample();

    // 示例3:读写锁示例
    readWriteLockExample();

    // 示例4:锁执行器示例
    lockExecutorExample();

    } finally {
    // 清理资源
    DistributedLockFactory.shutdownAll();
    }
    }

    /**
    * 基本锁使用示例
    */
    private static void basicLockExample() throws Exception {
    System.out.println("1. 基本锁使用示例:");

    // 获取锁实例
    CuratorDistributedLock lock = DistributedLockFactory.getLock("basic-example-lock");

    // 获取锁
    boolean locked = lock.lock(5000, TimeUnit.MILLISECONDS);
    if (locked) {
    try {
    System.out.println(" 获取锁成功,执行关键业务…");

    // 模拟业务处理
    Thread.sleep(1000);

    System.out.println(" 业务处理完成");

    } finally {
    // 释放锁
    lock.unlock();
    System.out.println(" 释放锁");
    }
    } else {
    System.out.println(" 获取锁超时");
    }

    // 查看锁状态
    CuratorDistributedLock.LockStatus status = lock.getStatus();
    System.out.println(" 锁状态: " + status);
    }

    /**
    * 订单号生成示例
    */
    private static void orderCodeExample() throws Exception {
    System.out.println("\\n2. 订单号生成示例:");

    OrderCodeGenerator generator = new OrderCodeGenerator("example-generator");

    // 生成单个订单号
    String orderCode1 = generator.generateUniqueOrderCode();
    String orderCode2 = generator.generateUniqueOrderCode();

    System.out.println(" 订单号1: " + orderCode1);
    System.out.println(" 订单号2: " + orderCode2);

    // 验证唯一性
    if (!orderCode1.equals(orderCode2)) {
    System.out.println(" ✓ 订单号唯一性验证通过");
    } else {
    System.err.println(" ✗ 订单号重复!");
    }

    // 批量生成订单号
    System.out.println("\\n 批量生成订单号:");
    java.util.List<String> batchCodes = generator.batchGenerateOrderCodes(5);
    for (int i = 0; i < batchCodes.size(); i++) {
    System.out.println(" 订单" + (i + 1) + ": " + batchCodes.get(i));
    }

    System.out.println(" 生成器ID: " + generator.getGeneratorId());
    System.out.println(" 当前序列号: " + generator.getCurrentSequence());
    }

    /**
    * 读写锁示例
    */
    private static void readWriteLockExample() throws Exception {
    System.out.println("\\n3. 读写锁示例:");

    // 获取读写锁
    DistributedLockFactory.ReadWriteLockPair rwLock =
    DistributedLockFactory.getReadWriteLock("rw-example");

    // 读锁示例(可以并发)
    System.out.println(" 读锁示例(多线程并发读):");
    for (int i = 0; i < 3; i++) {
    final int threadId = i;
    new Thread(() -> {
    try {
    String result = rwLock.executeWithReadLock(() -> {
    System.out.println(" 线程" + threadId + ": 读取数据");
    Thread.sleep(500);
    return "读取结果-" + threadId;
    });
    System.out.println(" 线程" + threadId + ": " + result);
    } catch (Exception e) {
    e.printStackTrace();
    }
    }).start();
    }

    Thread.sleep(2000);

    // 写锁示例(互斥)
    System.out.println("\\n 写锁示例(互斥写):");
    for (int i = 0; i < 2; i++) {
    final int threadId = i;
    new Thread(() -> {
    try {
    String result = rwLock.executeWithWriteLock(() -> {
    System.out.println(" 写线程" + threadId + ": 写入数据");
    Thread.sleep(1000);
    return "写入结果-" + threadId;
    });
    System.out.println(" 写线程" + threadId + ": " + result);
    } catch (Exception e) {
    e.printStackTrace();
    }
    }).start();
    }

    Thread.sleep(3000);
    }

    /**
    * 锁执行器示例
    */
    private static void lockExecutorExample() throws Exception {
    System.out.println("\\n4. 锁执行器示例:");

    CuratorDistributedLock lock = DistributedLockFactory.getLock("executor-example");

    // 使用锁执行器执行任务
    System.out.println(" 使用锁执行器:");
    Integer result = lock.executeWithLock(() -> {
    System.out.println(" 在锁保护下执行计算任务");
    int sum = 0;
    for (int i = 1; i <= 100; i++) {
    sum += i;
    }
    return sum;
    }, 10000, TimeUnit.MILLISECONDS);

    System.out.println(" 计算结果: " + result);
    System.out.println(" 锁状态: " + lock.getStatus());

    // 测试超时情况
    System.out.println("\\n 测试锁超时:");

    // 先让一个线程持有锁
    Thread holder = new Thread(() -> {
    try {
    CuratorDistributedLock sameLock =
    DistributedLockFactory.getLock("executor-example");
    sameLock.lock();
    System.out.println(" 线程A获取锁,持有3秒");
    Thread.sleep(3000);
    sameLock.unlock();
    System.out.println(" 线程A释放锁");
    } catch (Exception e) {
    e.printStackTrace();
    }
    });

    holder.start();
    Thread.sleep(100); // 确保holder先启动

    // 另一个线程尝试获取锁(应该超时)
    try {
    lock.executeWithLock(() -> {
    System.out.println(" 线程B获取锁成功(不应该执行到这里)");
    return null;
    }, 1000, TimeUnit.MILLISECONDS);

    System.err.println(" ✗ 应该抛出超时异常!");

    } catch (CuratorDistributedLock.LockTimeoutException e) {
    System.out.println(" ✓ 正确抛出超时异常: " + e.getMessage());
    }

    holder.join();
    }
    }

    四、配置文件和注意事项

    4.1 配置文件示例

    # curator-lock.properties
    # ZooKeeper配置
    zookeeper.address=192.168.6.100:2181
    zookeeper.namespace=aoyou
    zookeeper.session.timeout=15000
    zookeeper.connection.timeout=10000

    # 重试策略
    retry.base.sleep.time=1000
    retry.max.retries=3

    # 锁配置
    lock.default.timeout=5000
    lock.max.wait.time=30000

    # 业务锁配置
    lock.order.generator.path=/locks/order-generator
    lock.order.generator.timeout=3000
    lock.payment.processing.path=/locks/payment-processing
    lock.payment.processing.timeout=10000

    # 监控配置
    lock.monitor.enabled=true
    lock.monitor.interval=10000

    4.2 注意事项

  • Curator客户端管理:
      • Curator客户端为重量级对象,建议在应用中作为单例使用(如Spring Bean的@Singleton),避免频繁创建/关闭连接。
      • 生产环境中应配置合理的重试策略(如ExponentialBackoffRetry的重试次数建议设为3-5次,初始间隔设为1000ms),提高连接可靠性。
  • 锁路径设计:
      • 锁路径应按业务维度划分命名空间(如/aoyou/order_lock、/aoyou/inventory_lock),避免不同业务的锁冲突。
      • 锁路径的ZooKeeper节点为临时节点,Curator会自动管理节点的创建和删除,无需手动处理。
  • 锁的超时处理:
      • InterProcessMutex的acquire()方法支持超时参数(如acquire(10, TimeUnit.SECONDS)),建议在生产环境中使用超时机制,避免线程无限阻塞。
      • 若获取锁超时,应根据业务场景处理(如重试、抛出异常、降级处理等)。
  • 可重入性验证:
      • InterProcessMutex支持可重入,同一线程多次调用acquire()时需调用相同次数的release()才能完全释放锁,需注意在finally块中确保释放次数与获取次数一致,避免锁泄漏。
  • 测试环境优化:
      • 单元测试中可使用Curator的TestingServer(嵌入式ZooKeeper),避免依赖外部ZooKeeper服务,提高测试的可重复性和独立性。

    五、总结

    Curator分布式锁提供了比原生ZooKeeper锁更强大和易用的功能:

  • 易于使用:简洁的API设计,支持重入锁
  • 功能丰富:支持读写锁、信号量等多种锁类型
  • 可靠性高:完善的错误处理和重试机制
  • 性能优越:优化的锁实现,减少网络开销
  • 监控完善:详细的锁状态信息
  • 赞(0)
    未经允许不得转载:171主机测评 » 基于Curator实现分布式锁
    分享到: 更多 (0)

    评论 抢沙发

    • 昵称 (必填)
    • 邮箱 (必填)
    • 网址