ZooKeeper分布式锁与选主全过程:前驱Watch、Session失效和Fencing Token
ZooKeeper 临时顺序节点能构造排队锁和 Leader 选举,但“临时节点会自动删除”不等于外部资源绝不并发写。客户端可能在长时间 GC 后恢复,旧 Session 已经失效,新客户端已经获得锁;如果数据库仍接受旧客户端写入,就会出现两个实际执行者。
因此完整正确性链必须包含:顺序排队、前驱 Watch、Session 状态、安全退出、业务幂等,以及由外部资源校验的 Fencing Token。
一、学习目标
- 从临时顺序节点推导分布式互斥过程。
- 解释为什么只监听前驱节点可以避免羊群效应。
- 处理前驱节点在 Watch 注册前已删除的竞态。
- 区分锁节点删除、客户端线程恢复和外部资源所有权。
- 解释 Fencing Token 必须在哪里校验。
- 区分公平排队、严格公平和高吞吐锁。
- 区分 Leader 身份、任务执行和Exactly Once。
- 使用 Curator Recipe并理解它没有替业务解决什么。
二、一个生产锁需要哪些性质
| 性质 | 含义 | ZooKeeper节点机制是否单独足够 |
|---|---|---|
| Mutual Exclusion | 同一时刻只有一个合法持有者 | 在Session和节点视角可以协调 |
| Deadlock Recovery | 持有者崩溃后最终释放 | 临时节点在Session过期后删除 |
| Ordering | 等待者按稳定顺序竞争 | 顺序节点提供近似FIFO排队 |
| Ownership Check | 只有持有者能释放自己的锁 | Recipe和节点路径需校验 |
| Fencing | 旧持有者不能写外部资源 | 必须由外部资源校验递增Token |
| Idempotency | 重复执行不重复产生业务效果 | 需要业务键、唯一约束和状态机 |
锁解决的是进入临界区的协调,不会自动回滚数据库,不会让网络调用只执行一次,也不会让任务在Leader切换时Exactly Once。
三、锁目录的数据模型
对资源 inventory/SKU-1,可以使用稳定锁目录:
/locks/inventory/SKU-1
lock-0000000041
lock-0000000042
lock-0000000043每个竞争者创建 EPHEMERAL_SEQUENTIAL 节点。父路径必须稳定存在且受 ACL 保护,不能让普通客户端删除父路径、伪造节点或重置序列语义。
节点内容可包含诊断信息:
- 应用和实例ID。
- 线程或任务ID。
- 请求ID。
- 创建时间仅用于诊断。
- 不应存密码、Token等秘密。
节点内容不能作为可信锁顺序,顺序由ZooKeeper生成的节点序号决定。
四、获取锁完整算法
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 建立链:
41 持有锁
42 只Watch 41
43 只Watch 42
44 只Watch 43删除 41 时主要唤醒 42;42 删除时再唤醒 43。每次所有权转移只需要相邻等待者重新判断。
六、前驱删除与注册Watch的竞态怎样处理
错误步骤:
先读取子节点,发现前驱是41
41立刻删除
再对41注册Watch失败
客户端永远等待一个已不存在节点正确方式使用 exists(predecessor, watcher),让“检查存在”和“注册该存在状态的 Watch”由一个ZooKeeper操作表达:
- 若返回 Stat,说明前驱在注册时仍存在,可以等待。
- 若返回
null,说明前驱已经消失,立即重新读取并排序。 - Watch触发后也重新读取完整子节点,不直接假定自己一定获得锁。
每次唤醒都重新判断“自己是否最小”,这是应对重复通知、Session变化和并发删除的关键。
七、获取超时和取消
等待者达到超时、线程中断或业务请求取消时,应删除自己的临时顺序节点,并停止等待。否则虽然 Session 最终会清理节点,但存活客户端会留下无用排队节点,增加延迟和监控噪声。
取消与获得锁可能并发:
- 超时线程准备删除节点时,前驱恰好删除。
- Watch回调和取消逻辑同时唤醒。
- 删除返回ConnectionLoss,结果未知。
成熟 Recipe 应通过状态机保证只进入 ACQUIRED 或 CANCELLED 之一,并在未知结果时查询节点事实,而不是无条件重建一个新节点。
八、释放锁全过程
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造成双执行窗口
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,只接受不小于当前规则的新持有者请求,拒绝旧持有者。
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必须由外部资源校验
仅在客户端写:
if (myToken < latestToken) return;没有意义,因为暂停的旧客户端不知道最新Token。数据库、对象存储网关、任务执行器或设备网关必须在接收写入时原子比较并保存Token。
11.3 Fencing不能替代幂等
同一个合法Token的请求可能因网络重试重复到达。仍要使用业务请求ID、唯一约束和结果复用防止同一持有者重复执行。
十二、Curator InterProcessMutex Demo
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单实例/主从锁 |
|---|---|---|
| 所有权期限 | Session | Key TTL |
| 等待通知 | Watch前驱 | 常见轮询、Pub/Sub或客户端队列 |
| 排队 | 顺序节点近似FIFO | 常见实现默认非公平,Redisson有不同Recipe |
| 故障释放 | Session过期 | TTL到期或显式Lua释放 |
| 主要风险 | GC/断连导致Session过期,旧线程恢复 | TTL到期、续期失败、故障转移、误删 |
| 吞吐 | 协调写成本较高 | 通常更高 |
| Fencing | 仍然需要 | 同样需要 |
“ZooKeeper偏CP,所以锁绝对安全”是错误结论。共识保护的是ZooKeeper节点状态,不会自动阻止暂停线程对ZooKeeper之外的资源继续写。
十四、应用Leader选举与ZooKeeper集群选主不是一回事
| 选举 | 选谁 | 使用机制 |
|---|---|---|
| ZooKeeper Server Leader | 负责ZAB事务排序的Server | FastLeaderElection、epoch、zxid、Quorum |
| 业务应用Leader | 多个业务实例中的协调者 | 临时顺序节点、LeaderLatch/LeaderSelector |
业务 LeaderLatch 不参与ZooKeeper集群ZAB选主。两个概念都叫Leader,但职责和协议完全不同。
十五、业务选主完整过程
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
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,但不需要创建锁节点:
flowchart TD
A["管理员用expectedVersion更新配置"] --> B["ZooKeeper提交新version"]
B --> C["Watcher或CuratorCache收到变化提示"]
C --> D["应用重新读取完整数据和Stat"]
D --> E["校验后原子替换本地配置"]详细Watcher语义见Watcher与Session全过程。不要为了“确保只刷新一次”给配置读取再套分布式锁;版本裁决和幂等应用通常更合适。
十九、JDK 8 Demo:外部资源拒绝旧Token
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);
}
}预期输出:
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 锁长期不释放
- 查看锁目录所有节点、序号和
ephemeralOwner。 - 确认持有者Session是否仍有效,不要直接删最小节点。
- 检查持有者线程栈、GC Pause、外部调用和死循环。
- 检查Recipe是否可重入但release次数不足。
- 只有确认业务影响和所有权后才做人工处理,并保留审计。
21.2 大量等待者同时被唤醒
检查是否所有客户端都Watch最小节点、父节点children或统一配置节点;改为只Watch紧邻前驱。检查回调是否又触发全量昂贵操作。
21.3 两个实例都认为自己是Leader
- 对比它们的Session状态和Leader节点路径。
- 查看旧实例是否在Disconnected/Expired后仍继续任务。
- 检查任务是否响应取消和中断。
- 查看外部资源记录的最大Fencing Token。
- 若没有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。
