Skip to content

ZooKeeper分布式锁与选主全过程:前驱Watch、Session失效和Fencing Token

ZooKeeper 临时顺序节点能构造排队锁和 Leader 选举,但“临时节点会自动删除”不等于外部资源绝不并发写。客户端可能在长时间 GC 后恢复,旧 Session 已经失效,新客户端已经获得锁;如果数据库仍接受旧客户端写入,就会出现两个实际执行者。

因此完整正确性链必须包含:顺序排队、前驱 Watch、Session 状态、安全退出、业务幂等,以及由外部资源校验的 Fencing Token。

一、学习目标

  1. 从临时顺序节点推导分布式互斥过程。
  2. 解释为什么只监听前驱节点可以避免羊群效应。
  3. 处理前驱节点在 Watch 注册前已删除的竞态。
  4. 区分锁节点删除、客户端线程恢复和外部资源所有权。
  5. 解释 Fencing Token 必须在哪里校验。
  6. 区分公平排队、严格公平和高吞吐锁。
  7. 区分 Leader 身份、任务执行和Exactly Once。
  8. 使用 Curator Recipe并理解它没有替业务解决什么。

二、一个生产锁需要哪些性质

性质含义ZooKeeper节点机制是否单独足够
Mutual Exclusion同一时刻只有一个合法持有者在Session和节点视角可以协调
Deadlock Recovery持有者崩溃后最终释放临时节点在Session过期后删除
Ordering等待者按稳定顺序竞争顺序节点提供近似FIFO排队
Ownership Check只有持有者能释放自己的锁Recipe和节点路径需校验
Fencing旧持有者不能写外部资源必须由外部资源校验递增Token
Idempotency重复执行不重复产生业务效果需要业务键、唯一约束和状态机

锁解决的是进入临界区的协调,不会自动回滚数据库,不会让网络调用只执行一次,也不会让任务在Leader切换时Exactly Once。

三、锁目录的数据模型

对资源 inventory/SKU-1,可以使用稳定锁目录:

text
/locks/inventory/SKU-1
  lock-0000000041
  lock-0000000042
  lock-0000000043

每个竞争者创建 EPHEMERAL_SEQUENTIAL 节点。父路径必须稳定存在且受 ACL 保护,不能让普通客户端删除父路径、伪造节点或重置序列语义。

节点内容可包含诊断信息:

  • 应用和实例ID。
  • 线程或任务ID。
  • 请求ID。
  • 创建时间仅用于诊断。
  • 不应存密码、Token等秘密。

节点内容不能作为可信锁顺序,顺序由ZooKeeper生成的节点序号决定。

四、获取锁完整算法

mermaid
flowchart TD
    A["创建临时顺序节点myNode"] --> B["读取锁目录所有子节点"]
    B --> C["按顺序后缀排序"]
    C --> D{"myNode是否最小"}
    D -- "是" --> E["取得锁和Fencing Token"]
    D -- "否" --> F["找到紧邻前驱predecessor"]
    F --> G["用exists同时检查并注册删除Watch"]
    G --> H{"前驱是否仍存在"}
    H -- "否" --> B
    H -- "是" --> I["等待前驱删除、超时或Session事件"]
    I --> B
    E --> J["外部资源校验Token后执行业务"]

4.1 为什么是临时节点

客户端正常关闭或 Session 最终过期后,服务端删除临时节点,后继竞争者才能继续。短暂 TCP 断开不会立即删除,因此不会因为一次网络抖动立刻把锁交给其他人。

代价是:断连期间客户端不知道 Session 是否仍有效,应暂停高风险写;过期后必须永久放弃旧所有权。

4.2 为什么是顺序节点

顺序后缀为竞争者建立稳定队列,避免所有客户端同时抢同一个固定 znode 形成高冲突。最小节点持锁,其他节点按前驱关系等待。

它提供的是近似 FIFO:客户端调度、网络延迟、Session过期和创建请求到达 ZooKeeper 的顺序会影响最终编号,不能承诺严格按业务发起时间公平。

五、为什么只监听前驱节点

若所有等待者都 Watch 当前最小节点,锁释放时 N 个客户端同时被唤醒、读取子节点并竞争,只有一个成功,其余再次睡眠,形成羊群效应。

前驱 Watch 建立链:

text
41 持有锁
42 只Watch 41
43 只Watch 42
44 只Watch 43

删除 41 时主要唤醒 42;42 删除时再唤醒 43。每次所有权转移只需要相邻等待者重新判断。

六、前驱删除与注册Watch的竞态怎样处理

错误步骤:

text
先读取子节点,发现前驱是41
41立刻删除
再对41注册Watch失败
客户端永远等待一个已不存在节点

正确方式使用 exists(predecessor, watcher),让“检查存在”和“注册该存在状态的 Watch”由一个ZooKeeper操作表达:

  1. 若返回 Stat,说明前驱在注册时仍存在,可以等待。
  2. 若返回 null,说明前驱已经消失,立即重新读取并排序。
  3. Watch触发后也重新读取完整子节点,不直接假定自己一定获得锁。

每次唤醒都重新判断“自己是否最小”,这是应对重复通知、Session变化和并发删除的关键。

七、获取超时和取消

等待者达到超时、线程中断或业务请求取消时,应删除自己的临时顺序节点,并停止等待。否则虽然 Session 最终会清理节点,但存活客户端会留下无用排队节点,增加延迟和监控噪声。

取消与获得锁可能并发:

  • 超时线程准备删除节点时,前驱恰好删除。
  • Watch回调和取消逻辑同时唤醒。
  • 删除返回ConnectionLoss,结果未知。

成熟 Recipe 应通过状态机保证只进入 ACQUIREDCANCELLED 之一,并在未知结果时查询节点事实,而不是无条件重建一个新节点。

八、释放锁全过程

mermaid
flowchart TD
    A["业务临界区结束"] --> B["停止提交新的异步子任务"]
    B --> C["等待必须属于本次锁的操作收敛"]
    C --> D["删除自己的锁节点"]
    D --> E{"删除结果是否明确"}
    E -- "成功或NoNode" --> F["本地标记已释放"]
    E -- "ConnectionLoss" --> G["重连后查询节点和Session事实"]
    G --> F
    F --> H["后继收到前驱删除提示"]

不能通过删除“当前最小节点”释放锁;只能删除自己创建并拥有的节点。否则一个超时线程可能误删新持有者节点。

临界区必须使用 finally 释放,但若 Session 已过期,本地 release 可能得到节点不存在。这通常表示 ZooKeeper 已经撤销所有权,业务应记录失锁,而不是把异常解释成“锁仍在”。

九、断连、过期和锁所有权

事件ZooKeeper节点可能状态业务动作
Disconnected节点可能仍存在暂停高风险写,等待恢复或超时
原Session重连节点通常仍属于原Session重新读取队列确认仍最小
Session Expired临时节点最终删除永久撤销旧所有权,重新竞争
新Session建立旧节点不属于新Session创建新顺序节点,获得新Token

连接恢复回调不能直接执行 isOwner=true。必须确认:Session未过期、自己的节点仍存在、仍是最小节点、外部资源仍接受Token。

十、GC Pause造成双执行窗口

mermaid
flowchart TD
    A["Client A取得锁Token 101"] --> B["A发生长时间GC Pause"]
    B --> C["A的Session过期,临时节点删除"]
    C --> D["Client B取得锁Token 102"]
    D --> E["B向外部资源写入Token 102"]
    B --> F["A恢复并继续旧代码"]
    F --> G{"外部资源是否校验Token"}
    G -- "否" --> H["A可能覆盖B,出现双写"]
    G -- "是" --> I["拒绝Token 101的陈旧写入"]

这不是 ZooKeeper 共识失效。ZooKeeper 已经正确把锁交给 B;问题是外部数据库无法知道 A 的业务线程已经过期。

十一、Fencing Token原理

Fencing Token 是每次所有权授予时递增的序号。外部资源保存已接受的最大 Token,只接受不小于当前规则的新持有者请求,拒绝旧持有者。

sql
UPDATE resource_state
SET value = :new_value,
    fencing_token = :token
WHERE resource_id = :resource_id
  AND fencing_token < :token;

影响行数为 1 才表示写入被接受;为 0 时需要读取当前 Token,判断是陈旧持有者还是资源不存在。

11.1 Token从哪里来

可选来源:

  • 同一稳定锁父路径下的顺序节点序号,需理解序号范围、父节点生命周期和解析规则。
  • 数据库序列或原子递增版本。
  • 独立一致性服务产生的单调序号。

不要使用客户端本地时间、随机UUID或进程内 AtomicLong 作为跨进程Fencing Token:它们不能证明全局单调。

11.2 Token必须由外部资源校验

仅在客户端写:

java
if (myToken < latestToken) return;

没有意义,因为暂停的旧客户端不知道最新Token。数据库、对象存储网关、任务执行器或设备网关必须在接收写入时原子比较并保存Token。

11.3 Fencing不能替代幂等

同一个合法Token的请求可能因网络重试重复到达。仍要使用业务请求ID、唯一约束和结果复用防止同一持有者重复执行。

十二、Curator InterProcessMutex Demo

java
import org.apache.curator.framework.CuratorFramework;
import org.apache.curator.framework.CuratorFrameworkFactory;
import org.apache.curator.framework.recipes.locks.InterProcessMutex;
import org.apache.curator.retry.ExponentialBackoffRetry;

import java.util.concurrent.TimeUnit;

public class CuratorMutexDemo {
    public static void main(String[] args) throws Exception {
        CuratorFramework client = CuratorFrameworkFactory.newClient(
                "127.0.0.1:2181",
                10_000,
                5_000,
                new ExponentialBackoffRetry(500, 5));
        client.start();

        if (!client.blockUntilConnected(8, TimeUnit.SECONDS)) {
            throw new IllegalStateException("ZooKeeper connect timeout");
        }

        InterProcessMutex mutex = new InterProcessMutex(
                client, "/locks/cache/product-1001");
        boolean acquired = mutex.acquire(3, TimeUnit.SECONDS);
        if (!acquired) {
            System.out.println("lock busy, return or enqueue retry");
            client.close();
            return;
        }

        try {
            rebuildCacheIdempotently("product-1001");
        } finally {
            mutex.release();
            client.close();
        }
    }

    static void rebuildCacheIdempotently(String key) {
        System.out.println("rebuild=" + key);
    }
}

InterProcessMutex 具有可重入语义,调用方必须成对 release;释放线程和Recipe约束要按当前Curator版本遵守。需要非重入语义时可评估其他Recipe,不能仅凭类名猜行为。

Curator帮助处理节点、Watch和连接细节,但不会自动提供业务Fencing Token、数据库幂等和事务一致性。

十三、ZooKeeper锁与Redis锁准确比较

维度ZooKeeper临时顺序锁Redis单实例/主从锁
所有权期限SessionKey TTL
等待通知Watch前驱常见轮询、Pub/Sub或客户端队列
排队顺序节点近似FIFO常见实现默认非公平,Redisson有不同Recipe
故障释放Session过期TTL到期或显式Lua释放
主要风险GC/断连导致Session过期,旧线程恢复TTL到期、续期失败、故障转移、误删
吞吐协调写成本较高通常更高
Fencing仍然需要同样需要

“ZooKeeper偏CP,所以锁绝对安全”是错误结论。共识保护的是ZooKeeper节点状态,不会自动阻止暂停线程对ZooKeeper之外的资源继续写。

十四、应用Leader选举与ZooKeeper集群选主不是一回事

选举选谁使用机制
ZooKeeper Server Leader负责ZAB事务排序的ServerFastLeaderElection、epoch、zxid、Quorum
业务应用Leader多个业务实例中的协调者临时顺序节点、LeaderLatch/LeaderSelector

业务 LeaderLatch 不参与ZooKeeper集群ZAB选主。两个概念都叫Leader,但职责和协议完全不同。

十五、业务选主完整过程

mermaid
flowchart TD
    A["每个实例创建临时顺序候选节点"] --> B["按序号排序"]
    B --> C["最小节点成为业务Leader"]
    C --> D["取得新的Fencing Token"]
    D --> E["启动受Token保护的主任务"]
    E --> F{"Session断连、过期或主动让位"}
    F -- "断连" --> G["暂停高风险任务"]
    F -- "过期或让位" --> H["停止并清理本地任务"]
    H --> I["下一候选重新确认并成为Leader"]

Leader回调到达和旧任务真正停止之间可能有窗口。任务执行器必须可取消、可幂等,并让外部写操作携带Token。

十六、Curator LeaderLatch Demo

java
import org.apache.curator.framework.CuratorFramework;
import org.apache.curator.framework.recipes.leader.LeaderLatch;
import org.apache.curator.framework.recipes.leader.LeaderLatchListener;

import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.Future;
import java.util.concurrent.atomic.AtomicReference;

public final class SafeLeaderWorker implements AutoCloseable {
    private final LeaderLatch latch;
    private final ExecutorService executor = Executors.newSingleThreadExecutor();
    private final AtomicReference<Future<?>> running = new AtomicReference<Future<?>>();

    public SafeLeaderWorker(CuratorFramework client, String instanceId) {
        this.latch = new LeaderLatch(client, "/leaders/order-reconcile", instanceId);
        this.latch.addListener(new LeaderLatchListener() {
            @Override
            public void isLeader() {
                Future<?> future = executor.submit(() -> runIdempotentLoop());
                if (!running.compareAndSet(null, future)) future.cancel(true);
            }

            @Override
            public void notLeader() {
                Future<?> future = running.getAndSet(null);
                if (future != null) future.cancel(true);
            }
        });
    }

    public void start() throws Exception {
        latch.start();
    }

    private void runIdempotentLoop() {
        while (!Thread.currentThread().isInterrupted()) {
            // 每批任务仍需幂等、状态机和Fencing Token
        }
    }

    @Override
    public void close() throws Exception {
        Future<?> future = running.getAndSet(null);
        if (future != null) future.cancel(true);
        latch.close();
        executor.shutdownNow();
    }
}

这段代码展示了Leader身份变化与任务生命周期绑定,但生产实现还必须处理任务阻塞不响应中断、Token获取、任务事实表和进程强杀。

十七、Leader不等于Exactly Once

场景:旧 Leader 已生成账单但在写“已完成”前Session过期;新 Leader不知道业务是否完成,再次生成。

解决方案:

  • billing_date + tenant_id 建唯一任务键。
  • 任务状态机记录 PENDING/RUNNING/SUCCEEDED/FAILED
  • 执行步骤和结果可查询。
  • 写外部资源携带Fencing Token。
  • 新 Leader接管前扫描事实状态。
  • 定期对账,不以Leader回调日志作为成功事实。

十八、配置通知与锁的共同底层

配置通知同样依赖 znode、version、Watch和Session,但不需要创建锁节点:

mermaid
flowchart TD
    A["管理员用expectedVersion更新配置"] --> B["ZooKeeper提交新version"]
    B --> C["Watcher或CuratorCache收到变化提示"]
    C --> D["应用重新读取完整数据和Stat"]
    D --> E["校验后原子替换本地配置"]

详细Watcher语义见Watcher与Session全过程。不要为了“确保只刷新一次”给配置读取再套分布式锁;版本裁决和幂等应用通常更合适。

十九、JDK 8 Demo:外部资源拒绝旧Token

java
public class FencingTokenDemo {

    static final class ExternalResource {
        private long maxToken;
        private String value;

        synchronized boolean write(long token, String newValue) {
            if (token <= maxToken) {
                System.out.println("reject token=" + token + ", max=" + maxToken);
                return false;
            }
            maxToken = token;
            value = newValue;
            System.out.println("accept token=" + token + ", value=" + value);
            return true;
        }
    }

    public static void main(String[] args) {
        ExternalResource resource = new ExternalResource();
        resource.write(101L, "written-by-A");
        resource.write(102L, "written-by-B");
        resource.write(101L, "late-write-by-A");
        System.out.println("final=" + resource.value);
    }
}

预期输出:

text
accept token=101, value=written-by-A
accept token=102, value=written-by-B
reject token=101, max=102
final=written-by-B

示例使用 token <= maxToken 拒绝同Token重复写;如果业务允许同Token幂等重放,应再比较请求ID和结果摘要,而不是简单放行所有相同Token请求。

二十、商业场景选型

场景推荐方式原因
低频主备调度ZooKeeper选主 + 任务幂等 + Token协调清晰,可恢复接管
多实例缓存重建Redis或ZooKeeper锁 + 双重检查按频率和已有基础设施选择
订单扣库存数据库条件更新/唯一约束优先高频业务不应依赖全局粗锁串行化
设备主控ZooKeeper选主 + 设备网关校验Token必须阻止旧主恢复后发命令
分片分配Leader生成版本化分配计划Worker按epoch/version接受新计划
批量对账Leader负责扫描,任务表原子认领Leader切换不重复处理同一批次

二十一、生产排查Runbook

21.1 锁长期不释放

  1. 查看锁目录所有节点、序号和 ephemeralOwner
  2. 确认持有者Session是否仍有效,不要直接删最小节点。
  3. 检查持有者线程栈、GC Pause、外部调用和死循环。
  4. 检查Recipe是否可重入但release次数不足。
  5. 只有确认业务影响和所有权后才做人工处理,并保留审计。

21.2 大量等待者同时被唤醒

检查是否所有客户端都Watch最小节点、父节点children或统一配置节点;改为只Watch紧邻前驱。检查回调是否又触发全量昂贵操作。

21.3 两个实例都认为自己是Leader

  1. 对比它们的Session状态和Leader节点路径。
  2. 查看旧实例是否在Disconnected/Expired后仍继续任务。
  3. 检查任务是否响应取消和中断。
  4. 查看外部资源记录的最大Fencing Token。
  5. 若没有Token,只能依赖业务事实对账并补救双执行。

21.4 获取锁延迟高

查看等待节点数、临界区P95/P99、Session事件、ZooKeeper写延迟和回调执行器。锁持有时间长时,增加ZooKeeper节点不会提高临界区吞吐;应缩短临界区或按资源分片锁路径。

二十二、必须监控的指标

  • 每个锁路径等待节点数。
  • 获取等待时间、持有时间和超时率。
  • Session Disconnected/Expired次数。
  • 节点创建、删除和ConnectionLoss次数。
  • 前驱Watch触发与重新检查次数。
  • 当前Leader实例、任期/Token和切换次数。
  • 任务取消耗时和旧任务仍运行数量。
  • 外部资源拒绝旧Token次数。
  • 幂等命中、重复任务和人工接管次数。

二十三、常见错误及后果

错误后果
所有人Watch最小节点羊群效应和ZK读峰值
前驱不存在仍进入等待永久卡住
Session过期后继续临界区与新持有者双写
Token只在客户端检查暂停客户端不知道已有新Token
UUID作为Fencing Token无法比较新旧所有权
Leader回调等于任务成功切换窗口重复执行或漏执行
用粗锁串行所有订单吞吐下降、锁路径热点
人工删除最小节点解锁可能撤销合法持有者并制造双执行

二十四、面试标准回答

ZooKeeper排队锁通常在稳定父路径下创建临时顺序节点,序号最小者持锁,其他客户端只对紧邻前驱执行exists + Watch,前驱已经删除就立即重新排序,仍存在才等待,因此避免所有等待者监听最小节点产生羊群效应。临时节点只在Session过期后删除,断连期间所有权不确定;长GC可能让旧Session过期、新客户端取得锁,而旧线程恢复后继续写,所以外部数据库或网关必须原子校验单调递增Fencing Token,不能只靠客户端判断。业务Leader选举可复用顺序节点或Curator LeaderLatch,但Leader身份不保证任务Exactly Once,任务还需要幂等键、状态机、接管扫描和对账。

二十五、关联知识点

本章小结

ZooKeeper锁的完整算法不是“创建一个临时节点”:竞争者创建临时顺序节点、排序、只Watch前驱、处理前驱已删除竞态,并在Session事件和取消中维护状态机。它只能协调ZooKeeper中的合法所有权;要防止暂停旧线程写外部资源,必须使用外部强制校验的Fencing Token。选主同理,它只选出协调者,不会自动让业务任务Exactly Once。