本章定位:掌握 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 |
|---|---|---|
| create | data=hello, children=[] | 0 |
| setData | data=world, children=[] | 1 |
| create child | data=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 | 网络断开 | 等待自动重连 |
| Expired | Session 超时 | 重建 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:一次性的设计避免了通知风暴和内存泄漏。实现持续监听的方式:
- Watcher 回调中重新调用 getData/exists/getChildren 并传入自身
- 使用 Curator 的 NodeCache/TreeCache 自动处理循环注册
Q:原生 API 的异步回调在哪个线程执行?
A:在 EventThread 中执行。因此回调代码应尽量简短,不应包含阻塞操作。长时间处理应提交到独立的线程池。
小结
| 要点 | 说明 |
|---|---|
| 连接等待 | CountDownLatch 等待 SyncConnected |
| 同步 vs 异步 | 同步简单、异步高效(不阻塞) |
| 线程模型 | SendThread(IO)+ EventThread(事件) |
| Watcher 循环 | 回调中重新注册 |
| 连接管理 | 处理 Expired / Disconnected 状态 |
| 资源释放 | finally 中 close() |
原生 API 是理解 ZooKeeper 编程模型的最佳途径。下一章进入运维与监控,学习四字命令和监控体系。