乐途乐途
主页
  • 计算机基础

    • TCP/IP
    • Linux
    • HTTP
  • 数据库

    • SQL
    • MySQL 5.7
  • 编程语言

    • C
    • C++
    • Java SE
    • Python2
    • Python3
  • 数据格式

    • JSON
    • XML
  • 认证与安全

    • JWT
  • 工具

    • Markdown
  • Git

    • GitFlow
  • Quartz

    • Quartz
  • Java

    • Maven 入门
    • Maven 进阶
    • MyBatis
    • Spring
    • Spring MVC
  • Java

    • Spring Boot
    • Spring Cloud
    • Spring Cloud Alibaba
    • Spring Security
    • Spring AI
    • Spring Batch
    • Kafka
    • Java 设计模式
  • 缓存

    • Redis
  • 搜索引擎

    • Elasticsearch
  • 分布式协调

    • ZooKeeper
联系
阿里云
主页
  • 计算机基础

    • TCP/IP
    • Linux
    • HTTP
  • 数据库

    • SQL
    • MySQL 5.7
  • 编程语言

    • C
    • C++
    • Java SE
    • Python2
    • Python3
  • 数据格式

    • JSON
    • XML
  • 认证与安全

    • JWT
  • 工具

    • Markdown
  • Git

    • GitFlow
  • Quartz

    • Quartz
  • Java

    • Maven 入门
    • Maven 进阶
    • MyBatis
    • Spring
    • Spring MVC
  • Java

    • Spring Boot
    • Spring Cloud
    • Spring Cloud Alibaba
    • Spring Security
    • Spring AI
    • Spring Batch
    • Kafka
    • Java 设计模式
  • 缓存

    • Redis
  • 搜索引擎

    • Elasticsearch
  • 分布式协调

    • ZooKeeper
联系
阿里云
  • ZooKeeper 学习路径
  • 第1章 分布式协调与ZooKeeper概述

    • ZooKeeper 概述
  • 第2章 单机与集群搭建

    • 配置参数详解
    • 集群搭建
  • 第3章 数据模型与ZNode

    • ZNode 详解
    • 节点类型对比
    • 顺序节点
    • ACL 权限控制
  • 第4章 会话与Watcher机制

    • 会话机制
    • Watcher 机制
  • 第5章 ZAB协议与一致性保证

    • ZAB 协议
    • 一致性保证
    • 数据同步
  • 第6章 Leader选举

    • Leader 选举
  • 第7章 典型应用:分布式锁

    • 分布式锁
  • 第8章 典型应用:配置中心与命名服务

    • 配置中心
    • 命名服务
  • 第9章 客户端编程基础(Java原生API)

    • Java 原生 API 编程
  • 第10章 运维与监控

    • 四字命令
    • 监控体系
  • 第11章 面试考点与设计思想

    • 设计思想
    • 面试考点汇编

本章定位:掌握 Java 原生 ZooKeeper API 的 CRUD、Watcher、异步调用和连接管理。这是后续所有编程操作的基石。

定义与作用

ZooKeeper 官方提供 Java 客户端库(org.apache.zookeeper),通过 ZooKeeper 类暴露所有操作接口。原生 API 是 Curator 等高层封装的基础,理解它有助于深入掌握 ZooKeeper 的编程模型。

核心设计特点:

特点说明
异步为主所有操作提供同步和异步两个版本
Watcher 绑定Watcher 在操作时注册,非持久化
线程模型独立的 IO 线程(SendThread)+ 事件线程(EventThread)
连接管理自动重连,但 Session 状态需自行管理

核心原理

ZooKeeper 客户端线程模型

SendThread 负责所有网络 I/O 和心跳维持,EventThread 负责分发 Watcher 事件。Watcher 回调不能阻塞 EventThread。

API 方法分类

类别同步方法异步方法功能
创建create()create(..., cb, ctx)创建 ZNode
读取getData()getData(..., cb, ctx)获取数据 + Stat
更新setData()setData(..., cb, ctx)修改数据
删除delete()delete(..., cb, ctx)删除节点
存在检查exists()exists(..., cb, ctx)检查节点是否存在
子节点列表getChildren()getChildren(..., cb, ctx)获取子节点

完整示例

示例一:完整的 CRUD 操作

场景说明:从连接到销毁的完整 ZooKeeper 会话。

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

public class CRUDDemo {
    private static ZooKeeper zk;

    public static void main(String[] args) throws Exception {
        // 1. 连接 ZooKeeper
        CountDownLatch latch = new CountDownLatch(1);
        zk = new ZooKeeper("127.0.0.1:2181", 3000, event -> {
            if (event.getState() ==
                    Watcher.Event.KeeperState.SyncConnected) {
                latch.countDown();
            }
        });
        latch.await();
        System.out.println("Connected");

        // 2. 创建节点
        String path = zk.create("/crud-demo", "hello".getBytes(),
                ZooDefs.Ids.OPEN_ACL_UNSAFE,
                CreateMode.PERSISTENT);
        System.out.println("Created: " + path);

        // 3. 读取数据
        Stat stat = new Stat();
        byte[] data = zk.getData("/crud-demo", false, stat);
        System.out.println("Data: " + new String(data));
        System.out.println("Version: " + stat.getVersion());

        // 4. 更新数据(乐观锁)
        stat = zk.setData("/crud-demo", "world".getBytes(),
                stat.getVersion());
        System.out.println("Updated, new version: "
                + stat.getVersion());

        // 5. 创建子节点
        zk.create("/crud-demo/child", "child-data".getBytes(),
                ZooDefs.Ids.OPEN_ACL_UNSAFE,
                CreateMode.PERSISTENT);

        // 6. 获取子节点列表
        List<String> children =
                zk.getChildren("/crud-demo", false);
        System.out.println("Children: " + children);

        // 7. 删除(先删子节点再删父节点)
        zk.delete("/crud-demo/child", -1);
        zk.delete("/crud-demo", -1);
        System.out.println("Deleted");

        // 8. 关闭连接
        zk.close();
    }
}

执行结果:

Connected
Created: /crud-demo
Data: hello
Version: 0
Updated, new version: 1
Children: [child]
Deleted

操作前后对比:

操作/crud-demo 状态version
createdata=hello, children=[]0
setDatadata=world, children=[]1
create childdata=world, children=[child]cversion=1
delete all节点不存在—

示例二:异步 API 调用

场景说明:使用异步回调执行操作,不阻塞主线程。

public class AsyncDemo {
    public static void main(String[] args) throws Exception {
        ZooKeeper zk = new ZooKeeper(
                "127.0.0.1:2181", 3000, event -> {});

        // 异步创建
        zk.create("/async-node", "async-data".getBytes(),
                ZooDefs.Ids.OPEN_ACL_UNSAFE,
                CreateMode.PERSISTENT,
                (rc, path, ctx, name) -> {
                    System.out.println("Create result: " +
                            KeeperException.Code.get(rc));
                    System.out.println("Created path: " + name);
                },
                "create-context");

        // 异步读取
        zk.getData("/async-node", false,
                (rc, path, ctx, data, stat) -> {
                    System.out.println("Get result: " +
                            KeeperException.Code.get(rc));
                    System.out.println("Data: " +
                            new String(data));
                },
                "get-context");

        // 异步删除
        zk.delete("/async-node", -1,
                (rc, path, ctx) -> {
                    System.out.println("Delete result: " +
                            KeeperException.Code.get(rc));
                },
                "delete-context");

        Thread.sleep(2000);
        zk.close();
    }
}

执行结果(注意顺序可能是乱序的):

Create result: OK
Created path: /async-node
Get result: OK
Data: async-data
Delete result: OK

操作前后对比:

特性同步 API异步 API
调用方式直接返回结果通过回调获取结果
阻塞阻塞调用线程不阻塞
错误处理try-catch 异常rc 返回码判断
适用场景简单场景批量操作、高并发

示例三:连接状态监听与重连管理

场景说明:监听连接状态,处理 Session Expired。

public class ConnectionManager implements Watcher {
    private ZooKeeper zk;
    private CountDownLatch connectedSignal =
            new CountDownLatch(1);

    public ZooKeeper connect(String hosts) throws Exception {
        zk = new ZooKeeper(hosts, 5000, this);
        connectedSignal.await();
        return zk;
    }

    @Override
    public void process(WatchedEvent event) {
        System.out.printf("State: %s, Type: %s, Path: %s%n",
                event.getState(), event.getType(),
                event.getPath());

        switch (event.getState()) {
            case SyncConnected:
                connectedSignal.countDown();
                break;
            case Disconnected:
                System.out.println("Disconnected - waiting...");
                break;
            case Expired:
                System.out.println("Session Expired! " +
                        "Recreating connection...");
                try {
                    zk.close();
                    connect("127.0.0.1:2181");
                } catch (Exception e) {
                    e.printStackTrace();
                }
                break;
            case AuthFailed:
                System.out.println("Authentication failed!");
                break;
        }
    }

    public static void main(String[] args) throws Exception {
        ConnectionManager manager = new ConnectionManager();
        ZooKeeper zk = manager.connect("127.0.0.1:2181");

        // 使用 zk 进行操作...
        zk.create("/reconnect-test", null,
                ZooDefs.Ids.OPEN_ACL_UNSAFE,
                CreateMode.EPHEMERAL);

        // 模拟 Session 过期(不发送请求 >30s)
        Thread.sleep(35000);
        // 输出:State: Expired...
        // 自动重建连接
    }
}

操作前后对比:

连接状态触发条件处理策略
SyncConnected连接成功释放 CountDownLatch
Disconnected网络断开等待自动重连
ExpiredSession 超时重建 ZooKeeper 实例
AuthFailed认证失败检查凭据

易错场景与面试考点

易错场景

1. Watcher 回调中执行阻塞操作

// 错误:在 EventThread 中调用同步操作
public void process(WatchedEvent event) {
    byte[] data = zk.getData("/path", false, null); // 可能阻塞
}

应使用异步版本:

public void process(WatchedEvent event) {
    zk.getData("/path", false, (rc, path, ctx, data, stat) -> {
        // 处理结果
    }, null);
}

2. 忘记等待连接建立

// 错误:构造函数异步,可能还未连接
ZooKeeper zk = new ZooKeeper("...", 3000, watcher);
zk.create("/node", ...);  // ConnectionLossException

正确:使用 CountDownLatch 等待 SyncConnected 事件。

3. ZooKeeper 实例未 close

ZooKeeper 对象持有独立的 IO 和 Event 线程。不 close 会导致线程泄漏和连接未释放。应在 finally 块中 zk.close()。

面试高频题

Q:ZooKeeper 客户端的 SendThread 和 EventThread 分别做什么?

A:

  • SendThread:管理 Socket 连接、发送请求、接收响应、处理心跳(ping)
  • EventThread:从事件队列中消费 WatchedEvent 并回调 process() 方法

两个线程分离保证了 Watcher 回调不会阻塞网络 I/O。

Q:为什么原生 API 的 Watcher 是一次性的?如何实现持续监听?

A:一次性的设计避免了通知风暴和内存泄漏。实现持续监听的方式:

  1. Watcher 回调中重新调用 getData/exists/getChildren 并传入自身
  2. 使用 Curator 的 NodeCache/TreeCache 自动处理循环注册

Q:原生 API 的异步回调在哪个线程执行?

A:在 EventThread 中执行。因此回调代码应尽量简短,不应包含阻塞操作。长时间处理应提交到独立的线程池。

小结

要点说明
连接等待CountDownLatch 等待 SyncConnected
同步 vs 异步同步简单、异步高效(不阻塞)
线程模型SendThread(IO)+ EventThread(事件)
Watcher 循环回调中重新注册
连接管理处理 Expired / Disconnected 状态
资源释放finally 中 close()

原生 API 是理解 ZooKeeper 编程模型的最佳途径。下一章进入运维与监控,学习四字命令和监控体系。