第一章:Zookeeper 基础概念 第一章:Zookeeper 基础概念 1.1 Zookeeper 的定义与定位 Zookeeper 官方定义为一个分布式协调服务,它暴露了一组简单的原语,分布式应用程序可以基于它们实现更高层次的同步、配置维护及分组服务等。更通俗地来说,Zookeeper 就像一个分布式系统的“注册中心”和“协调器”,它帮助不同的服务器节点协同工作,保证数据一致性和系统稳定性。 在分布式系统中,我们常常面临以下挑战: 数据一致性: 多个服务器节点需要访问和修改共享数据,如何保证数据在各个节点之间的一致性? 配置管理: 分布式应用的配置信息通常需要在集群中共享,如何高效、可靠地进行配置管理和动态更新?
Zookeeper 官方定义为一个分布式协调服务,它暴露了一组简单的原语,分布式应用程序可以基于它们实现更高层次的同步、配置维护及分组服务等。更通俗地来说,Zookeeper 就像一个分布式系统的“注册中心”和“协调器”,它帮助不同的服务器节点协同工作,保证数据一致性和系统稳定性。
在分布式系统中,我们常常面临以下挑战:
数据一致性: 多个服务器节点需要访问和修改共享数据,如何保证数据在各个节点之间的一致性?
配置管理: 分布式应用的配置信息通常需要在集群中共享,如何高效、可靠地进行配置管理和动态更新?
服务发现: 服务提供者和消费者需要互相发现对方,如何实现动态的服务注册和发现?
领导者选举: 在分布式集群中,需要选举一个领导者节点来协调工作,如何可靠地进行领导者选举?
分布式锁: 在并发环境下,如何控制对共享资源的访问,避免数据冲突?
Zookeeper 正是为了解决这些分布式系统挑战而诞生的。它通过提供一个共享的、分层命名空间(类似于文件系统),并提供数据访问、监听机制等,使得分布式应用能够构建在其之上,实现各种协调功能。
总结 Zookeeper 的关键特性:
分布式协调: 核心功能,解决分布式系统中的一致性、同步等问题。
数据模型: 基于类似文件系统的树形结构存储数据,称为 ZNode。
高可用性: 通过集群部署(Ensemble)保证服务的可用性,即使部分节点故障,服务依然可用。
数据一致性: 保证客户端读到的数据在一定程度上是一致的(最终一致性,通过 Zab 协议保证)。
顺序访问: 客户端的所有操作都按照顺序执行,保证操作的有序性。
实时性: 通过 Watch 机制,客户端可以实时监控数据的变化。
Zookeeper 集群通常由一组服务器组成,称为 Ensemble。Ensemble 中的服务器角色主要分为三种:Leader、Follower 和 Observer。
组件详解:
Ensemble (集群): Zookeeper 的核心组成部分,由多个 Zookeeper 服务器组成,通常为奇数个(推荐 3、5、7 等)。Ensemble 保证了 Zookeeper 的高可用性和数据一致性。
Leader (领导者): Ensemble 中选举产生一个 Leader 服务器,负责处理客户端的写请求,并协调 Follower 服务器进行数据同步。在一个 Zookeeper 集群中,同一时刻只有一个 Leader。
Follower (跟随者): Follower 服务器接收 Leader 服务器的数据同步,并处理客户端的读请求。当 Leader 服务器宕机时,Follower 服务器会参与新的 Leader 选举。
Observer (观察者): Observer 服务器与 Follower 服务器类似,也接收 Leader 服务器的数据同步并处理客户端的读请求。但 Observer 不参与 Leader 选举,也不参与写操作的投票过程,因此可以提升集群的读性能,但牺牲了一定的写性能。
Client (客户端): 应用程序通过 Zookeeper 客户端库连接到 Zookeeper Ensemble,进行数据操作和监听。
工作流程简述:
客户端连接: 客户端连接到 Ensemble 中的任意一台 Zookeeper 服务器。
读请求: 读请求可以直接由连接的服务器(Leader、Follower 或 Observer)处理。
写请求: 写请求会被转发给 Leader 服务器处理。
数据同步: Leader 服务器将写操作同步给 Follower 和 Observer 服务器。
数据持久化: Zookeeper 服务器将数据持久化到磁盘,保证数据可靠性。
Watch 机制: 客户端可以注册 Watcher 监听 ZNode 的变化,当 ZNode 数据或状态发生改变时,Zookeeper 服务器会通知客户端。
Zab 协议 (Zookeeper Atomic Broadcast):
Zookeeper 使用 Zab 协议来保证数据在 Ensemble 中的一致性。Zab 协议是一种基于 Paxos 算法的改进协议,专门用于 Zookeeper 的数据同步和 Leader 选举。Zab 协议保证了以下特性:
原子广播: Leader 发起的事务(写操作)会被原子地广播到所有 Follower 服务器,要么全部成功,要么全部失败。
顺序广播: 事务按照 Leader 发起的顺序广播到 Follower 服务器,保证事务的顺序性。
单一领导者: 在任何时候,集群中只有一个 Leader 服务器负责处理写请求。
Zookeeper 的数据模型是一个层次化的命名空间,类似于标准的文件系统。每个节点称为 ZNode (Zookeeper Node)。ZNode 可以存储少量的数据(默认最大 1MB),并且可以拥有子节点,形成一个树形结构。
ZNode 的类型:
ZNode 主要分为两大类:持久节点 (Persistent) 和 临时节点 (Ephemeral)。每种类型又可以细分为 顺序节点 (Sequential) 和 非顺序节点。
持久节点 (Persistent):
持久节点: 一旦创建,除非显式删除,否则一直存在。即使创建该节点的客户端会话结束,节点依然存在。
持久顺序节点: 在持久节点的基础上,Zookeeper 会自动为节点名称追加一个单调递增的序号。例如,创建节点 /locks/lock-,Zookeeper 可能会将其重命名为 /locks/lock-0000000001、/locks/lock-0000000002 等。顺序节点常用于实现分布式锁、队列等场景。
临时节点 (Ephemeral):
临时节点: 生命周期与客户端会话绑定。当创建该节点的客户端会话结束(例如客户端断开连接),Zookeeper 会自动删除该临时节点。
临时顺序节点: 在临时节点的基础上,Zookeeper 也会自动为节点名称追加一个单调递增的序号。
ZNode 的状态信息 (Stat):
每个 ZNode 除了存储数据外,还维护着一些元数据信息,称为 Stat。Stat 包含了 ZNode 的版本号、创建时间、修改时间、子节点数量等信息。客户端可以通过 getData、exists 等操作获取 ZNode 的 Stat 信息。
常用的 Stat 属性:
czxid (Creation Transaction ID): ZNode 创建时的事务 ID。
mzxid (Modification Transaction ID): ZNode 最后一次修改时的事务 ID。
pzxid (Parent Transaction ID): 子节点列表最后一次修改时的事务 ID。
ctime (Creation Time): ZNode 创建时间。
mtime (Modification Time): ZNode 最后一次修改时间。
version (Data Version): ZNode 数据版本号,每次数据更新版本号都会递增。
cversion (Children Version): 子节点版本号,子节点列表每次修改版本号都会递增。
aversion (ACL Version): ACL 版本号,ACL 列表每次修改版本号都会递增。
ephemeralOwner (Ephemeral Owner): 如果 ZNode 是临时节点,则该属性表示创建该节点的会话 ID。如果是持久节点,则为 0。
dataLength (Data Length): ZNode 存储的数据长度。
numChildren (Number of Children): ZNode 的子节点数量。
客户端要与 Zookeeper 服务器进行交互,首先需要建立一个 会话 (Session)。会话是客户端与 Zookeeper 集群之间的连接,客户端的所有操作都必须在会话的上下文中进行。
会话的生命周期:
会话建立: 客户端通过 Zookeeper 客户端库连接到 Zookeeper 服务器,建立会话。
会话维持: 客户端和服务器之间通过心跳机制保持会话的活跃状态。客户端会定期向服务器发送心跳包,服务器收到心跳包后会更新会话的超时时间。
会话超时: 如果在会话超时时间内,服务器没有收到客户端的心跳包,则认为会话失效。
会话结束: 会话失效或客户端主动断开连接,会话结束。会话结束后,与该会话相关的临时节点会被自动删除。
会话超时时间 (Session Timeout):
会话超时时间是客户端在创建会话时指定的,表示客户端与服务器之间允许的最大失联时间。如果在超时时间内,客户端没有发送心跳包,服务器就会认为会话失效。会话超时时间需要在服务器配置的 minSessionTimeout 和 maxSessionTimeout 之间。
会话 ID (Session ID):
每个会话都有一个唯一的会话 ID,由 Zookeeper 服务器分配。会话 ID 用于标识客户端会话,以及关联临时节点。
Watch (监听) 是 Zookeeper 提供的一种重要的事件通知机制。客户端可以注册 Watcher 监听某个 ZNode 的数据变化或子节点变化。当被监听的 ZNode 发生变化时,Zookeeper 服务器会通知客户端,客户端可以根据通知进行相应的处理。
Watch 的特性:
一次性触发: Watch 是一次性触发的。一旦 Watch 被触发,它就会被移除。如果客户端需要持续监听,需要重新注册 Watcher。
异步通知: Watch 通知是异步的,客户端收到通知时,ZNode 的变化可能已经发生了一段时间。
顺序一致性: Watch 通知保证顺序一致性,即客户端收到的 Watch 通知顺序与 ZNode 变化的顺序一致。
轻量级: Watch 机制是轻量级的,不会对 Zookeeper 服务器造成过大的性能压力。
Watch 的类型:
客户端可以监听以下类型的事件:
NodeDataChanged: 监听 ZNode 数据内容的变化。
NodeChildrenChanged: 监听 ZNode 子节点列表的变化。
NodeCreated: 监听 ZNode 的创建事件(通常用于监听父节点下新创建的子节点)。
NodeDeleted: 监听 ZNode 的删除事件。
Watch 工作流程:
客户端注册 Watcher: 客户端在读取 ZNode 数据或获取子节点列表时,可以同时注册 Watcher。
服务器存储 Watcher: Zookeeper 服务器将 Watcher 信息存储在被监听的 ZNode 关联的 Watcher 列表中。
ZNode 发生变化: 当被监听的 ZNode 数据或子节点列表发生变化时,服务器会遍历 Watcher 列表,并向注册了 Watcher 的客户端发送通知。
客户端接收通知: 客户端收到 Watch 通知后,可以执行相应的处理逻辑,例如重新获取最新的数据或子节点列表。
代码实践:Zookeeper 基础操作
以下 Java 代码示例演示了 Zookeeper 的基本操作,包括连接 Zookeeper 集群、创建 ZNode、获取 ZNode 数据、设置 Watcher 和删除 ZNode。
import org.apache.zookeeper.*; import org.apache.zookeeper.data.Stat; import java.io.IOException; import java.util.List; import java.util.concurrent.CountDownLatch; public class ZookeeperBasicOperations { private static final String CONNECT_STRING = "127.0.0.1:2181"; // Zookeeper 集群地址 private static final int SESSION_TIMEOUT = 5000; // 会话超时时间 private static CountDownLatch connectedSignal = new CountDownLatch(1); public static void main(String[] args) throws IOException, InterruptedException, KeeperException { ZooKeeper zooKeeper = connectZookeeper(); // 1. 创建持久节点 createPersistentNode(zooKeeper, "/my_persistent_node", "Persistent Node Data"); // 2. 创建临时节点 createEphemeralNode(zooKeeper, "/my_ephemeral_node", "Ephemeral Node Data"); // 3. 获取节点数据并设置 Watcher getDataAndWatch(zooKeeper, "/my_persistent_node"); // 4. 修改节点数据 setData(zooKeeper, "/my_persistent_node", "Updated Persistent Node Data"); // 5. 获取子节点列表 getChildren(zooKeeper, "/"); // 6. 删除节点 deleteNode(zooKeeper, "/my_ephemeral_node"); // 只能删除叶子节点,需要先删除子节点才能删除父节点 zooKeeper.close(); } // 连接 Zookeeper public static ZooKeeper connectZookeeper() throws IOException, InterruptedException { ZooKeeper zooKeeper = 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("Connected to Zookeeper!"); } } }); connectedSignal.await(); // 等待连接成功 return zooKeeper; } // 创建持久节点 public static void createPersistentNode(ZooKeeper zooKeeper, String path, String data) throws KeeperException, InterruptedException { zooKeeper.create(path, data.getBytes(), ZooDefs.Ids.OPEN_ACL_UNSAFE, CreateMode.PERSISTENT); System.out.println("Created persistent node: " + path); } // 创建临时节点 public static void createEphemeralNode(ZooKeeper zooKeeper, String path, String data) throws KeeperException, InterruptedException { zooKeeper.create(path, data.getBytes(), ZooDefs.Ids.OPEN_ACL_UNSAFE, CreateMode.EPHEMERAL); System.out.println("Created ephemeral node: " + path); } // 获取节点数据并设置 Watcher public static void getDataAndWatch(ZooKeeper zooKeeper, String path) throws KeeperException, InterruptedException { Stat stat = new Stat(); byte[] data = zooKeeper.getData(path, new Watcher() { @Override public void process(WatchedEvent event) { if (event.getType() == Event.EventType.NodeDataChanged) { try { getDataAndWatch(zooKeeper, path); // 循环注册 Watcher,实现持续监听 System.out.println("Node data changed, new data: " + new String(zooKeeper.getData(path, false, null))); } catch (KeeperException e) { e.printStackTrace(); } catch (InterruptedException e) { e.printStackTrace(); } } } }, stat); System.out.println("Data of node " + path + ": " + new String(data) + ", version: " + stat.getVersion()); } // 修改节点数据 public static void setData(ZooKeeper zooKeeper, String path, String data) throws KeeperException, InterruptedException { Stat stat = zooKeeper.setData(path, data.getBytes(), -1); // version -1 表示忽略版本检查,强制更新 System.out.println("Set data for node " + path + ", new version: " + stat.getVersion()); } // 获取子节点列表 public static void getChildren(ZooKeeper zooKeeper, String path) throws KeeperException, InterruptedException { List<String> children = zooKeeper.getChildren(path, false); System.out.println("Children of node " + path + ": " + children); } // 删除节点 public static void deleteNode(ZooKeeper zooKeeper, String path) throws KeeperException, InterruptedException { zooKeeper.delete(path, -1); // version -1 表示忽略版本检查,强制删除 System.out.println("Deleted node: " + path); } }
代码详解:
connectZookeeper() 方法:
创建 ZooKeeper 对象,传入 Zookeeper 集群地址 CONNECT_STRING、会话超时时间 SESSION_TIMEOUT 和一个 Watcher 对象。
Watcher 对象用于监听连接状态事件。当连接状态变为 SyncConnected 时,表示连接成功,CountDownLatch 计数器减一,解除 await() 方法的阻塞。
connectedSignal.await() 用于阻塞主线程,直到连接成功。
返回 ZooKeeper 对象,用于后续操作。
createPersistentNode() 方法:
使用 zooKeeper.create() 方法创建持久节点。
参数解释:
path: 节点路径。
data.getBytes(): 节点数据,需要转换为字节数组。
ZooDefs.Ids.OPEN_ACL_UNSAFE: ACL 权限设置,OPEN_ACL_UNSAFE 表示允许所有客户端访问,生产环境需要根据实际情况设置更安全的 ACL。
CreateMode.PERSISTENT: 创建模式为持久节点。
createEphemeralNode() 方法:
createPersistentNode() 方法类似,只是 CreateMode 参数设置为 CreateMode.EPHEMERAL,创建临时节点。getDataAndWatch() 方法:
使用 zooKeeper.getData() 方法获取节点数据,并同时注册 Watcher 监听数据变化事件。
参数解释:
path: 节点路径。
Watcher 对象:匿名内部类实现 Watcher 接口,重写 process() 方法处理 Watch 事件。
在 process() 方法中,判断事件类型是否为 Event.EventType.NodeDataChanged,如果是,则表示节点数据发生变化。
在事件处理逻辑中,再次调用 getDataAndWatch() 方法,重新注册 Watcher,实现持续监听。
获取最新的节点数据并打印。
stat: 用于接收节点 Stat 信息的 Stat 对象。
zooKeeper.getData() 方法返回节点数据字节数组,需要转换为 String 类型进行打印。
setData() 方法:
使用 zooKeeper.setData() 方法修改节点数据。
参数解释:
path: 节点路径。
data.getBytes(): 新的节点数据。
-1: 版本号,-1 表示忽略版本检查,强制更新。实际应用中,可以使用 getData() 返回的 Stat 对象中的 version 属性进行乐观锁控制,避免并发修改冲突。
getChildren() 方法:
使用 zooKeeper.getChildren() 方法获取子节点列表。
参数解释:
path: 父节点路径。
false: 是否设置 Watcher 监听子节点列表变化,这里设置为 false 表示不监听。
deleteNode() 方法:
使用 zooKeeper.delete() 方法删除节点。
参数解释:
path: 节点路径。
-1: 版本号,-1 表示忽略版本检查,强制删除。
运行代码前准备:
安装 Zookeeper: 确保本地或远程部署了 Zookeeper 集群。
引入 Zookeeper 客户端库: 在 Maven 或 Gradle 项目中引入 Zookeeper 客户端依赖(例如 org.apache.zookeeper:zookeeper:3.6.3)。
修改 CONNECT_STRING: 将 CONNECT_STRING 修改为实际的 Zookeeper 集群地址。
运行代码后预期结果:
控制台会输出 Zookeeper 连接成功信息,以及节点创建、数据获取、数据修改、子节点列表获取、节点删除等操作的日志信息。如果修改了 /my_persistent_node 节点的数据,会触发 Watcher,并输出节点数据变化的通知信息。