🌺The Begin🌺点点关注,收藏不迷路🌺

摘要:在分布式系统中,多个进程或线程需要互斥访问共享资源时,传统的线程锁已经无能为力,必须引入分布式锁。ZooKeeper 凭借其顺序节点、临时节点和 Watcher 机制,成为了实现分布式锁的理想选择。本文将深入剖析基于 ZooKeeper 实现分布式锁的原理,包括排他锁和读写锁的实现方式,并通过完整代码示例帮助读者掌握这一关键技术。

一、分布式锁概述

1.1 什么是分布式锁?

分布式锁是一种用于协调分布式系统中多个进程对共享资源进行互斥访问的机制。在分布式环境下,传统的基于内存的锁(如 Java 的 synchronizedReentrantLock)无法跨进程工作,因此需要借助外部协调服务来实现。

分布式锁应用场景

请求锁

请求锁

请求锁

获得锁

访问共享资源

等待

等待

进程1

分布式锁服务

进程2

进程3

共享资源

1.2 分布式锁的基本要求

要求 说明
互斥性 任何时刻只能有一个客户端持有锁
避免死锁 即使持有锁的客户端崩溃,锁也能自动释放
可重入性 同一个客户端可以多次获取同一把锁
高性能 获取和释放锁的开销要小
容错性 锁服务本身要高可用

1.3 为什么选择 ZooKeeper 实现分布式锁?

特性 作用
临时节点 客户端崩溃时自动释放锁,避免死锁
顺序节点 实现公平锁,按请求顺序获取锁
Watcher 机制 客户端可以监听锁状态,避免频繁轮询
强一致性 ZAB 协议保证锁状态一致

二、基于 ZooKeeper 的排他锁实现

排他锁(Exclusive Lock)是最简单的分布式锁,也称为写锁。它保证同一时刻只有一个客户端能持有锁。

2.1 实现原理

排他锁实现流程

客户端请求锁

创建临时顺序节点
/locks/lock-

获取/locks下所有子节点

自己是否是最小节点?

获得锁,执行业务

监听前一个节点

前一个节点释放

2.2 完整代码实现

import org.apache.zookeeper.*;
import org.apache.zookeeper.data.Stat;
import java.io.IOException;
import java.util.Collections;
import java.util.List;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.TimeUnit;

public class ZooKeeperDistributedLock implements Watcher {
    private ZooKeeper zk;
    private String lockPath = "/locks";
    private String lockName = "lock-";
    private String currentNodePath;
    private String waitNodePath;
    private CountDownLatch connectedLatch = new CountDownLatch(1);
    private CountDownLatch lockLatch;
    
    public ZooKeeperDistributedLock(String connectString, String lockPath) throws IOException {
        this.lockPath = lockPath;
        this.zk = new ZooKeeper(connectString, 5000, this);
        try {
            connectedLatch.await();
            // 确保锁根节点存在
            Stat stat = zk.exists(lockPath, false);
            if (stat == null) {
                zk.create(lockPath, new byte[0], ZooDefs.Ids.OPEN_ACL_UNSAFE, CreateMode.PERSISTENT);
            }
        } catch (Exception e) {
            e.printStackTrace();
        }
    }
    
    @Override
    public void process(WatchedEvent event) {
        if (event.getState() == Event.KeeperState.SyncConnected) {
            connectedLatch.countDown();
        }
        
        // 如果监听到前一个节点被删除,释放锁等待
        if (event.getType() == Event.EventType.NodeDeleted && event.getPath().equals(waitNodePath)) {
            if (lockLatch != null) {
                lockLatch.countDown();
            }
        }
    }
    
    /**
     * 尝试获取锁(阻塞式)
     */
    public boolean lock() throws Exception {
        if (tryLock()) {
            return true;
        }
        return waitForLock();
    }
    
    /**
     * 尝试获取锁(非阻塞式)
     */
    public boolean tryLock() throws Exception {
        // 创建临时顺序节点
        currentNodePath = zk.create(lockPath + "/" + lockName, new byte[0],
                ZooDefs.Ids.OPEN_ACL_UNSAFE, CreateMode.EPHEMERAL_SEQUENTIAL);
        
        // 获取所有子节点
        List<String> children = zk.getChildren(lockPath, false);
        Collections.sort(children);
        
        // 判断自己是否是最小节点
        String currentNodeName = currentNodePath.substring(lockPath.length() + 1);
        int index = children.indexOf(currentNodeName);
        
        if (index == 0) {
            // 是最小节点,获得锁
            return true;
        } else {
            // 不是最小节点,监听前一个节点
            waitNodePath = lockPath + "/" + children.get(index - 1);
            return false;
        }
    }
    
    /**
     * 等待锁
     */
    private boolean waitForLock() throws Exception {
        lockLatch = new CountDownLatch(1);
        
        // 检查前一个节点是否存在
        Stat stat = zk.exists(waitNodePath, this);
        if (stat == null) {
            // 前一个节点已不存在,重新尝试获取锁
            return tryLock();
        }
        
        // 等待前一个节点释放
        lockLatch.await();
        
        // 再次尝试获取锁
        return tryLock();
    }
    
    /**
     * 释放锁
     */
    public void unlock() throws Exception {
        if (currentNodePath != null) {
            zk.delete(currentNodePath, -1);
            currentNodePath = null;
        }
    }
    
    /**
     * 带超时的尝试获取锁
     */
    public boolean tryLock(long timeout, TimeUnit unit) throws Exception {
        if (tryLock()) {
            return true;
        }
        return waitForLockTimeout(timeout, unit);
    }
    
    private boolean waitForLockTimeout(long timeout, TimeUnit unit) throws Exception {
        lockLatch = new CountDownLatch(1);
        
        Stat stat = zk.exists(waitNodePath, this);
        if (stat == null) {
            return tryLock();
        }
        
        boolean awaitResult = lockLatch.await(timeout, unit);
        if (!awaitResult) {
            // 超时,删除自己创建的节点
            unlock();
            return false;
        }
        
        return tryLock();
    }
}

2.3 使用示例

public class DistributedLockExample {
    private static final String ZK_SERVER = "localhost:2181";
    private static final String LOCK_PATH = "/distributed-locks";
    
    public static void main(String[] args) throws Exception {
        // 创建两个客户端模拟竞争
        ZooKeeperDistributedLock lock1 = new ZooKeeperDistributedLock(ZK_SERVER, LOCK_PATH);
        ZooKeeperDistributedLock lock2 = new ZooKeeperDistributedLock(ZK_SERVER, LOCK_PATH);
        
        // 线程1获取锁
        new Thread(() -> {
            try {
                System.out.println("线程1尝试获取锁");
                if (lock1.lock()) {
                    System.out.println("线程1获得锁,执行业务...");
                    Thread.sleep(5000); // 模拟业务处理
                    lock1.unlock();
                    System.out.println("线程1释放锁");
                }
            } catch (Exception e) {
                e.printStackTrace();
            }
        }).start();
        
        // 线程2获取锁
        new Thread(() -> {
            try {
                Thread.sleep(1000); // 确保线程1先执行
                System.out.println("线程2尝试获取锁");
                if (lock2.lock()) {
                    System.out.println("线程2获得锁,执行业务...");
                    lock2.unlock();
                    System.out.println("线程2释放锁");
                }
            } catch (Exception e) {
                e.printStackTrace();
            }
        }).start();
        
        Thread.sleep(10000);
    }
}

三、基于 ZooKeeper 的读写锁实现

读写锁(Read/Write Lock)比排他锁更精细,允许多个客户端同时读,但写操作必须互斥。

3.1 读写锁规则

锁类型 规则
读锁 如果没有写锁(没有以 “write-” 开头的节点且序号更小),可以获得读锁
写锁 如果自己是最小的节点(没有序号更小的任何节点),可以获得写锁

3.2 实现原理

读写锁判断逻辑

读锁

写锁

有更小写节点

没有

有更小节点

没有

创建临时顺序节点

锁类型判断

检查是否有写节点比自己小

检查是否有任何节点比自己小

监听最小写节点

获得读锁

监听最小节点

获得写锁

3.3 读写锁代码实现

import org.apache.zookeeper.*;
import org.apache.zookeeper.data.Stat;
import java.io.IOException;
import java.util.Collections;
import java.util.List;
import java.util.concurrent.CountDownLatch;

public class ZooKeeperReadWriteLock implements Watcher {
    private ZooKeeper zk;
    private String lockPath = "/rwlocks";
    private String currentNodePath;
    private CountDownLatch connectedLatch = new CountDownLatch(1);
    private CountDownLatch lockLatch;
    private LockType lockType;
    private String waitNodePath;
    
    public enum LockType {
        READ, WRITE
    }
    
    public ZooKeeperReadWriteLock(String connectString, String lockPath) throws IOException {
        this.lockPath = lockPath;
        this.zk = new ZooKeeper(connectString, 5000, this);
        try {
            connectedLatch.await();
            Stat stat = zk.exists(lockPath, false);
            if (stat == null) {
                zk.create(lockPath, new byte[0], ZooDefs.Ids.OPEN_ACL_UNSAFE, CreateMode.PERSISTENT);
            }
        } catch (Exception e) {
            e.printStackTrace();
        }
    }
    
    @Override
    public void process(WatchedEvent event) {
        if (event.getState() == Event.KeeperState.SyncConnected) {
            connectedLatch.countDown();
        }
        
        if (event.getType() == Event.EventType.NodeDeleted && event.getPath().equals(waitNodePath)) {
            if (lockLatch != null) {
                lockLatch.countDown();
            }
        }
    }
    
    /**
     * 获取读锁
     */
    public boolean readLock() throws Exception {
        this.lockType = LockType.READ;
        return acquireLock();
    }
    
    /**
     * 获取写锁
     */
    public boolean writeLock() throws Exception {
        this.lockType = LockType.WRITE;
        return acquireLock();
    }
    
    private boolean acquireLock() throws Exception {
        String prefix = lockType == LockType.READ ? "read-" : "write-";
        currentNodePath = zk.create(lockPath + "/" + prefix, new byte[0],
                ZooDefs.Ids.OPEN_ACL_UNSAFE, CreateMode.EPHEMERAL_SEQUENTIAL);
        
        if (checkCanLock()) {
            return true;
        }
        
        return waitForLock();
    }
    
    private boolean checkCanLock() throws Exception {
        List<String> children = zk.getChildren(lockPath, false);
        Collections.sort(children);
        
        String currentNodeName = currentNodePath.substring(lockPath.length() + 1);
        int index = children.indexOf(currentNodeName);
        
        // 检查是否有比自己小的节点
        for (int i = 0; i < index; i++) {
            String node = children.get(i);
            
            if (lockType == LockType.WRITE) {
                // 写锁:只要前面有任意节点,就不能获得锁
                waitNodePath = lockPath + "/" + node;
                return false;
            } else {
                // 读锁:如果前面有写节点,不能获得锁
                if (node.startsWith("write-")) {
                    waitNodePath = lockPath + "/" + node;
                    return false;
                }
            }
        }
        
        // 没有需要等待的节点
        return true;
    }
    
    private boolean waitForLock() throws Exception {
        lockLatch = new CountDownLatch(1);
        
        Stat stat = zk.exists(waitNodePath, this);
        if (stat == null) {
            // 等待的节点已不存在,重新检查
            return checkCanLock();
        }
        
        lockLatch.await();
        return checkCanLock();
    }
    
    public void unlock() throws Exception {
        if (currentNodePath != null) {
            zk.delete(currentNodePath, -1);
            currentNodePath = null;
        }
    }
}

四、锁的释放与故障处理

4.1 正常释放

客户端主动调用 unlock() 方法删除自己创建的节点:

public void unlock() throws Exception {
    if (currentNodePath != null) {
        // 删除临时节点
        zk.delete(currentNodePath, -1);
        currentNodePath = null;
    }
}

4.2 异常释放

如果持有锁的客户端崩溃或网络断开,ZooKeeper 的临时节点会在会话超时后自动删除:

其他等待客户端 ZooKeeper 持有锁的客户端 其他等待客户端 ZooKeeper 持有锁的客户端 客户端崩溃 创建临时顺序节点 /lock-000001 获得锁 等待会话超时 删除临时节点 /lock-000001 发送 Watcher 通知 重新检查锁状态

4.3 避免死锁的保障

机制 作用
临时节点 客户端崩溃时自动删除,锁自动释放
会话超时 网络分区时,锁最终会被释放
版本控制 防止误删其他客户端的锁

五、完整示例:订单服务中的分布式锁应用

5.1 业务场景

在电商系统中,需要防止用户重复下单,可以使用分布式锁对同一用户的订单创建进行互斥。

public class OrderService {
    private ZooKeeperDistributedLock lock;
    private String userId;
    
    public OrderService(String userId, String zkServer) throws IOException {
        this.userId = userId;
        // 每个用户有自己的锁路径
        this.lock = new ZooKeeperDistributedLock(zkServer, "/order-locks/" + userId);
    }
    
    /**
     * 创建订单
     */
    public boolean createOrder(String orderInfo) throws Exception {
        boolean locked = false;
        try {
            // 尝试获取锁,超时时间3秒
            locked = lock.tryLock(3, TimeUnit.SECONDS);
            
            if (!locked) {
                System.out.println("用户 " + userId + " 正在创建订单,请稍后");
                return false;
            }
            
            // 执行业务逻辑
            System.out.println("用户 " + userId + " 开始创建订单: " + orderInfo);
            Thread.sleep(2000); // 模拟业务处理
            System.out.println("用户 " + userId + " 订单创建成功");
            
            return true;
            
        } finally {
            if (locked) {
                lock.unlock();
            }
        }
    }
    
    public static void main(String[] args) {
        String zkServer = "localhost:2181";
        
        // 模拟同一个用户的并发下单
        for (int i = 0; i < 5; i++) {
            final int threadId = i;
            new Thread(() -> {
                try {
                    OrderService orderService = new OrderService("user-123", zkServer);
                    orderService.createOrder("订单-" + threadId);
                } catch (Exception e) {
                    e.printStackTrace();
                }
            }).start();
        }
    }
}

六、ZooKeeper 分布式锁的优缺点

6.1 优点

优点 说明
可靠性高 ZooKeeper 集群保证锁状态一致
自动释放 临时节点机制避免死锁
公平锁 顺序节点保证请求顺序
支持读写锁 可以实现更细粒度的锁控制
事件通知 Watcher 机制避免轮询

6.2 缺点

缺点 说明
性能开销 每次锁操作需要与 ZooKeeper 交互
羊群效应 大量客户端同时唤醒可能导致性能波动
依赖外部服务 ZooKeeper 本身也需要高可用部署
复杂性 比 Redis 分布式锁实现更复杂

6.3 与其他分布式锁的对比

对比维度 ZooKeeper 锁 Redis 锁 数据库锁
一致性 强一致 最终一致 强一致
性能 中等
可靠性 中等
自动释放 ✅(需设置超时)

七、最佳实践建议

7.1 锁路径命名规范

/locks/{业务}/{资源ID}
示例:
/locks/order/user-123      # 用户订单锁
/locks/stock/product-456    # 商品库存锁
/locks/seckill/activity-789 # 秒杀活动锁

7.2 超时设置

// 合理设置锁超时时间
public boolean tryLock(long timeout, TimeUnit unit) {
    // 超时时间应根据业务执行时间估算
    // 通常设置为业务执行时间的3-5倍
}

7.3 异常处理

public void doWithLock(Runnable task) {
    boolean locked = false;
    try {
        locked = lock.tryLock(5, TimeUnit.SECONDS);
        if (!locked) {
            throw new RuntimeException("获取锁超时");
        }
        
        task.run();
        
    } catch (InterruptedException e) {
        Thread.currentThread().interrupt();
        throw new RuntimeException("线程被中断", e);
    } catch (Exception e) {
        throw new RuntimeException("业务执行异常", e);
    } finally {
        if (locked) {
            try {
                lock.unlock();
            } catch (Exception e) {
                // 记录日志但不要抛出,避免掩盖业务异常
                e.printStackTrace();
            }
        }
    }
}

八、总结

8.1 核心要点回顾

  1. 利用临时顺序节点:实现公平锁和自动释放
  2. 利用 Watcher 机制:避免轮询,提高效率
  3. 利用最小节点判断:确定锁的持有者
  4. 利用节点前缀区分:实现读写锁
  5. 利用版本控制:防止误删锁

8.2 实现流程图

请求获取锁

创建临时顺序节点

获取所有子节点

排序比较

是否是最小节点

获得锁
执行业务

释放锁
删除节点

监听前一个节点

前一个节点删除

8.3 一句话总结

基于 ZooKeeper 的分布式锁,通过临时顺序节点实现公平性和自动释放,利用Watcher 机制实现事件通知,结合节点前缀和排序实现读写锁,最终构建出一个高可靠、可扩展的分布式互斥解决方案。

在这里插入图片描述


🌺The End🌺点点关注,收藏不迷路🌺
Logo

汇聚全球AI编程工具,助力开发者即刻编程。

更多推荐