3.2 Curator 常用 API 第三章:Zookeeper 客户端 API 使用 3.2 Curator 常用 API 详解与实践 Apache Curator 是一个用于简化 Apache ZooKeeper 客户端开发的 Java 库。它在原生 ZooKeeper 客户端 API 的基础上进行了封装,提供了更高层次的抽象和更便捷的工具,极大地简化了 ZooKeeper 的使用,并提高了开发效率和应用稳定性。本章节将深入探讨 Curator 中常用的 API,并通过代码示例和图表进行详细解析。 3.2.1 CuratorFramework:Curator 的核心入口 Curator 最核心的组件是 接口。
Apache Curator 是一个用于简化 Apache ZooKeeper 客户端开发的 Java 库。它在原生 ZooKeeper 客户端 API 的基础上进行了封装,提供了更高层次的抽象和更便捷的工具,极大地简化了 ZooKeeper 的使用,并提高了开发效率和应用稳定性。本章节将深入探讨 Curator 中常用的 API,并通过代码示例和图表进行详细解析。
Curator 最核心的组件是 CuratorFramework 接口。它是所有 Curator API 的入口,代表了一个到 ZooKeeper 集群的连接。CuratorFramework 实例的创建通常通过 CuratorFrameworkFactory 工厂类完成。
1. 创建 CuratorFramework 实例
Curator 提供了多种方式来创建 CuratorFramework 实例,最常用的方式是通过 CuratorFrameworkFactory 的 builder 方法,它允许我们以链式调用的方式配置客户端。
import org.apache.curator.framework.CuratorFramework; import org.apache.curator.framework.CuratorFrameworkFactory; import org.apache.curator.retry.ExponentialBackoffRetry; public class CuratorClientExample { public static void main(String[] args) throws Exception { // ZooKeeper 连接地址 String connectString = "localhost:2181"; // 重试策略:初始 sleep 时间为 1 秒,最大重试次数为 3 次 ExponentialBackoffRetry retryPolicy = new ExponentialBackoffRetry(1000, 3); // 使用工厂类构建 CuratorFramework 实例 CuratorFramework client = CuratorFrameworkFactory.builder() .connectString(connectString) // 连接地址 .retryPolicy(retryPolicy) // 重试策略 .sessionTimeoutMs(60000) // 会话超时时间 (毫秒) .connectionTimeoutMs(15000) // 连接超时时间 (毫秒) .namespace("my_namespace") // 命名空间 (可选) .build(); // 启动客户端 client.start(); System.out.println("Curator client started successfully!"); // 在这里可以进行后续的 ZooKeeper 操作... // 关闭客户端 (在程序结束前) client.close(); System.out.println("Curator client closed."); } }
代码详解:
CuratorFrameworkFactory.builder(): 使用建造者模式创建 CuratorFramework。
.connectString(connectString): 设置 ZooKeeper 集群的连接地址。可以指定多个地址,用逗号分隔,Curator 会自动进行故障转移。
.retryPolicy(retryPolicy): 配置重试策略。ExponentialBackoffRetry 是一种常用的策略,它会随着重试次数的增加,指数级地增加重试之间的等待时间。
.sessionTimeoutMs(60000): 设置会话超时时间。ZooKeeper 客户端与服务端之间的会话如果超过这个时间没有心跳,服务端会认为会话失效。
.connectionTimeoutMs(15000): 设置连接超时时间。客户端尝试连接 ZooKeeper 服务端的最长等待时间。
.namespace("my_namespace"): 设置命名空间。所有后续的 Curator 操作都将在这个命名空间下进行,相当于在 ZooKeeper 根路径下创建了一个子目录。使用命名空间可以隔离不同应用的数据,避免冲突。
.build(): 完成配置,构建 CuratorFramework 实例。
.start(): 启动 Curator 客户端,建立与 ZooKeeper 集群的连接。
.close(): 关闭 Curator 客户端,释放资源。
Mermaid 图示:CuratorFramework 创建流程
2. 重试策略 (RetryPolicy)
Curator 提供了多种内置的重试策略,位于 org.apache.curator.retry 包下,常用的有:
ExponentialBackoffRetry: 指数退避重试。每次重试之间的等待时间呈指数增长,避免在 ZooKeeper 服务端压力过大时造成更大的负担。
RetryNTimes: 固定次数重试。尝试指定次数的重试,每次重试之间等待固定的时间。
RetryOneTime: 只重试一次。
RetryForever: 永远重试,直到成功。
UntilElapsedRetry: 在指定时间内重试,直到时间耗尽。
选择合适的重试策略对于保证应用的健壮性非常重要。ExponentialBackoffRetry 通常是一个较为稳妥的选择。
CuratorFramework 提供了丰富的 API 来操作 ZooKeeper 节点,包括创建 (Create)、读取 (Read)、更新 (Update)、删除 (Delete) 操作 (CRUD),以及检查节点是否存在、监听节点变化等功能。
1. 创建节点 (Create Operations)
Curator 使用 create() API 来创建 ZooKeeper 节点。它提供了多种选项来控制节点的创建行为。
import org.apache.curator.framework.api.CreateMode; public class CreateNodeExample { public static void main(String[] args) throws Exception { CuratorFramework client = CuratorClientExample.getClient(); // 获取 Curator 客户端实例 String path = "/my_node"; byte[] data = "Hello Curator!".getBytes(); try { // 1. 创建持久节点 client.create().forPath(path, data); System.out.println("Created persistent node: " + path); // 2. 创建临时节点 String ephemeralPath = "/ephemeral_node"; client.create().withMode(CreateMode.EPHEMERAL).forPath(ephemeralPath, data); System.out.println("Created ephemeral node: " + ephemeralPath); // 3. 创建顺序节点 String sequentialPath = "/sequential_node"; String createdSequentialPath = client.create().withMode(CreateMode.PERSISTENT_SEQUENTIAL).forPath(sequentialPath, data); System.out.println("Created sequential node: " + createdSequentialPath); // 4. 创建节点,如果父节点不存在则自动创建 String parentPath = "/parent_node/child_node"; client.create().creatingParentsIfNeeded().forPath(parentPath, data); System.out.println("Created node with parents: " + parentPath); } catch (Exception e) { e.printStackTrace(); } finally { CuratorClientExample.closeClient(client); // 关闭客户端 } } }
代码详解:
client.create(): 获取 CreateBuilder 接口,开始构建创建操作。
.forPath(path, data): 指定节点路径和节点数据。
.withMode(CreateMode.xxx): 设置节点类型,CreateMode 枚举定义了以下几种类型:
PERSISTENT: 持久节点,客户端断开连接后节点仍然存在。
EPHEMERAL: 临时节点,客户端断开连接后节点自动删除。
PERSISTENT_SEQUENTIAL: 持久顺序节点,持久节点的基础上,ZooKeeper 会在节点路径后自动追加一个单调递增的数字。
EPHEMERAL_SEQUENTIAL: 临时顺序节点,临时节点的基础上,ZooKeeper 会在节点路径后自动追加一个单调递增的数字。
.creatingParentsIfNeeded(): 如果父节点不存在,则自动创建父节点。这在创建深层路径节点时非常方便。
Mermaid 图示:节点创建流程
2. 读取节点数据 (Read Operations)
Curator 使用 getData() API 来读取 ZooKeeper 节点的数据。
public class GetDataExample { public static void main(String[] args) throws Exception { CuratorFramework client = CuratorClientExample.getClient(); String path = "/my_node"; try { byte[] data = client.getData().forPath(path); System.out.println("Data of node " + path + ": " + new String(data)); // 检查节点是否存在 Stat stat = client.checkExists().forPath(path); if (stat != null) { System.out.println("Node " + path + " exists."); } else { System.out.println("Node " + path + " does not exist."); } // 获取子节点列表 List<String> children = client.getChildren().forPath("/"); // 获取根节点下的子节点 System.out.println("Children of root node: " + children); } catch (Exception e) { e.printStackTrace(); } finally { CuratorClientExample.closeClient(client); } } }
代码详解:
client.getData().forPath(path): 获取指定路径节点的数据。返回值为 byte[]。
client.checkExists().forPath(path): 检查节点是否存在。返回值为 Stat 对象,如果节点不存在则返回 null。Stat 对象包含了节点的元数据信息,如版本号、创建时间等。
client.getChildren().forPath(path): 获取指定路径节点的子节点列表。返回值为 List<String>,包含了子节点的名称。
Mermaid 图示:节点数据读取流程
3. 更新节点数据 (Update Operations)
Curator 使用 setData() API 来更新 ZooKeeper 节点的数据。
public class SetDataExample { public static void main(String[] args) throws Exception { CuratorFramework client = CuratorClientExample.getClient(); String path = "/my_node"; byte[] newData = "Updated Curator Data!".getBytes(); try { // 更新节点数据 client.setData().forPath(path, newData); System.out.println("Updated data of node " + path); // 获取更新后的数据进行验证 byte[] updatedData = client.getData().forPath(path); System.out.println("New data of node " + path + ": " + new String(updatedData)); } catch (Exception e) { e.printStackTrace(); } finally { CuratorClientExample.closeClient(client); } } }
代码详解:
client.setData().forPath(path, newData): 设置指定路径节点的数据为 newData。Mermaid 图示:节点数据更新流程
4. 删除节点 (Delete Operations)
Curator 使用 delete() API 来删除 ZooKeeper 节点。
public class DeleteNodeExample { public static void main(String[] args) throws Exception { CuratorFramework client = CuratorClientExample.getClient(); String path = "/my_node"; try { // 1. 删除节点 client.delete().forPath(path); System.out.println("Deleted node: " + path); // 2. 强制删除子节点 (如果存在) String parentPath = "/parent_node"; client.delete().deletingChildrenIfNeeded().forPath(parentPath); System.out.println("Deleted node and its children: " + parentPath); // 3. 保证删除成功 (Guaranteed Delete) String guaranteedPath = "/guaranteed_node"; client.create().forPath(guaranteedPath); // 先创建节点 client.delete().guaranteed().forPath(guaranteedPath); System.out.println("Guaranteed deletion of node: " + guaranteedPath); } catch (Exception e) { e.printStackTrace(); } finally { CuratorClientExample.closeClient(client); } } }
代码详解:
client.delete().forPath(path): 删除指定路径的节点。如果节点有子节点,删除操作会失败。
.deletingChildrenIfNeeded(): 如果节点有子节点,则递归删除所有子节点。
.guaranteed(): 保证删除操作最终成功。即使删除操作失败 (例如,网络问题),Curator 会在后台不断重试,直到节点被成功删除。这对于关键节点的删除非常重要。
Mermaid 图示:节点删除流程
Curator 提供了异步 API,允许在后台执行 ZooKeeper 操作,而不会阻塞当前线程。异步操作可以提高应用的吞吐量和响应速度。
import org.apache.curator.framework.api.BackgroundCallback; import org.apache.curator.framework.api.CuratorEvent; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; public class AsyncOperationExample { public static void main(String[] args) throws Exception { CuratorFramework client = CuratorClientExample.getClient(); String path = "/async_node"; byte[] data = "Async Data".getBytes(); // 创建一个线程池用于执行回调 ExecutorService executor = Executors.newFixedThreadPool(2); try { // 异步创建节点 client.create().inBackground(new BackgroundCallback() { @Override public void processResult(CuratorFramework client, CuratorEvent event) throws Exception { System.out.println("Async create operation completed."); System.out.println("Event type: " + event.getType()); System.out.println("Result code: " + event.getResultCode()); System.out.println("Path: " + event.getPath()); // 在这里可以处理异步操作的结果 } }, executor).forPath(path, data); System.out.println("Async create operation submitted..."); // 主线程可以继续执行其他任务,异步操作在后台执行 Thread.sleep(2000); // 模拟主线程执行其他任务 } catch (Exception e) { e.printStackTrace(); } finally { executor.shutdown(); // 关闭线程池 CuratorClientExample.closeClient(client); } } }
代码详解:
.inBackground(BackgroundCallback callback, ExecutorService executor): 指定异步操作的回调函数 BackgroundCallback 和执行回调的线程池 ExecutorService。
BackgroundCallback 接口: 定义了 processResult 方法,用于处理异步操作的结果。CuratorEvent 对象包含了操作的结果信息,如事件类型、结果码、节点路径等。
ExecutorService: 用于执行回调函数的线程池。如果不指定 ExecutorService,Curator 会使用默认的线程池。
Mermaid 图示:异步操作流程
Curator 提供了强大的监听器机制,可以监听 ZooKeeper 节点的变化,例如节点数据变化、子节点列表变化等。Curator 提供了多种类型的监听器,简化了 ZooKeeper Watcher 的使用。
1. NodeCache:节点数据监听
NodeCache 用于监听单个节点的数据变化。当节点数据发生变化时,NodeCache 会自动更新缓存,并触发事件通知。
import org.apache.curator.framework.recipes.cache.NodeCache; import org.apache.curator.framework.recipes.cache.NodeCacheListener; public class NodeCacheExample { public static void main(String[] args) throws Exception { CuratorFramework client = CuratorClientExample.getClient(); String path = "/node_to_watch"; client.create().forPath(path, "Initial Data".getBytes()); // 创建节点 NodeCache nodeCache = new NodeCache(client, path); nodeCache.getListenable().addListener(new NodeCacheListener() { @Override public void nodeChanged() throws Exception { System.out.println("Node data changed for: " + path); System.out.println("New data: " + new String(nodeCache.getCurrentData().getData())); } }); nodeCache.start(); System.out.println("NodeCache started for: " + path); // 模拟节点数据变化 client.setData().forPath(path, "Updated Data".getBytes()); Thread.sleep(2000); // 等待事件触发 nodeCache.close(); CuratorClientExample.closeClient(client); } }
代码详解:
NodeCache nodeCache = new NodeCache(client, path): 创建 NodeCache 实例,指定要监听的节点路径。
nodeCache.getListenable().addListener(NodeCacheListener listener): 添加 NodeCacheListener 监听器。nodeChanged() 方法在节点数据变化时被调用。
nodeCache.start(): 启动 NodeCache,开始监听节点变化。
nodeCache.getCurrentData(): 获取当前节点的缓存数据。
2. PathChildrenCache:子节点列表监听
PathChildrenCache 用于监听指定路径下子节点列表的变化,包括子节点的创建、删除、更新等事件。
import org.apache.curator.framework.recipes.cache.PathChildrenCache; import org.apache.curator.framework.recipes.cache.PathChildrenCacheEvent; import org.apache.curator.framework.recipes.cache.PathChildrenCacheListener; public class PathChildrenCacheExample { public static void main(String[] args) throws Exception { CuratorFramework client = CuratorClientExample.getClient(); String parentPath = "/parent_to_watch"; client.create().forPath(parentPath); // 创建父节点 PathChildrenCache childrenCache = new PathChildrenCache(client, parentPath, true); // 第三个参数表示是否缓存节点数据 childrenCache.getListenable().addListener(new PathChildrenCacheListener() { @Override public void childEvent(CuratorFramework client, PathChildrenCacheEvent event) throws Exception { System.out.println("Child event received: " + event.getType()); switch (event.getType()) { case CHILD_ADDED: System.out.println("Child added: " + event.getData().getPath()); break; case CHILD_UPDATED: System.out.println("Child updated: " + event.getData().getPath()); break; case CHILD_REMOVED: System.out.println("Child removed: " + event.getData().getPath()); break; default: // 其他事件类型 break; } } }); childrenCache.start(); System.out.println("PathChildrenCache started for: " + parentPath); // 模拟子节点变化 client.create().forPath(parentPath + "/child1", "Child 1 Data".getBytes()); Thread.sleep(1000); client.setData().forPath(parentPath + "/child1", "Updated Child 1 Data".getBytes()); Thread.sleep(1000); client.delete().forPath(parentPath + "/child1"); Thread.sleep(2000); // 等待事件触发 childrenCache.close(); CuratorClientExample.closeClient(client); } }
代码详解:
PathChildrenCache childrenCache = new PathChildrenCache(client, parentPath, true): 创建 PathChildrenCache 实例,指定要监听的父节点路径。第三个参数 true 表示缓存子节点的数据。
childrenCache.getListenable().addListener(PathChildrenCacheListener listener): 添加 PathChildrenCacheListener 监听器。childEvent() 方法在子节点发生变化时被调用。
PathChildrenCacheEvent.getType(): 获取事件类型,PathChildrenCacheEvent.Type 枚举定义了子节点变化的类型,如 CHILD_ADDED、CHILD_UPDATED、CHILD_REMOVED 等。
event.getData(): 获取事件相关的数据,例如 ChildData 对象包含了子节点的路径和数据。
3. TreeCache:树形节点监听
TreeCache 是功能更强大的监听器,它可以监听指定路径及其所有子孙节点的变化,形成一个树形结构的缓存。
import org.apache.curator.framework.recipes.cache.TreeCache; import org.apache.curator.framework.recipes.cache.TreeCacheEvent; import org.apache.curator.framework.recipes.cache.TreeCacheListener; public class TreeCacheExample { public static void main(String[] args) throws Exception { CuratorFramework client = CuratorClientExample.getClient(); String rootPath = "/tree_root"; client.create().forPath(rootPath); // 创建根节点 TreeCache treeCache = new TreeCache(client, rootPath); treeCache.getListenable().addListener(new TreeCacheListener() { @Override public void childEvent(CuratorFramework client, TreeCacheEvent event) throws Exception { System.out.println("TreeCache event received: " + event.getType()); switch (event.getType()) { case NODE_ADDED: System.out.println("Node added: " + event.getData().getPath()); break; case NODE_UPDATED: System.out.println("Node updated: " + event.getData().getPath()); break; case NODE_REMOVED: System.out.println("Node removed: " + event.getData().getPath()); break; default: // 其他事件类型 break; } } }); treeCache.start(); System.out.println("TreeCache started for: " + rootPath); // 模拟树形结构节点变化 client.create().creatingParentsIfNeeded().forPath(rootPath + "/level1/level2", "Data in Level 2".getBytes()); Thread.sleep(1000); client.setData().forPath(rootPath + "/level1/level2", "Updated Data in Level 2".getBytes()); Thread.sleep(1000); client.delete().deletingChildrenIfNeeded().forPath(rootPath + "/level1"); // 删除 level1 及其子节点 Thread.sleep(2000); // 等待事件触发 treeCache.close(); CuratorClientExample.closeClient(client); } }
代码详解:
TreeCache treeCache = new TreeCache(client, rootPath): 创建 TreeCache 实例,指定要监听的根路径。
TreeCacheListener: TreeCacheListener 与 PathChildrenCacheListener 类似,但事件类型为 TreeCacheEvent.Type,例如 NODE_ADDED、NODE_UPDATED、NODE_REMOVED 等,适用于树形结构的节点变化监听。
Mermaid 图示:监听器工作流程 (以 NodeCache 为例)
Curator 提供了事务 API,允许将多个 ZooKeeper 操作作为一个原子事务执行。事务保证了多个操作要么全部成功,要么全部失败,保持数据的一致性。