3.2 Curator 常用 API


文档摘要

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 最核心的组件是 接口。

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 最核心的组件是 CuratorFramework 接口。它是所有 Curator API 的入口,代表了一个到 ZooKeeper 集群的连接。CuratorFramework 实例的创建通常通过 CuratorFrameworkFactory 工厂类完成。

1. 创建 CuratorFramework 实例

Curator 提供了多种方式来创建 CuratorFramework 实例,最常用的方式是通过 CuratorFrameworkFactorybuilder 方法,它允许我们以链式调用的方式配置客户端。

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 通常是一个较为稳妥的选择。

3.2.2 常用 API 操作:CRUD、检查、监听

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 对象,如果节点不存在则返回 nullStat 对象包含了节点的元数据信息,如版本号、创建时间等。

  • 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 图示:节点删除流程

3.2.3 异步操作 (Asynchronous Operations)

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 图示:异步操作流程

3.2.4 监听器 (Listeners / Watches)

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_ADDEDCHILD_UPDATEDCHILD_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: TreeCacheListenerPathChildrenCacheListener 类似,但事件类型为 TreeCacheEvent.Type,例如 NODE_ADDEDNODE_UPDATEDNODE_REMOVED 等,适用于树形结构的节点变化监听。

Mermaid 图示:监听器工作流程 (以 NodeCache 为例)

3.2.5 事务 (Transactions)

Curator 提供了事务 API,允许将多个 ZooKeeper 操作作为一个原子事务执行。事务保证了多个操作要么全部成功,要么全部失败,保持数据的一致性。


作者与出处
原作者: 灏天文库
来源:灏天文库
整理: 灏天文库整理
由灏天文库平台收录,内容或由平台用户上传,仅供学习交流
发布者: 作者: 灏天文库 转发
评论区 (0)
U