大家好欢迎来到我的技术博客 在这里我会分享学习笔记、实战经验与技术思考力求用简单的方式讲清楚复杂的问题。 本文将围绕Zookeeper这个话题展开希望能为你带来一些启发或实用的参考。 无论你是刚入门的新手还是正在进阶的开发者希望你都能有所收获文章目录Zookeeper - 基于 Java API 的客户端连接实操开发 引言环境准备 创建 Maven 项目可选建立 Zookeeper 客户端连接 示例代码代码说明创建节点 示例代码代码说明获取节点数据 示例代码代码说明更新节点数据 示例代码代码说明删除节点 ️示例代码代码说明节点监听机制 示例代码代码说明Zookeeper 客户端连接状态管理 示例代码代码说明使用 Curator 框架简化开发 示例代码使用 Curator代码说明总结 Zookeeper - 基于 Java API 的客户端连接实操开发 引言Apache Zookeeper 是一个分布式协调服务广泛用于分布式系统的协调与管理。它提供了诸如命名服务、分布式同步、集群管理等功能是构建高可用、可扩展的分布式系统的重要工具之一。在实际开发中Zookeeper 的 Java API 是开发者最常用的接口之一。本文将详细介绍如何使用 Java API 实现 Zookeeper 客户端的连接、节点操作、监听机制等核心功能并通过代码示例帮助读者更好地理解和实践。环境准备 在开始之前我们需要准备好以下环境Java 8 或更高版本Zookeeper 3.5 或更高版本IDE如 IntelliJ IDEA 或 EclipseMaven 项目可选你可以从 Zookeeper 官方网站 下载并安装 Zookeeper。安装完成后确保 Zookeeper 服务已经启动。创建 Maven 项目可选如果你使用 Maven 来管理项目依赖可以在pom.xml文件中添加以下依赖项dependenciesdependencygroupIdorg.apache.zookeeper/groupIdartifactIdzookeeper/artifactIdversion3.7.1/version/dependency/dependencies⚠️ 注意请根据你使用的 Zookeeper 版本选择合适的依赖版本。建立 Zookeeper 客户端连接 Zookeeper 客户端的连接是通过ZooKeeper类来实现的。我们可以通过构造函数创建一个客户端实例并连接到 Zookeeper 服务器。示例代码importorg.apache.zookeeper.WatchedEvent;importorg.apache.zookeeper.Watcher;importorg.apache.zookeeper.ZooKeeper;importjava.io.IOException;publicclassZookeeperClientExample{publicstaticvoidmain(String[]args){StringhostPortlocalhost:2181;// Zookeeper 服务器地址intsessionTimeout3000;// 会话超时时间try{ZooKeeperzooKeepernewZooKeeper(hostPort,sessionTimeout,newWatcher(){Overridepublicvoidprocess(WatchedEventevent){// 监听事件回调System.out.println(Received event: event.getType());}});System.out.println(Connected to Zookeeper server);// 阻止主线程退出保持连接Thread.sleep(Long.MAX_VALUE);zooKeeper.close();}catch(IOException|InterruptedExceptione){e.printStackTrace();}}}代码说明ZooKeeper构造函数的参数说明hostPortZookeeper 服务器地址格式为host:port。sessionTimeout会话超时时间单位为毫秒。watcher监听器用于监听 Zookeeper 事件。process方法是Watcher接口的实现用于处理 Zookeeper 事件。Thread.sleep(Long.MAX_VALUE)用于保持主线程不退出从而保持与 Zookeeper 的连接。创建节点 在 Zookeeper 中节点ZNode是数据存储的基本单位。我们可以通过create方法创建节点。示例代码importorg.apache.zookeeper.CreateMode;importorg.apache.zookeeper.ZooDefs;importorg.apache.zookeeper.ZooKeeper;importorg.apache.zookeeper.data.Stat;importjava.io.IOException;importjava.util.concurrent.CountDownLatch;publicclassCreateZNodeExample{publicstaticvoidmain(String[]args)throwsIOException,InterruptedException{StringhostPortlocalhost:2181;intsessionTimeout3000;CountDownLatchconnectedSignalnewCountDownLatch(1);ZooKeeperzooKeepernewZooKeeper(hostPort,sessionTimeout,event-{if(event.getState()Event.KeeperState.SyncConnected){connectedSignal.countDown();}});connectedSignal.await();Stringpath/exampleNode;byte[]dataHello Zookeeper.getBytes();// 创建持久节点StringcreatedPathzooKeeper.create(path,data,ZooDefs.Ids.OPEN_ACL_UNSAFE,CreateMode.PERSISTENT);System.out.println(Created node at path: createdPath);zooKeeper.close();}}代码说明CreateMode.PERSISTENT创建一个持久节点。ZooDefs.Ids.OPEN_ACL_UNSAFE设置节点的 ACL访问控制列表这里使用的是开放权限。create方法返回新创建节点的路径。获取节点数据 我们可以使用getData方法获取节点的数据。示例代码importorg.apache.zookeeper.ZooKeeper;importorg.apache.zookeeper.data.Stat;importjava.io.IOException;importjava.util.concurrent.CountDownLatch;publicclassGetZNodeDataExample{publicstaticvoidmain(String[]args)throwsIOException,InterruptedException{StringhostPortlocalhost:2181;intsessionTimeout3000;CountDownLatchconnectedSignalnewCountDownLatch(1);ZooKeeperzooKeepernewZooKeeper(hostPort,sessionTimeout,event-{if(event.getState()Event.KeeperState.SyncConnected){connectedSignal.countDown();}});connectedSignal.await();Stringpath/exampleNode;StatstatnewStat();byte[]datazooKeeper.getData(path,false,stat);System.out.println(Node data: newString(data));zooKeeper.close();}}代码说明getData方法的参数说明path节点路径。watch是否设置监听器。stat用于存储节点的元数据。更新节点数据 我们可以使用setData方法更新节点的数据。示例代码importorg.apache.zookeeper.ZooKeeper;importorg.apache.zookeeper.data.Stat;importjava.io.IOException;importjava.util.concurrent.CountDownLatch;publicclassUpdateZNodeDataExample{publicstaticvoidmain(String[]args)throwsIOException,InterruptedException{StringhostPortlocalhost:2181;intsessionTimeout3000;CountDownLatchconnectedSignalnewCountDownLatch(1);ZooKeeperzooKeepernewZooKeeper(hostPort,sessionTimeout,event-{if(event.getState()Event.KeeperState.SyncConnected){connectedSignal.countDown();}});connectedSignal.await();Stringpath/exampleNode;byte[]newDataUpdated Data.getBytes();// 更新节点数据StatstatzooKeeper.setData(path,newData,-1);// -1 表示忽略版本号System.out.println(Node updated with version: stat.getVersion());zooKeeper.close();}}代码说明setData方法的参数说明path节点路径。data新的数据。version节点的版本号-1 表示忽略版本号。删除节点 ️我们可以使用delete方法删除节点。示例代码importorg.apache.zookeeper.ZooKeeper;importorg.apache.zookeeper.KeeperException;importjava.io.IOException;importjava.util.concurrent.CountDownLatch;publicclassDeleteZNodeExample{publicstaticvoidmain(String[]args)throwsIOException,InterruptedException,KeeperException{StringhostPortlocalhost:2181;intsessionTimeout3000;CountDownLatchconnectedSignalnewCountDownLatch(1);ZooKeeperzooKeepernewZooKeeper(hostPort,sessionTimeout,event-{if(event.getState()Event.KeeperState.SyncConnected){connectedSignal.countDown();}});connectedSignal.await();Stringpath/exampleNode;// 删除节点zooKeeper.delete(path,-1);// -1 表示忽略版本号System.out.println(Node deleted);zooKeeper.close();}}代码说明delete方法的参数说明path节点路径。version节点的版本号-1 表示忽略版本号。节点监听机制 Zookeeper 提供了强大的监听机制允许客户端监听节点的变化。我们可以通过exists、getData、getChildren等方法设置监听器。示例代码importorg.apache.zookeeper.WatchedEvent;importorg.apache.zookeeper.Watcher;importorg.apache.zookeeper.ZooKeeper;importorg.apache.zookeeper.data.Stat;importjava.io.IOException;importjava.util.concurrent.CountDownLatch;publicclassWatcherExample{publicstaticvoidmain(String[]args)throwsIOException,InterruptedException{StringhostPortlocalhost:2181;intsessionTimeout3000;CountDownLatchconnectedSignalnewCountDownLatch(1);ZooKeeperzooKeepernewZooKeeper(hostPort,sessionTimeout,newWatcher(){Overridepublicvoidprocess(WatchedEventevent){System.out.println(Received event: event.getType());if(event.getType()Event.EventType.Noneevent.getState()Event.KeeperState.SyncConnected){connectedSignal.countDown();}elseif(event.getType()Event.EventType.NodeDataChanged){System.out.println(Node data changed);}}});connectedSignal.await();Stringpath/exampleNode;// 设置监听器StatstatzooKeeper.exists(path,true);if(stat!null){System.out.println(Node exists);}else{System.out.println(Node does not exist);}// 阻止主线程退出保持连接Thread.sleep(Long.MAX_VALUE);zooKeeper.close();}}代码说明exists方法的参数说明path节点路径。watch是否设置监听器。当节点数据发生变化时监听器会收到NodeDataChanged事件。Zookeeper 客户端连接状态管理 在实际应用中Zookeeper 客户端可能会因为网络问题或服务器故障而断开连接。我们需要处理这些情况并实现自动重连机制。示例代码importorg.apache.zookeeper.WatchedEvent;importorg.apache.zookeeper.Watcher;importorg.apache.zookeeper.ZooKeeper;importorg.apache.zookeeper.KeeperException;importjava.io.IOException;importjava.util.concurrent.CountDownLatch;publicclassReconnectExample{privatestaticZooKeeperzooKeeper;publicstaticvoidmain(String[]args){StringhostPortlocalhost:2181;intsessionTimeout3000;connectToZookeeper(hostPort,sessionTimeout);// 模拟断开连接try{Thread.sleep(5000);System.out.println(Simulating connection loss...);zooKeeper.close();Thread.sleep(5000);connectToZookeeper(hostPort,sessionTimeout);}catch(InterruptedException|IOExceptione){e.printStackTrace();}}privatestaticvoidconnectToZookeeper(StringhostPort,intsessionTimeout){CountDownLatchconnectedSignalnewCountDownLatch(1);try{zooKeepernewZooKeeper(hostPort,sessionTimeout,event-{if(event.getState()Watcher.Event.KeeperState.SyncConnected){connectedSignal.countDown();}elseif(event.getState()Watcher.Event.KeeperState.Disconnected){System.out.println(Disconnected from Zookeeper);}elseif(event.getState()Watcher.Event.KeeperState.Expired){System.out.println(Session expired);try{connectToZookeeper(hostPort,sessionTimeout);}catch(IOException|InterruptedExceptione){e.printStackTrace();}}});connectedSignal.await();System.out.println(Connected to Zookeeper);}catch(IOException|InterruptedExceptione){e.printStackTrace();}}}代码说明当客户端断开连接或会话过期时监听器会收到相应的事件。在会话过期的情况下我们尝试重新连接到 Zookeeper 服务器。使用 Curator 框架简化开发 虽然 Zookeeper 提供了原生的 Java API但在实际开发中我们通常会使用 Curator 框架来简化开发。Curator 是 Apache 提供的一个 Zookeeper 客户端框架它封装了原生 API并提供了更高级的功能。示例代码使用 Curatorimportorg.apache.curator.framework.CuratorFramework;importorg.apache.curator.framework.CuratorFrameworkFactory;importorg.apache.curator.retry.ExponentialBackoffRetry;publicclassCuratorExample{publicstaticvoidmain(String[]args)throwsException{StringhostPortlocalhost:2181;// 创建 Curator 客户端CuratorFrameworkclientCuratorFrameworkFactory.newClient(hostPort,newExponentialBackoffRetry(1000,3));client.start();Stringpath/curatorNode;byte[]dataHello Curator.getBytes();// 创建节点client.create().forPath(path,data);System.out.println(Node created);// 获取节点数据byte[]retrievedDataclient.getData().forPath(path);System.out.println(Node data: newString(retrievedData));// 更新节点数据byte[]newDataUpdated with Curator.getBytes();client.setData().forPath(path,newData);System.out.println(Node updated);// 删除节点client.delete().forPath(path);System.out.println(Node deleted);client.close();}}代码说明CuratorFrameworkFactory.newClient创建 Curator 客户端。ExponentialBackoffRetry设置重试策略。create()、getData()、setData()、delete()Curator 提供的简化 API。总结 通过本文的学习我们了解了如何使用 Zookeeper 的 Java API 实现客户端连接、节点操作、监听机制等功能。我们还介绍了如何使用 Curator 框架简化开发并提供了相应的代码示例。希望本文能帮助读者更好地理解和实践 Zookeeper 的 Java 客户端开发。如果你对 Zookeeper 感兴趣可以进一步阅读 Zookeeper 官方文档了解更多高级功能和最佳实践。同时Curator 官方文档 也是学习 Curator 框架的好资源。Happy coding! 感谢你读到这里 技术之路没有捷径但每一次阅读、思考和实践都在悄悄拉近你与目标的距离。 如果本文对你有帮助不妨 点赞、收藏、分享给更多需要的朋友 欢迎在评论区留下你的想法、疑问或建议我会一一回复我们一起交流、共同成长 关注我不错过下一篇干货我们下期再见✨