3.4 原生 Zookeeper API 常用 API 第三章:Zookeeper 客户端 API 使用 - 3.4 原生 Zookeeper API 常用 API 详解与代码实践 3.4 原生 Zookeeper API 常用 API 我们将重点关注以下几个核心 API 领域: 连接管理 (Connection Management):如何建立和管理与 Zookeeper 集群的连接,包括同步和异步连接方式。 节点操作 (Node Operations): Zookeeper 的核心概念是节点 (ZNode),我们将学习如何创建、删除、检查节点是否存在等操作。
我们将重点关注以下几个核心 API 领域:
连接管理 (Connection Management):如何建立和管理与 Zookeeper 集群的连接,包括同步和异步连接方式。
节点操作 (Node Operations): Zookeeper 的核心概念是节点 (ZNode),我们将学习如何创建、删除、检查节点是否存在等操作。
数据操作 (Data Operations): 如何读写 ZNode 节点的数据,这是 Zookeeper 存储和共享配置信息的关键。
监听机制 (Watcher Mechanism): Zookeeper 的事件通知机制,允许客户端监听节点的变化,实现实时的配置更新和分布式协调。
ACL 权限控制 (Access Control Lists): Zookeeper 的权限管理机制,保障数据的安全性。
为了更好地理解和实践,我们将结合 Java 代码示例,并使用 Mermaid 图形化工具来辅助说明 Zookeeper 的工作原理和 API 的使用流程。
与 Zookeeper 集群建立连接是进行任何操作的前提。原生 Zookeeper API 提供了 ZooKeeper 类来实现连接管理。
1. ZooKeeper 构造方法:建立连接
ZooKeeper 类提供了多个构造方法,但最常用的构造方法如下:
public ZooKeeper(String connectString, int sessionTimeout, Watcher watcher) throws IOException;
connectString: Zookeeper 集群的连接地址字符串。可以指定多个地址,以逗号分隔,例如 "192.168.1.100:2181,192.168.1.101:2181,192.168.1.102:2181"。客户端会尝试连接列表中的地址,直到成功连接到集群中的一个节点。
sessionTimeout: 会话超时时间,单位毫秒。如果客户端在指定时间内没有向服务器发送心跳,服务器会认为客户端连接断开,并清理会话相关的资源。
watcher: 一个 Watcher 接口的实现类,用于接收来自 Zookeeper 服务器的事件通知。例如,连接状态变化、节点数据变化等。
代码示例:同步连接
import org.apache.zookeeper.*; import java.io.IOException; import java.util.concurrent.CountDownLatch; public class ConnectionDemo { private static final String CONNECT_STRING = "127.0.0.1:2181"; // 替换为你的 Zookeeper 地址 private static final int SESSION_TIMEOUT = 5000; public static void main(String[] args) throws IOException, InterruptedException { final CountDownLatch connectedSignal = new CountDownLatch(1); ZooKeeper zk = new ZooKeeper(CONNECT_STRING, SESSION_TIMEOUT, new Watcher() { @Override public void process(WatchedEvent event) { if (event.getState() == Event.KeeperState.SyncConnected) { connectedSignal.countDown(); System.out.println("已连接到 Zookeeper!"); } else if (event.getState() == Event.KeeperState.Disconnected) { System.out.println("与 Zookeeper 断开连接!"); } } }); connectedSignal.await(); // 等待连接建立成功 System.out.println("Zookeeper 状态: " + zk.getState()); zk.close(); // 关闭连接 } }
代码详解:
CountDownLatch connectedSignal = new CountDownLatch(1);: 使用 CountDownLatch 实现同步等待,确保在连接建立成功后再进行后续操作。
new ZooKeeper(...): 创建 ZooKeeper 实例,传入连接字符串、会话超时时间和 Watcher 实例。
Watcher 实现: 匿名内部类实现了 Watcher 接口的 process 方法。
event.getState() == Event.KeeperState.SyncConnected: 当连接状态变为 SyncConnected 时,表示连接建立成功,调用 connectedSignal.countDown() 释放等待。
event.getState() == Event.KeeperState.Disconnected: 当连接断开时,打印日志信息。
connectedSignal.await();: 主线程等待 connectedSignal 计数器变为 0,即等待连接建立成功。
zk.getState(): 获取当前 Zookeeper 连接状态。
zk.close(): 关闭 Zookeeper 连接,释放资源。
Mermaid 图示:连接建立流程
2. ZooKeeper.close():关闭连接
使用 ZooKeeper.close() 方法可以关闭与 Zookeeper 集群的连接。这是一个重要的操作,应该在不再需要使用 Zookeeper 连接时及时关闭,释放资源。
Zookeeper 的数据模型是树状结构的命名空间,每个节点称为 ZNode。原生 API 提供了丰富的节点操作方法。
1. create():创建节点
create() 方法用于在 Zookeeper 中创建节点。常用的 create() 方法签名如下:
public String create(String path, byte[] data, List<ACL> acl, CreateMode createMode) throws KeeperException, InterruptedException;
path: 要创建的节点路径,例如 "/myNode", "/app/config"。
data: 节点存储的数据,字节数组形式。可以为空 (null)。
acl: 访问控制列表 (ACL),用于设置节点的访问权限。可以使用 ZooDefs.Ids 中预定义的 ACL 方案,例如 ZooDefs.Ids.OPEN_ACL_UNSAFE (开放权限)。
createMode: 创建模式,定义节点的类型和特性。常用的创建模式包括:
CreateMode.PERSISTENT: 持久节点,创建后一直存在,即使创建该节点的客户端断开连接。
CreateMode.EPHEMERAL: 临时节点,与创建该节点的客户端会话绑定。当会话结束 (客户端断开连接或会话超时),临时节点会被自动删除。
CreateMode.PERSISTENT_SEQUENTIAL: 持久顺序节点,具有持久节点的特性,并且在节点路径末尾自动追加一个单调递增的序列号。
CreateMode.EPHEMERAL_SEQUENTIAL: 临时顺序节点,具有临时节点的特性,并且在节点路径末尾自动追加一个单调递增的序列号。
代码示例:创建持久节点和临时节点
import org.apache.zookeeper.*; import org.apache.zookeeper.data.Stat; import java.io.IOException; import java.util.concurrent.CountDownLatch; public class CreateNodeDemo { private static final String CONNECT_STRING = "127.0.0.1:2181"; private static final int SESSION_TIMEOUT = 5000; private static ZooKeeper zk; private static final CountDownLatch connectedSignal = new CountDownLatch(1); public static void main(String[] args) throws IOException, InterruptedException, KeeperException { zk = new ZooKeeper(CONNECT_STRING, SESSION_TIMEOUT, new Watcher() { @Override public void process(WatchedEvent event) { if (event.getState() == Event.KeeperState.SyncConnected) { connectedSignal.countDown(); } } }); connectedSignal.await(); // 创建持久节点 String persistentPath = "/persistentNode"; String persistentData = "This is a persistent node"; String createdPersistentPath = zk.create(persistentPath, persistentData.getBytes(), ZooDefs.Ids.OPEN_ACL_UNSAFE, CreateMode.PERSISTENT); System.out.println("创建持久节点成功,路径: " + createdPersistentPath); // 创建临时节点 String ephemeralPath = "/ephemeralNode"; String ephemeralData = "This is an ephemeral node"; String createdEphemeralPath = zk.create(ephemeralPath, ephemeralData.getBytes(), ZooDefs.Ids.OPEN_ACL_UNSAFE, CreateMode.EPHEMERAL); System.out.println("创建临时节点成功,路径: " + createdEphemeralPath); zk.close(); } }
代码详解:
zk.create(persistentPath, persistentData.getBytes(), ZooDefs.Ids.OPEN_ACL_UNSAFE, CreateMode.PERSISTENT): 创建持久节点 /persistentNode,数据为 "This is a persistent node",ACL 为开放权限,创建模式为 PERSISTENT。
zk.create(ephemeralPath, ephemeralData.getBytes(), ZooDefs.Ids.OPEN_ACL_UNSAFE, CreateMode.EPHEMERAL): 创建临时节点 /ephemeralNode,数据为 "This is an ephemeral node",ACL 为开放权限,创建模式为 EPHEMERAL。
2. delete():删除节点
delete() 方法用于删除 Zookeeper 中的节点。方法签名如下:
public void delete(String path, int version) throws InterruptedException, KeeperException;
path: 要删除的节点路径。
version: 数据版本号。Zookeeper 使用版本号来实现乐观锁机制。
如果 version 为 -1,则表示忽略版本检查,直接删除节点。
如果 version 大于等于 0,则只有当节点的当前版本号与指定的 version 相匹配时,才能删除节点。
代码示例:删除节点
import org.apache.zookeeper.*; import java.io.IOException; import java.util.concurrent.CountDownLatch; public class DeleteNodeDemo { private static final String CONNECT_STRING = "127.0.0.1:2181"; private static final int SESSION_TIMEOUT = 5000; private static ZooKeeper zk; private static final CountDownLatch connectedSignal = new CountDownLatch(1); public static void main(String[] args) throws IOException, InterruptedException, KeeperException { zk = new ZooKeeper(CONNECT_STRING, SESSION_TIMEOUT, new Watcher() { @Override public void process(WatchedEvent event) { if (event.getState() == Event.KeeperState.SyncConnected) { connectedSignal.countDown(); } } }); connectedSignal.await(); String pathToDelete = "/persistentNode"; // 假设节点已存在 // 删除节点,忽略版本检查 zk.delete(pathToDelete, -1); System.out.println("节点删除成功,路径: " + pathToDelete); zk.close(); } }
代码详解:
zk.delete(pathToDelete, -1): 删除路径为 /persistentNode 的节点,版本号设置为 -1,表示忽略版本检查。3. exists():检查节点是否存在
exists() 方法用于检查指定路径的节点是否存在。方法签名如下:
public Stat exists(String path, boolean watch) throws KeeperException, InterruptedException;
path: 要检查的节点路径。
watch: 是否设置 Watcher 监听节点的存在事件。
如果设置为 true,当节点被创建、删除或节点数据发生变化时,会触发 Watcher 的 process 方法。
如果设置为 false,则不设置 Watcher。
exists() 方法返回一个 Stat 对象,包含节点的元数据信息,例如版本号、创建时间、修改时间等。如果节点不存在,则返回 null。
代码示例:检查节点是否存在
import org.apache.zookeeper.*; import org.apache.zookeeper.data.Stat; import java.io.IOException; import java.util.concurrent.CountDownLatch; public class ExistsNodeDemo { private static final String CONNECT_STRING = "127.0.0.1:2181"; private static final int SESSION_TIMEOUT = 5000; private static ZooKeeper zk; private static final CountDownLatch connectedSignal = new CountDownLatch(1); public static void main(String[] args) throws IOException, InterruptedException, KeeperException { zk = new ZooKeeper(CONNECT_STRING, SESSION_TIMEOUT, new Watcher() { @Override public void process(WatchedEvent event) { if (event.getState() == Event.KeeperState.SyncConnected) { connectedSignal.countDown(); } } }); connectedSignal.await(); String pathToCheck = "/persistentNode"; Stat stat = zk.exists(pathToCheck, false); if (stat != null) { System.out.println("节点存在,路径: " + pathToCheck); System.out.println("节点版本号: " + stat.getVersion()); } else { System.out.println("节点不存在,路径: " + pathToCheck); } zk.close(); } }
代码详解:
zk.exists(pathToCheck, false): 检查路径为 /persistentNode 的节点是否存在,不设置 Watcher。
stat != null: 判断 exists() 方法的返回值是否为 null,如果不是 null,则表示节点存在。
4. getChildren():获取子节点列表
getChildren() 方法用于获取指定节点路径下的所有子节点列表。方法签名如下:
public List<String> getChildren(String path, boolean watch) throws KeeperException, InterruptedException;
path: 父节点路径。
watch: 是否设置 Watcher 监听子节点的变化事件。
getChildren() 方法返回一个 List<String>,包含所有子节点的名称。
代码示例:获取子节点列表
import org.apache.zookeeper.*; import java.io.IOException; import java.util.List; import java.util.concurrent.CountDownLatch; public class GetChildrenDemo { private static final String CONNECT_STRING = "127.0.0.1:2181"; private static final int SESSION_TIMEOUT = 5000; private static ZooKeeper zk; private static final CountDownLatch connectedSignal = new CountDownLatch(1); public static void main(String[] args) throws IOException, InterruptedException, KeeperException { zk = new ZooKeeper(CONNECT_STRING, SESSION_TIMEOUT, new Watcher() { @Override public void process(WatchedEvent event) { if (event.getState() == Event.KeeperState.SyncConnected) { connectedSignal.countDown(); } } }); connectedSignal.await(); String parentPath = "/"; // 获取根节点下的子节点 List<String> children = zk.getChildren(parentPath, false); System.out.println("子节点列表 (路径: " + parentPath + "):"); for (String child : children) { System.out.println(child); } zk.close(); } }
代码详解:
zk.getChildren(parentPath, false): 获取路径为 / (根节点) 的子节点列表,不设置 Watcher。
遍历 children 列表: 打印获取到的子节点名称。
Zookeeper 的节点可以存储少量数据,这些数据通常用于配置信息、状态信息等。
1. getData():获取节点数据
getData() 方法用于获取指定节点路径的数据。方法签名如下:
public byte[] getData(String path, boolean watch, Stat stat) throws KeeperException, InterruptedException;
path: 节点路径。
watch: 是否设置 Watcher 监听节点数据变化事件.
stat: 用于接收节点元数据信息的 Stat 对象。可以为 null,表示不需要获取元数据。
getData() 方法返回节点的数据,字节数组形式。
代码示例:获取节点数据
import org.apache.zookeeper.*; import org.apache.zookeeper.data.Stat; import java.io.IOException; import java.util.concurrent.CountDownLatch; public class GetDataDemo { private static final String CONNECT_STRING = "127.0.0.1:2181"; private static final int SESSION_TIMEOUT = 5000; private static ZooKeeper zk; private static final CountDownLatch connectedSignal = new CountDownLatch(1); public static void main(String[] args) throws IOException, InterruptedException, KeeperException { zk = new ZooKeeper(CONNECT_STRING, SESSION_TIMEOUT, new Watcher() { @Override public void process(WatchedEvent event) { if (event.getState() == Event.KeeperState.SyncConnected) { connectedSignal.countDown(); } } }); connectedSignal.await(); String pathGetData = "/persistentNode"; // 假设节点已存在并有数据 Stat stat = new Stat(); byte[] data = zk.getData(pathGetData, false, stat); String value = new String(data); System.out.println("节点数据 (路径: " + pathGetData + "): " + value); System.out.println("节点版本号: " + stat.getVersion()); zk.close(); } }
代码详解:
Stat stat = new Stat();: 创建一个 Stat 对象用于接收节点元数据。
zk.getData(pathGetData, false, stat): 获取路径为 /persistentNode 的节点数据,不设置 Watcher,并将元数据信息写入 stat 对象。
String value = new String(data);: 将字节数组数据转换为字符串。
stat.getVersion(): 获取节点的版本号。
2. setData():设置节点数据
setData() 方法用于设置指定节点路径的数据。方法签名如下:
public Stat setData(String path, byte[] data, int version) throws KeeperException, InterruptedException;
path: 节点路径。
data: 要设置的数据,字节数组形式。
version: 数据版本号,用于乐观锁机制。
如果 version 为 -1,则表示忽略版本检查,直接设置数据。
如果 version 大于等于 0,则只有当节点的当前版本号与指定的 version 相匹配时,才能设置数据。
setData() 方法返回一个 Stat 对象,包含更新后的节点元数据信息。
代码示例:设置节点数据
import org.apache.zookeeper.*; import org.apache.zookeeper.data.Stat; import java.io.IOException; import java.util.concurrent.CountDownLatch; public class SetDataDemo { private static final String CONNECT_STRING = "127.0.0.1:2181"; private static final int SESSION_TIMEOUT = 5000; private static ZooKeeper zk; private static final CountDownLatch connectedSignal = new CountDownLatch(1); public static void main(String[] args) throws IOException, InterruptedException, KeeperException { zk = new ZooKeeper(CONNECT_STRING, SESSION_TIMEOUT, new Watcher() { @Override public void process(WatchedEvent event) { if (event.getState() == Event.KeeperState.SyncConnected) { connectedSignal.countDown(); } } }); connectedSignal.await(); String pathSetData = "/persistentNode"; // 假设节点已存在 String newData = "Updated data for persistent node"; Stat stat = zk.setData(pathSetData, newData.getBytes(), -1); System.out.println("节点数据设置成功 (路径: " + pathSetData + ")"); System.out.println("更新后节点版本号: " + stat.getVersion()); zk.close(); } }
代码详解:
zk.setData(pathSetData, newData.getBytes(), -1): 设置路径为 /persistentNode 的节点数据为 newData,版本号设置为 -1,表示忽略版本检查。
stat.getVersion(): 获取更新后的节点版本号。
Watcher 是 Zookeeper 的核心特性之一,它允许客户端注册监听器,当 Zookeeper 上的数据发生变化时,服务器会主动通知客户端。
1. Watcher 接口
Watcher 接口定义了处理事件通知的方法:
public interface Watcher { void process(WatchedEvent event); }
process(WatchedEvent event): 当监听的事件发生时,Zookeeper 服务器会调用客户端注册的 Watcher 的 process 方法,并将事件信息封装在 WatchedEvent 对象中传递给客户端。WatchedEvent 对象包含以下信息:
EventType getType(): 事件类型,例如 NodeCreated, NodeDeleted, NodeDataChanged, NodeChildrenChanged。
KeeperState getState(): 连接状态,例如 SyncConnected, Disconnected, Expired.
String getPath(): 发生事件的节点路径。
2. 设置 Watcher 的 API
在之前的 API 中,例如 exists(), getData(), getChildren(), 都提供了 watch 参数,用于设置 Watcher。
exists(String path, boolean watch): 监听节点的存在事件 (节点创建、删除、数据变化)。
getData(String path, boolean watch, Stat stat): 监听节点的数据变化事件。
getChildren(String path, boolean watch): 监听子节点的变化事件 (子节点创建、删除)。
注意:
一次性触发: Watcher 是一次性触发的。当事件被触发后,Watcher 会被移除。如果需要持续监听,需要在 process 方法中重新注册 Watcher。
异步通知: 事件通知是异步的,客户端接收到事件通知的顺序可能与事件发生的顺序不一致。
会话相关: Watcher 与客户端会话绑定。当会话失效或过期时,注册在该会话上的所有 Watcher 都会失效。
代码示例:监听节点数据变化
import org.apache.zookeeper.*; import org.apache.zookeeper.data.Stat; import java.io.IOException; import java.util.concurrent.CountDownLatch; public class WatcherDemo { private static final String CONNECT_STRING = "127.0.0.1:2181"; private static final int SESSION_TIMEOUT = 5000; private static ZooKeeper zk; private static final CountDownLatch connectedSignal = new CountDownLatch(1); public static void main(String[] args) throws IOException, InterruptedException, KeeperException { zk = new ZooKeeper(CONNECT_STRING, SESSION_TIMEOUT, new Watcher() { @Override public void process(WatchedEvent event) { if (event.getState() == Event.KeeperState.SyncConnected) { connectedSignal.countDown(); } } }); connectedSignal.await(); String watchPath = "/watchNode"; createNodeIfNotExists(watchPath, "Initial Data"); // 注册数据变化 Watcher zk.getData(watchPath, new Watcher() { @Override public void process(WatchedEvent event) { System.out.println("Watcher 触发,事件类型: " + event.getType()); if (event.getType() == Event.EventType.NodeDataChanged) { try { byte[] data = zk.getData(watchPath, this, null); // 再次注册 Watcher System.out.println("最新数据: " + new String(data)); } catch (KeeperException e) { e.printStackTrace(); } catch (InterruptedException e) { e.printStackTrace(); } } } }, null); System.out.println("已注册 Watcher 监听节点数据变化: " + watchPath); // 模拟数据变化 Thread.sleep(5000); // 等待一段时间 setData(watchPath, "Updated Data"); Thread.sleep(10000); // 保持程序运行,等待 Watcher 触发 zk.close(); } private static void createNodeIfNotExists(String path, String data) throws KeeperException, InterruptedException { if (zk.exists(path, false) == null) { zk.create(path, data.getBytes(), ZooDefs.Ids.OPEN_ACL_UNSAFE, CreateMode.PERSISTENT); System.out.println("节点创建成功: " + path); } } private static void setData(String path, String data) throws KeeperException, InterruptedException { zk.setData(path, data.getBytes(), -1); System.out.println("节点数据更新成功: " + path); } }
代码详解:
zk.getData(watchPath, new Watcher() { ... }, null): 在 getData() 方法中注册 Watcher 监听 /watchNode 节点的数据变化事件。
Watcher.process(WatchedEvent event): Watcher 的 process 方法被触发。
event.getType() == Event.EventType.NodeDataChanged: 判断事件类型是否为 NodeDataChanged。
zk.getData(watchPath, this, null): 在 process 方法中再次调用 getData() 注册 Watcher,实现持续监听。this 关键字表示将当前 Watcher 实例再次注册。
setData(watchPath, "Updated Data"): 模拟修改节点数据,触发 Watcher 事件。
Mermaid 图示:Watcher 工作流程
原生 Zookeeper API 也提供了异步操作的方式,以提高客户端的性能和响应速度。异步 API 的方法通常以 Async 结尾,例如 createAsync(), getDataAsync(), setDataAsync() 等。
异步 API 的使用方式:
异步 API 方法不会立即返回结果,而是立即返回,并将结果通过回调函数 (Callback) 异步通知给客户端。
常用的异步回调接口:
AsyncCallback.StringCallback: 用于 createAsync() 方法,回调参数为创建的节点路径。
AsyncCallback.DataCallback: 用于 getDataAsync() 方法,回调参数为节点数据和 Stat 对象。
AsyncCallback.StatCallback: 用于 setDataAsync() 和 deleteAsync() 方法,回调参数为 Stat 对象。
AsyncCallback.ChildrenCallback: 用于 getChildrenAsync() 方法,回调参数为子节点列表。
AsyncCallback.VoidCallback: 用于不需要返回结果的操作,例如 existsAsync(),回调参数为返回码 (rc)。