Skip to content

ZooKeeper Watcher与Session全过程:重连、过期、事件窗口和生产恢复

ZooKeeper 的 Watcher 和 Session 经常被概括成“节点变化就通知客户端、客户端挂了临时节点自动删除”。这句话方向正确,但不足以指导生产:断开连接不等于 Session 过期;Watch 不是永久事件日志;客户端恢复运行不代表还拥有旧锁;节点在断连期间创建又删除,客户端可能只知道“状态需要重读”,而不是收到两条可靠事件。

本页从协议状态和失败窗口出发,讲清一次连接如何建立、心跳怎样维持 Session、重连怎样恢复、过期为什么不可逆、Watch 如何注册和触发、怎样构造不会因事件缺口而永久错误的客户端。

一、学习目标

  1. 区分 TCP 连接、ZooKeeper Session 和业务锁所有权。
  2. 解释 Session ID、Session Password、协商超时和心跳的作用。
  3. 区分 DisconnectedSyncConnectedExpiredAuthFailed
  4. 解释临时节点为什么短暂断线不删除、Session 过期才删除。
  5. 区分标准一次性 Watch、Persistent Watch 和 Persistent Recursive Watch。
  6. 解释 existsgetDatagetChildren 各自注册什么事件。
  7. 处理“读取数据”和“重新注册 Watch”之间的变化窗口。
  8. 设计断连后的安全停写、快照重建和生产排查流程。

二、先区分三个生命周期

生命周期标识失效后果
TCP连接到某台ZooKeeper Server的Socket可重连其他Server,Session可能仍有效
ZooKeeper SessionsessionId和认证凭据过期后不可恢复,临时节点和标准Watch所有权失效
业务所有权锁、Leader身份、任务执行权即使客户端线程恢复,也必须重新证明所有权
mermaid
flowchart TD
    A["客户端进程"] --> B["连接ZooKeeper Server A"]
    B --> C["建立或恢复ZooKeeper Session"]
    C --> D["Session拥有临时节点和Watch"]
    D --> E["业务把临时节点解释为锁或Leader身份"]

网络断开只直接破坏 TCP 连接。Session 是否过期由服务端在协商超时窗口内是否重新收到有效客户端活动决定。业务所有权还可能受 Fencing Token、外部数据库和上层状态机约束。

三、Session建立全过程

mermaid
flowchart TD
    A["客户端选择集群地址"] --> B["建立TCP连接"]
    B --> C["发送ConnectRequest和期望Session Timeout"]
    C --> D["Server按允许范围协商Timeout"]
    D --> E["新会话生成sessionId与会话凭据"]
    E --> F["返回ConnectResponse"]
    F --> G["客户端进入SyncConnected"]

3.1 Session ID和Session Password做什么

  • sessionId 标识会话。
  • 会话凭据用于防止任意客户端仅凭 sessionId 接管会话。
  • 重连时客户端携带原 Session 信息和已见事务位置,与另一台 Server 尝试恢复同一会话。
  • Session 过期后,旧 ID 不能重新变回有效会话;客户端只能建立新 Session。

不要把 session password 输出到日志、监控标签或接口响应,它属于敏感认证信息。

3.2 Session Timeout为什么是协商值

客户端提交期望超时,Server 会根据集群允许的最小、最大 Session Timeout 进行约束。ZooKeeper 常见默认边界与 tickTime 相关,例如最小约为两倍 tickTime、最大约为二十倍,但管理员可以显式配置,必须以当前集群配置和握手后的实际值为准。

客户端配置 60 秒不代表服务端一定接受 60 秒;排查时应查看协商结果。

四、心跳与会话活性

客户端库会在没有正常请求时发送 Ping 等活动,防止 Session 因完全静默而过期。心跳不是业务线程手工每秒调用一次 exists(),而是客户端 I/O 线程和协议实现的职责。

mermaid
flowchart TD
    A["客户端I/O线程维护Session"] --> B{"近期是否有正常请求"}
    B -- "有" --> C["正常请求也体现会话活动"]
    B -- "没有" --> D["按客户端节奏发送Ping"]
    C --> E["Server更新会话活性"]
    D --> E
    E --> F{"超过协商Timeout仍无有效活动"}
    F -- "否" --> A
    F -- "是" --> G["服务端判定Session过期"]

如果 JVM Stop-The-World、CPU 长时间饥饿、容器冻结或网络中断使 I/O 线程无法活动,即使业务线程没有崩溃,Session 也可能在服务端过期。

五、连接状态必须按“是否终态”理解

状态含义是否可继续认为拥有锁
SyncConnected已连接可服务Server,会话当前有效仍需检查业务锁节点和Token
Disconnected暂时没有连接,尚不能确定Session是否过期高风险写操作应暂停或进入安全模式
ConnectedReadOnly连接到允许只读服务的Server,具体需集群启用不能执行写入,业务语义需单独设计
Expired服务端已判定旧Session过期绝对不能继续使用旧锁或Leader身份
AuthFailed认证失败不能继续访问受保护数据,应终止或人工修复
Closed客户端已关闭旧客户端不再工作

实际枚举名称和可见状态会因原生客户端、Curator和ZooKeeper版本略有不同。设计重点是:Disconnected 是不确定状态,Expired 是不可逆终态。

六、Disconnected与Expired的完整时间线

mermaid
flowchart TD
    A["连接正常且Session有效"] --> B["网络断开或Server重启"]
    B --> C["客户端收到Disconnected"]
    C --> D{"在Session Timeout内重连成功"}
    D -- "是" --> E["恢复原Session并进入SyncConnected"]
    E --> F["临时节点仍属于原Session"]
    D -- "否" --> G["服务端判定Session过期"]
    G --> H["删除该Session临时节点"]
    H --> I["触发相关Watch变化"]
    I --> J["客户端最终收到Expired并建立新Session"]

6.1 断连期间为什么是状态未知

客户端看不到集群当前状态:

  • Session 可能仍有效。
  • Session 可能刚刚过期。
  • 锁节点可能仍存在。
  • 锁节点可能已经被删除并由其他客户端接管。

因此断连客户端不能说“我进程还活着,所以锁仍属于我”。对不可并发写入的数据库、对象存储、设备控制等资源,断连时应暂停高风险副作用,或让外部资源使用 Fencing Token 拒绝旧持有者。

七、Session过期后到底发生什么

服务端完成过期处理后:

  1. 旧 Session 不能恢复。
  2. 该 Session 创建的临时节点被删除。
  3. 其他客户端对这些节点设置的 Watch 可能被触发。
  4. 旧客户端的 Watch 和认证上下文不能继续按原会话使用。
  5. Curator 等高层库通常会建立新 Session,并重建它管理的Recipe状态。
  6. 业务必须重新竞争锁、重新参加选主、重新注册临时节点。

新建 Session 后节点名称可能变化,尤其是临时顺序节点。业务不能缓存旧节点路径并继续宣称所有权。

八、GC Pause为什么会让活进程丢锁

假设 Session Timeout 为 15 秒:

text
T0        客户端持有锁,Token=101
T1        JVM发生20秒Stop-The-World
T1+15s    ZooKeeper判定Session过期并删除锁节点
T1+16s    其他客户端获得锁,Token=102
T1+20s    旧JVM恢复,业务线程继续执行

如果外部数据库不校验 Token,Token=101 的旧客户端可能与 Token=102 的新客户端同时写入。这就是“临时节点自动删除”仍不能完全保护外部资源的原因。

相关原理继续看ZooKeeper锁、选主与Fencing Token

九、Watch是什么,不是什么

Watch 是“当被观察的 znode 状态发生特定变化时,向客户端发送变化提示”的机制。

它不是:

  • Kafka 式持久事件日志。
  • 每个变化永久保留的消息队列。
  • 业务数据传输通道。
  • 自动递归监听整棵树的传统默认机制。
  • 保证业务回调只执行一次的幂等系统。

正确思想:

Watch只负责提醒“状态可能变了”;客户端收到提示后重新读取当前事实,按版本或快照对账。

十、不同读取API注册什么Watch

API观察对象常见触发事件
exists(path, watch)节点是否存在和Stat创建、删除、数据变化
getData(path, watch, stat)节点数据数据变化、节点删除
getChildren(path, watch)直接子节点集合子节点新增或删除、父节点删除

标准 child watch 通常只观察直接子节点名称变化,不自动观察每个子节点的数据变化,也不递归观察孙节点。

为什么exists很特殊

对当前不存在的节点调用 exists 可以注册“节点未来被创建”的 Watch。getData 对不存在节点会失败,不能以同样方式观察未来创建。

十一、标准Watch为什么是一次性触发

标准 Watch 被对应事件触发后,需要重新读取并重新注册。一次性设计降低服务端长期维护复杂度和重复通知压力,但把状态收敛责任交给客户端。

错误写法的思想:

text
启动时注册一次Watch

以后永远依赖同一个Watch收到所有变化

第一次触发后若不重新注册,后续变化不会再通知。

十二、读取与重新注册之间有没有竞态

使用 getData(path, watcher, stat) 时,“读取当前数据”和“注册下一次 Watch”由一个 ZooKeeper 请求完成,可以避免应用自己先 getData(false)、再单独调用 exists(true) 造成明显空窗。

但客户端仍不能把 Watch 当事件流:

  • Watch 触发只告诉你发生过相关变化。
  • 多次快速变化可能在重新读取时已经收敛成最新状态。
  • 节点可能在断连期间创建后又删除。
  • 回调处理期间还可能发生下一次变化。

安全循环是:

mermaid
flowchart TD
    A["读取当前数据并同时注册Watch"] --> B["把数据版本应用到本地快照"]
    B --> C["等待Watch变化提示"]
    C --> D["把回调当作失效信号"]
    D --> A

客户端关心的是最终当前状态,不是执着于逐条重放中间瞬态。

十三、Watch顺序保证应怎样理解

ZooKeeper 对同一客户端的操作、异步响应和 Watch 事件提供有序语义,使客户端不会在完全任意顺序中观察同一事务变化。但这个保证不等于:

  • 所有客户端在同一物理时刻收到事件。
  • 回调业务执行完成顺序与事务顺序一致。
  • 网络断开后仍能获得每个中间事件。
  • Watch 回调可以执行长业务而不影响客户端事件线程。

回调应快速记录失效并把重读任务交给受控执行器;不能在 ZooKeeper 事件线程中调用慢 HTTP、执行大 SQL 或无限重试。

十四、Persistent Watch和Persistent Recursive Watch

ZooKeeper 3.6 引入持久 Watch 能力,常见模式包括:

  • Persistent Watch:持续观察目标节点相关变化,不因一次触发自动移除。
  • Persistent Recursive Watch:持续观察目标路径及其后代变化。

版本边界必须明确:

  • 老 ZooKeeper Server 不支持时不能使用。
  • 客户端库也必须支持相应 API。
  • 持久 Watch 仍不是持久事件日志;客户端长时间离线后仍应重建快照。
  • 递归 Watch 可能产生大量事件,必须控制路径规模和回调成本。
  • 不同事件类型和语义应按锁定版本官方文档验证。

不能把 Curator Cache 与 ZooKeeper Persistent Watch 直接画等号。Curator Cache 是客户端更高层的本地缓存 Recipe,底层可能根据版本使用不同监听和重建方式。

十五、为什么生产更常使用Curator Cache

直接手写 Watch 容易遗漏:

  • 初始数据加载。
  • Watch 重新注册。
  • Session 重建。
  • 子节点新增后的数据监听。
  • 本地快照线程安全。
  • 回调和关闭生命周期。

Curator 提供多种 Cache/Recipe。新项目应根据当前 Curator 版本选择推荐 API,旧版 PathChildrenCacheNodeCacheTreeCache 与较新的 CuratorCache 有版本边界。

概念示例:

java
CuratorCache cache = CuratorCache.build(client, "/config/order-service");
cache.listenable().addListener((type, oldData, data) -> {
    // 回调只投递“需要重载”的轻量任务,不执行慢业务
    reloadExecutor.execute(() -> reloadFromZooKeeper());
});
cache.start();

Cache 提供的是最终收敛的客户端视图,不是数据库事务隔离。业务仍应使用 znode version、配置 revision 或内容摘要防止旧回调覆盖新值。

十六、重连和新Session恢复动作不同

场景会话正确动作
短暂断连后恢复原Session仍有效重读关键状态,确认缓存收敛
Session Expired旧Session永久失效建立新Session、重建临时节点、重新竞争锁和Leader
AuthFailed无法访问受保护节点停止业务并修复认证,不应无限重试
集群只读连接只能读禁止依赖写入完成的业务流程

高层客户端可能自动重连,但“连接恢复”不能自动证明业务所有权恢复。业务需要监听连接状态并管理安全门。

十七、业务安全门设计

核心写任务可以维护一个本地门状态:

text
CONNECTED + 当前节点仍是最小顺序节点 + Token仍被外部资源接受
    => ALLOW_WRITE

DISCONNECTED / EXPIRED / 未重新竞争
    => DENY_WRITE

Disconnected 是否立即停写取决于业务风险:

  • 只读缓存计算可继续一小段受控时间。
  • 扣款、设备控制、主库写入等应偏保守停写。
  • 有可靠 Fencing Token 的外部资源可以最终拒绝旧请求,但客户端自身仍应尽快停止。

十八、JDK 8 Demo:Session状态驱动的写入安全门

下面的无依赖 Demo 模拟连接状态和锁 Token。重点不是实现 ZooKeeper 客户端,而是证明:Session 过期后,即使网络重新连接,也不能恢复旧 Token 的写权限。

java
public class SessionSafetyGateDemo {

    enum State {
        CONNECTED, DISCONNECTED, EXPIRED
    }

    static final class Gate {
        private State state = State.DISCONNECTED;
        private long ownedToken = -1L;

        synchronized void onConnected() {
            state = State.CONNECTED;
        }

        synchronized void onDisconnected() {
            state = State.DISCONNECTED;
        }

        synchronized void onExpired() {
            state = State.EXPIRED;
            ownedToken = -1L;
        }

        synchronized void acquire(long token) {
            if (state != State.CONNECTED) {
                throw new IllegalStateException("session is not connected");
            }
            ownedToken = token;
        }

        synchronized boolean allowWrite(long token) {
            return state == State.CONNECTED && ownedToken == token;
        }
    }

    public static void main(String[] args) {
        Gate gate = new Gate();
        gate.onConnected();
        gate.acquire(101L);
        System.out.println("beforeDisconnect=" + gate.allowWrite(101L));

        gate.onDisconnected();
        System.out.println("duringDisconnect=" + gate.allowWrite(101L));

        gate.onExpired();
        gate.onConnected();
        System.out.println("afterNewSessionWithOldToken=" + gate.allowWrite(101L));

        gate.acquire(102L);
        System.out.println("afterReacquire=" + gate.allowWrite(102L));
    }
}

预期输出:

text
beforeDisconnect=true
duringDisconnect=false
afterNewSessionWithOldToken=false
afterReacquire=true

真实系统中的 Token 必须由可单调比较的外部裁决生成,并由数据库、存储服务或设备网关拒绝旧 Token;仅在 JVM 内保存数值无法形成跨进程 Fencing。

十九、Curator配置监听Demo与版本裁决

java
import org.apache.curator.framework.CuratorFramework;
import org.apache.curator.framework.recipes.cache.CuratorCache;
import org.apache.curator.framework.recipes.cache.CuratorCacheListener;
import org.apache.zookeeper.data.Stat;

import java.nio.charset.StandardCharsets;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.atomic.AtomicInteger;

public class VersionedConfigWatcher {
    private final CuratorFramework client;
    private final AtomicInteger appliedVersion = new AtomicInteger(-1);
    private final ExecutorService executor = Executors.newSingleThreadExecutor();

    public VersionedConfigWatcher(CuratorFramework client) {
        this.client = client;
    }

    public CuratorCache start(String path) {
        CuratorCache cache = CuratorCache.build(client, path);
        CuratorCacheListener listener = CuratorCacheListener.builder()
                .forAll((type, oldData, data) -> executor.execute(() -> reload(path)))
                .build();
        cache.listenable().addListener(listener);
        cache.start();
        executor.execute(() -> reload(path));
        return cache;
    }

    private void reload(String path) {
        try {
            Stat stat = new Stat();
            byte[] bytes = client.getData().storingStatIn(stat).forPath(path);
            int version = stat.getVersion();
            while (true) {
                int current = appliedVersion.get();
                if (version <= current) return;
                if (appliedVersion.compareAndSet(current, version)) {
                    apply(new String(bytes, StandardCharsets.UTF_8), version);
                    return;
                }
            }
        } catch (Exception ex) {
            // 真实项目:分类重试、指标、告警;不能静默吞掉
            throw new IllegalStateException("reload config failed", ex);
        }
    }

    private void apply(String value, int version) {
        System.out.println("apply version=" + version + ", value=" + value);
    }
}

这个 Demo 展示“Watch提示 + 重新读取事实 + version防旧覆盖”。但 apply 若有外部副作用,还应先校验配置、构建新对象,成功后原子替换,失败时保留旧配置。

二十、Watcher在Dubbo注册发现中的位置

传统 Dubbo + ZooKeeper 的简化链路:

mermaid
flowchart TD
    A["Consumer订阅providers路径"] --> B["读取当前Provider子节点"]
    B --> C["注册子节点变化Watch"]
    C --> D["Provider注册或Session过期"]
    D --> E["ZooKeeper发出变化提示"]
    E --> F["Dubbo客户端重新读取完整地址集"]
    F --> G["RegistryDirectory对账Invoker快照"]

Watch 只触发重新发现,真正线程安全地复用、创建和销毁 Invoker 的逻辑在 Dubbo Directory。详细过程见Dubbo注册发现与Directory刷新

二十一、商业场景

21.1 动态配置

  • Watch 只触发重载。
  • 每次都读取完整配置和 Stat version。
  • 先解析、校验、预构建,再原子替换。
  • 新配置无效时保留旧配置并告警。

21.2 Leader任务

  • Disconnected 时暂停不可重复任务。
  • Expired 后永久撤销旧Leader身份。
  • 新Session重新竞选并取得新Fencing Token。
  • 外部事实源拒绝旧Token。

21.3 服务注册

  • 短暂断线不立即删除临时注册节点。
  • Session过期后客户端重新注册。
  • Consumer仍需超时和熔断应对陈旧地址窗口。

二十二、生产排查Runbook

22.1 频繁Disconnected

  1. 记录客户端、Server地址、Session ID脱敏值、协商Timeout和时间。
  2. 检查网络丢包、NAT、负载均衡、防火墙空闲连接策略。
  3. 查看ZooKeeper Server连接数、延迟、请求队列和GC。
  4. 查看客户端 JVM GC Pause、CPU throttling、线程调度和容器冻结。
  5. 检查是否所有客户端同时切换,判断Server故障还是局部网络。

22.2 频繁Session Expired

  1. 比较最大GC Pause、网络中断时长与协商Session Timeout。
  2. 检查客户端事件/I/O线程是否被阻塞。
  3. 检查集群Leader切换、Server过载和磁盘延迟。
  4. 不要只把Timeout无限调大;先解决停顿和容量问题。
  5. 验证业务在Expired后确实撤销旧锁和Leader身份。

22.3 配置变化没有生效

  1. 确认路径、chroot、ACL和环境。
  2. 读取当前数据和 Stat version,而不是只看回调日志。
  3. 检查标准Watch是否在首次触发后重新注册。
  4. 检查Cache是否完成初始化、是否已关闭。
  5. 检查回调执行器积压和业务校验失败。
  6. 确认旧版本回调没有覆盖新配置。

22.4 回调重复或顺序异常

Watch 回调不是业务幂等边界。记录 znode version/zxid 或业务 revision,只应用比当前更新的状态;对外部副作用使用业务幂等键。不要根据“回调次数”增减金额或库存。

二十三、必须监控的指标

  • 当前连接状态和状态持续时间。
  • Session Expired、AuthFailed次数。
  • 协商Session Timeout。
  • 重连次数和重连耗时。
  • Watch注册数、触发数、重载失败数。
  • 本地快照版本、年龄和最后成功重载时间。
  • 回调执行器队列长度。
  • 临时节点重建次数。
  • 业务安全门拒绝次数。
  • 旧Fencing Token被外部资源拒绝次数。

二十四、常见错误及后果

错误为什么错后果
TCP断开等于Session过期Session可在超时内跨Server恢复频繁误删业务状态或重复竞选
进程活着就仍持有锁GC Pause期间Session可能已过期两个持有者并发写
Watch只注册一次标准Watch一次触发后移除后续变化永久不再感知
Watch是可靠事件日志中间瞬态不应依赖逐条事件恢复本地状态与事实永久分叉
回调里执行慢HTTP阻塞事件处理和后续重载配置延迟、连接问题扩大
新Session继续使用旧节点路径临时顺序节点所有权已变化错误宣称锁或Leader身份
CuratorCache等于强一致本地库本地视图有传播和断连窗口用陈旧配置做高风险决策

二十五、面试标准回答

ZooKeeper Session 是跨单条TCP连接的逻辑会话。客户端短暂断连后可以在协商Timeout内连接其他Server并恢复原Session,因此临时节点不会在断线瞬间删除;服务端确认Session过期后,旧Session才不可恢复并删除其临时节点。标准Watch通常是一次性变化提示,existsgetDatagetChildren观察的事件范围不同,触发后客户端应重新读取完整当前状态并重新注册,而不能把Watch当作持久事件日志。ZooKeeper 3.6引入Persistent和Persistent Recursive Watch,但仍不能替代离线后的快照重建。锁和选主客户端在Disconnected时应进入安全模式,Expired后必须重新竞争,并使用Fencing Token防止GC暂停后恢复的旧持有者继续写外部资源。

二十六、关联知识点

本章小结

Session解决客户端在多次连接之间的逻辑身份,Watch解决状态变化提示,但二者都不是业务正确性的最终裁决。真正可靠的客户端必须把Disconnected视为不确定状态,把Expired视为旧所有权不可恢复的终态;把Watch当作缓存失效信号,每次重读事实并按版本原子发布;把锁和Leader身份交给Fencing Token与外部资源再次验证。只有这样,网络抖动、GC Pause、Server切换和事件窗口才不会演化成重复执行或双写。