第八章:Zookeeper 源码分析 (可选) 第八章:Zookeeper 源码分析 (可选) 8.1 Zookeeper 源码环境搭建与概览 在开始源码分析之前,我们需要搭建一个 Zookeeper 源码阅读和编译环境。 8.1.1 源码获取 Zookeeper 源码托管在 Apache 软件基金会的 Git 仓库中。您可以通过以下命令克隆源码: 建议切换到一个稳定的发布版本分支,例如 或者 ,避免 master 分支的频繁变动影响学习。 8.1.2 编译环境准备 Zookeeper 使用 Ant 和 Maven 构建工具进行编译。确保您的系统中安装了 JDK (Java Development Kit) 和 Ant 或 Maven。
第八章:Zookeeper 源码分析 (可选)
8.1 Zookeeper 源码环境搭建与概览
在开始源码分析之前,我们需要搭建一个 Zookeeper 源码阅读和编译环境。
8.1.1 源码获取
Zookeeper 源码托管在 Apache 软件基金会的 Git 仓库中。您可以通过以下命令克隆源码:
git clone https://github.com/apache/zookeeper.git cd zookeeper
建议切换到一个稳定的发布版本分支,例如 release-3.8.0 或者 release-3.7.0,避免 master 分支的频繁变动影响学习。
git checkout release-3.8.0
8.1.2 编译环境准备
Zookeeper 使用 Ant 和 Maven 构建工具进行编译。确保您的系统中安装了 JDK (Java Development Kit) 和 Ant 或 Maven。
JDK: Zookeeper 通常需要 JDK 8 或更高版本。
Ant/Maven: 根据 Zookeeper 版本选择合适的构建工具。较新的版本倾向于使用 Maven。
8.1.3 源码编译
在 Zookeeper 源码根目录下,执行以下命令进行编译 (以 Maven 为例):
mvn clean install -DskipTests
-DskipTests 参数可以跳过测试用例的编译和执行,加快编译速度。编译成功后,会在 zookeeper-server/target/zookeeper-server-*.jar 目录下生成 Zookeeper 服务端 JAR 包。
8.1.4 源码结构概览
Zookeeper 源码结构清晰,主要模块包括:
src/java/main/org/apache/zookeeper: Zookeeper 核心 Java 代码。
client: 客户端 API 相关代码。
server: 服务端核心代码,包括请求处理、数据管理、会话管理、选举等。
data: 数据模型相关,如 ZNode 数据结构。
txn: 事务处理相关。
Watcher: Watcher 机制相关。
jute: Zookeeper 自定义的序列化/反序列化框架。
metrics: 指标监控相关。
proto: 协议定义 (Protocol Buffers)。
auth: 认证授权相关。
common: 通用工具类。
src/java/test/org/apache/zookeeper: 单元测试代码。
src/conf: 配置文件模板。
src/bin: 启动脚本。
src/etc: 示例配置文件。
docs: 文档。
contrib: 贡献代码,一些扩展功能。
8.2 请求处理流程源码分析
Zookeeper 的请求处理流程是其核心工作流程之一。客户端的请求如何被接收、处理、以及响应返回,都由这个流程控制。我们来分析一下请求处理流程的关键源码。
8.2.1 请求接收入口:NIOServerCnxnFactory
Zookeeper 使用 Netty 或 NIO 作为网络通信框架。NIOServerCnxnFactory (如果使用 NIO) 负责监听端口,接收客户端连接,并为每个连接创建 NIOServerCnxn 对象。NIOServerCnxn 代表一个客户端连接,负责接收客户端请求。
8.2.2 请求分发:ZooKeeperServer 和 RequestProcessor 链
NIOServerCnxn 接收到请求后,会将请求交给 ZooKeeperServer 进行处理。ZooKeeperServer 是 Zookeeper 服务端的核心类,负责协调各个模块完成请求处理。
ZooKeeperServer 内部维护了一个 RequestProcessor 链,请求会依次经过这个链上的处理器进行处理。典型的 RequestProcessor 链包括:
PrepRequestProcessor: 预处理请求,例如检查权限、设置请求事务 ID (TxnId) 等。
SyncRequestProcessor: 将事务性请求 (例如 create, set, delete) 写入事务日志 (Transaction Log) 并同步到磁盘。
FinalRequestProcessor: 最终处理请求,将请求应用到内存数据树 (DataTree),并发送响应给客户端。
代码实践:查看 PrepRequestProcessor 的源码
打开 PrepRequestProcessor.java 文件,可以看到 processRequest() 方法,该方法负责预处理请求。例如,以下代码片段展示了 PrepRequestProcessor 如何检查权限:
public void processRequest(Request request) { // ... if (request.type != OpCode.auth && request.type != OpCode.closeSession) { if (!auth скипCheck && !auth скипCheckAny) { //权限检查 if (!ZooKeeperServer.getClientACLManager().checkPermissions( zks, request.authInfo, request.cnxn, request.getPath(), request.getOpCode(), request.getPermTo())) { // ... 权限检查失败处理 ... } } } // ... nextProcessor.processRequest(request); //传递给下一个处理器 }
这段代码首先判断请求类型是否是 auth 或 closeSession,如果是则跳过权限检查。否则,调用 ZooKeeperServer.getClientACLManager().checkPermissions() 方法进行权限检查。如果权限检查失败,则进行相应的错误处理。最后,调用 nextProcessor.processRequest(request) 将请求传递给下一个处理器。
8.2.3 事务日志同步:SyncRequestProcessor
SyncRequestProcessor 负责将事务性请求写入事务日志,保证数据持久性和一致性。Zookeeper 使用事务日志来记录所有对 ZNode 数据的修改操作。
代码实践:查看 SyncRequestProcessor 的源码
打开 SyncRequestProcessor.java 文件,可以看到 run() 方法,该方法是一个后台线程,负责从请求队列中取出请求并处理。以下代码片段展示了 SyncRequestProcessor 如何将请求写入事务日志:
public void run() { try { // ... while (true) { Request si = null; try { si = queuedRequests.take(); //从队列中获取请求 } catch (InterruptedException e) { // ... 中断处理 ... break; } if (si == requestOfDeath) { //退出信号 break; } txnLog.append(HdrTxnList.createTxnList(si)); //写入事务日志 if (zks.isLeader()) { //如果是 Leader zks.getLeader().processTxn(si); // Leader 同步给 Follower } syncCount++; if (syncCount > syncThreshold) { //达到同步阈值 txnLog.sync(); //同步到磁盘 syncCount = 0; } pRequest.processRequest(si); //传递给 FinalRequestProcessor } } finally { // ... 清理工作 ... } }
这段代码在一个循环中不断从 queuedRequests 队列中获取请求。对于每个请求,调用 txnLog.append(HdrTxnList.createTxnList(si)) 将请求写入事务日志。如果当前服务器是 Leader,还会调用 zks.getLeader().processTxn(si) 将事务同步给 Follower。当同步计数器 syncCount 达到阈值 syncThreshold 时,调用 txnLog.sync() 将事务日志同步到磁盘。最后,调用 pRequest.processRequest(si) 将请求传递给 FinalRequestProcessor。
8.2.4 最终处理和响应:FinalRequestProcessor
FinalRequestProcessor 是请求处理链的最后一个环节,负责将请求应用到内存数据树 DataTree,并发送响应给客户端。
代码实践:查看 FinalRequestProcessor 的源码
打开 FinalRequestProcessor.java 文件,可以看到 processRequest() 方法。以下代码片段展示了 FinalRequestProcessor 如何处理 create 请求:
public void processRequest(Request request) { // ... switch (request.getOpCode()) { case OpCode.create: { CreateRequest createRequest = new CreateRequest(); CreateResponse createResponse = new CreateResponse(); try { ByteBufferInputStream.byteBuffer2Record(request.request, createRequest); String path = createRequest.getPath(); byte data[] = createRequest.getData(); List<ACL> acl = createRequest.getAcl(); int flags = createRequest.getFlags(); String newPath = zks.createNode(path, data, acl, flags, request.getSession()); //创建 ZNode createResponse.setPath(newPath); zks.serverStats().incrementRequestsProcessed(); sendResponse(request.cnxn, request.sessionId, request.getXid(), request.getType(), ReturnCode.OK, createResponse); //发送成功响应 } catch (KeeperException e) { // ... 异常处理 ... sendResponse(request.cnxn, request.sessionId, request.getXid(), request.getType(), e.code(), null); //发送失败响应 } break; } // ... 其他 OpCode 处理 ... } }
这段代码根据请求的 OpCode 类型进行不同的处理。对于 OpCode.create 请求,它首先从请求中解析出路径、数据、ACL 和标志等参数,然后调用 zks.createNode() 方法在 DataTree 中创建 ZNode。创建成功后,构造 CreateResponse 对象,并调用 sendResponse() 方法发送成功响应给客户端。如果创建过程中发生异常,则捕获异常并发送失败响应。
8.3 Leader 选举源码分析
Leader 选举是 Zookeeper 实现高可用性和一致性的关键机制。当集群中 Leader 宕机或者启动时,需要进行 Leader 选举,选出一个新的 Leader 来负责处理客户端请求和协调集群状态。
8.3.1 选举算法:Fast Leader Election
Zookeeper 使用 Fast Leader Election 算法进行 Leader 选举。该算法基于 Zab 协议,旨在快速、可靠地选出 Leader。
8.3.2 选举流程概览
Leader 选举流程大致如下:
服务器启动或 Leader 宕机: 集群中的服务器启动时或者当前 Leader 宕机时,触发 Leader 选举。
广播选举信息: 每个服务器广播自己的选举信息,包括服务器 ID、事务 ID (Zxid) 等。
投票和计数: 每个服务器根据收到的选举信息进行投票,并统计投票结果。
选出 Leader: 根据投票结果,选出一个服务器作为 Leader。通常选择 Zxid 最大、服务器 ID 最大的服务器作为 Leader。
状态同步: 新的 Leader 与其他服务器进行状态同步,确保数据一致性。
8.3.3 关键类:FastLeaderElection 和 QuorumPeer
FastLeaderElection: 负责实现 Fast Leader Election 算法的核心逻辑。
QuorumPeer: 代表一个 Zookeeper 服务器节点,参与 Leader 选举,并维护选举状态。
代码实践:查看 FastLeaderElection 的源码
打开 FastLeaderElection.java 文件,可以看到 lookForLeader() 方法,该方法是 Leader 选举的核心方法。以下代码片段展示了 lookForLeader() 方法的主要逻辑:
public Vote lookForLeader() throws InterruptedException { // ... 初始化 ... while ((!this.shutdown) && (getElectionState() != ElectionState.LEADING) && (getElectionState() != ElectionState.OBSERVING) && (getElectionState() != ElectionState.FOLLOWING)) { //选举循环 // ... 发送选举信息 ... sendNotifications(); // ... 接收选举信息 ... recvqueue.clear(); while ((recvqueue.size() == 0) && (!shutdown) && (getElectionState() != ElectionState.LEADING) && (getElectionState() != ElectionState.OBSERVING) && (getElectionState() != ElectionState.FOLLOWING)) { incomingPacket(); //接收选举信息 } // ... 处理接收到的选举信息 ... if ((getElectionState() != ElectionState.LEADING) && (getElectionState() != ElectionState.OBSERVING) && (getElectionState() != ElectionState.FOLLOWING)) { processMessages(recvqueue); //处理接收到的选举信息 } // ... 判断是否选出 Leader ... Vote currentVote = getVote(); if (pingChecker != null) pingChecker.reset(); if ((getElectionState() == ElectionState.LEADING) || (getElectionState() == ElectionState.OBSERVING) || (getElectionState() == ElectionState.FOLLOWING)) { setLastVote(currentVote); break; } } // ... 返回选举结果 ... return self.getVote(); }
这段代码在一个循环中不断进行 Leader 选举。循环主要包括以下步骤:
发送选举信息: 调用 sendNotifications() 方法广播自己的选举信息。
接收选举信息: 调用 incomingPacket() 方法接收其他服务器的选举信息,并将信息放入 recvqueue 队列。
处理接收到的选举信息: 调用 processMessages(recvqueue) 方法处理接收到的选举信息,更新自己的投票和选举状态。
判断是否选出 Leader: 检查选举状态,如果已经选出 Leader (状态变为 LEADING, OBSERVING 或 FOLLOWING),则退出循环。
8.3.4 投票规则
在 FastLeaderElection 算法中,投票规则至关重要。每个服务器会根据收到的选举信息,结合自身的投票,决定投票给哪个服务器。主要的投票规则包括:
优先投票给 Zxid 更大的服务器: Zxid 代表服务器的事务 ID,Zxid 更大的服务器通常拥有更新的数据。
如果 Zxid 相同,则投票给服务器 ID 更大的服务器: 服务器 ID 用于区分服务器,ID 更大的服务器具有更高的优先级。
8.4 数据模型 DataTree 源码分析
DataTree 是 Zookeeper 内存数据树的实现,负责存储和管理 ZNode 数据。所有客户端对 ZNode 的操作最终都会反映到 DataTree 上。
8.4.1 数据结构:树形结构
DataTree 使用树形结构来组织 ZNode。每个 ZNode 由一个路径唯一标识,并存储数据、ACL 信息、以及子节点列表。
8.4.2 关键类:DataTree 和 ZKDatabase
DataTree: 负责内存数据树的构建、维护和操作。
ZKDatabase: 负责 DataTree 的持久化和加载,包括事务日志和快照 (Snapshot)。
代码实践:查看 DataTree 的源码
打开 DataTree.java 文件,可以看到 createNode() 方法,该方法负责在 DataTree 中创建 ZNode。以下代码片段展示了 createNode() 方法的主要逻辑:
public String createNode(final String path, byte data[], List<ACL> acl, int flags, long sessionId, String parentPath, long createTime) throws KeeperException.NoNodeException, KeeperException.NodeExistsException, KeeperException.InvalidACLException, KeeperException.NoAuthException, KeeperException.BadArgumentsException { // ... 权限检查 ... final ZNode znode = createZNode(path, data, acl, flags, createTime); //创建 ZNode 对象 nodes.put(path, znode); //添加到 nodes 映射表 if (!parentPath.equals("/")) { //如果不是根节点 ZNode parent = nodes.get(parentPath); parent.addChild(znode.getName()); //添加到父节点的子节点列表 } // ... Watcher 触发 ... return path; }
这段代码首先进行权限检查,然后调用 createZNode() 方法创建 ZNode 对象。接着,将 ZNode 对象添加到 nodes 映射表,nodes 是一个 ConcurrentHashMap<String, ZNode>,用于存储所有 ZNode,key 为路径,value 为 ZNode 对象。如果创建的不是根节点,还需要将新创建的 ZNode 添加到父节点的子节点列表中。最后,触发相关的 Watcher。
8.5 会话管理 SessionTracker 源码分析
Zookeeper 使用会话 (Session) 来管理客户端连接。客户端连接到 Zookeeper 集群后,会创建一个会话,并在会话有效期内与 Zookeeper 进行交互。
8.5.1 会话生命周期
会话具有生命周期,包括:
创建: 客户端连接到 Zookeeper 服务器时创建会话。
激活: 会话创建后处于激活状态,客户端可以发送请求。
保持活跃: 客户端需要定期发送心跳 (Ping) 包来保持会话活跃。
超时: 如果在会话超时时间内客户端没有发送心跳包,会话将被服务端超时关闭。
关闭: 客户端可以主动关闭会话,或者服务端超时关闭会话。
8.5.2 关键类:SessionTrackerImpl
SessionTrackerImpl 负责管理会话的创建、激活、保持活跃、超时和关闭等操作。
代码实践:查看 SessionTrackerImpl 的源码
打开 SessionTrackerImpl.java 文件,可以看到 startSession() 方法,该方法负责启动一个会话。以下代码片段展示了 startSession() 方法的主要逻辑:
public Session startSession(int sessionTimeout) { long sessionId = nextSessionId.getAndDecrement(); //生成新的 SessionId Session session = new SessionImpl(sessionTimeout, sessionId); //创建 Session 对象 sessionsById.put(sessionId, session); //添加到 sessionsById 映射表 sessionsWithTimeout.add(session); //添加到 sessionsWithTimeout 集合 if (sessionTimeout > 0) { //如果会话超时时间大于 0 sessionExpiryQueue.schedule(session, sessionTimeout, TimeUnit.MILLISECONDS); //加入会话过期队列 } return session; }
这段代码首先使用 nextSessionId.getAndDecrement() 生成一个新的会话 ID。然后,创建一个 SessionImpl 对象,并将其添加到 sessionsById 映射表和 sessionsWithTimeout 集合中。sessionsById 是一个 ConcurrentHashMap<Long, Session>,用于根据 SessionId 查找 Session 对象。sessionsWithTimeout 是一个 HashSet<Session>,用于快速遍历所有需要进行超时检查的 Session。如果会话超时时间大于 0,则将 Session 加入 sessionExpiryQueue 会话过期队列,由 SessionExpiryQueue 负责超时检查和会话关闭。
8.6 Watcher 机制源码分析
Watcher 机制是 Zookeeper 的一个重要特性,允许客户端注册 Watcher 监听 ZNode 的变化,当 ZNode 发生变化时,服务端会通知客户端。
8.6.1 Watcher 类型
Watcher 可以监听以下类型的事件:
NodeCreated: ZNode 被创建。
NodeDeleted: ZNode 被删除。
NodeDataChanged: ZNode 数据被修改。
NodeChildrenChanged: ZNode 子节点列表发生变化。
8.6.2 Watcher 触发流程
Watcher 触发流程大致如下:
客户端注册 Watcher: 客户端在读取 ZNode 数据或获取子节点列表时,可以注册 Watcher。
服务端存储 Watcher: 服务端将 Watcher 信息存储在 DataTree 中。
ZNode 发生变化: 当 ZNode 发生变化 (例如创建、删除、数据修改、子节点变化) 时,服务端会查找注册在该 ZNode 上的 Watcher。
触发 Watcher: 服务端向注册了 Watcher 的客户端发送 Watcher 通知。
客户端处理 Watcher 通知: 客户端接收到 Watcher 通知后,执行相应的处理逻辑。
8.6.3 关键类:WatcherManager 和 DataTree
WatcherManager: 负责管理 Watcher 的注册和触发。
DataTree: 负责存储 Watcher 信息,并在 ZNode 发生变化时通知 WatcherManager。
代码实践:查看 WatcherManager 的源码
打开 WatcherManager.java 文件,可以看到 triggerWatch() 方法,该方法负责触发 Watcher。以下代码片段展示了 triggerWatch() 方法的主要逻辑:
public void triggerWatch(String path, Watcher.Event.EventType type, Set<Watcher> watchers, boolean persistent, int rc) { if (watchers == null || watchers.isEmpty()) { return; //没有 Watcher 注册 } for (Watcher w : watchers) { //遍历 Watcher 集合 WatchedEvent event = new WatchedEvent(type, KeeperState.SyncConnected, path); //创建 WatchedEvent 对象 if (persistent) { //如果是持久 Watcher event.setWrapper(new RecoverableWatcher(w)); //使用 RecoverableWatcher 包装 } else { removeWatcher(path, w); //移除临时 Watcher } w.process(event); //调用 Watcher 的 process() 方法处理事件 } }
这段代码首先判断是否有 Watcher 注册在该路径上,如果没有则直接返回。否则,遍历 Watcher 集合,为每个 Watcher 创建一个 WatchedEvent 对象,表示 Watcher 事件。如果是持久 Watcher,使用 RecoverableWatcher 进行包装,以便在连接断开重连后能够恢复 Watcher。如果是临时 Watcher,则在触发后移除。最后,调用 Watcher 的 process() 方法处理事件,将 Watcher 通知发送给客户端。
8.7 总结与展望
本章我们对 Zookeeper 源码进行了初步的探索,分析了请求处理流程、Leader 选举、数据模型、会话管理和 Watcher 机制等关键模块的源码。通过源码分析,我们能够更深入地理解 Zookeeper 的内部工作原理,为更好地使用和优化 Zookeeper 打下基础。
Zookeeper 源码博大精深,本章只是冰山一角。希望通过本章的引导,读者能够对 Zookeeper 源码产生兴趣,并继续深入学习和探索。源码分析是一个持续学习的过程,随着对源码理解的深入,您将能够更好地掌握 Zookeeper 的精髓,并在实际应用中发挥其强大的功能。
后续学习方向:
深入研究 Zab 协议: Zab 协议是 Zookeeper 一致性协议的基础,深入理解 Zab 协议对于理解 Zookeeper 的一致性保证至关重要。
分析 Zookeeper 的持久化机制: Zookeeper 使用事务日志和快照来实现数据持久化,深入分析其持久化机制可以更好地理解 Zookeeper 的数据可靠性。
研究 Zookeeper 的性能优化: Zookeeper 在高并发场景下可能面临性能瓶颈,研究 Zookeeper 的性能优化方法,例如调整配置参数、优化代码逻辑等,可以提升 Zookeeper 的性能。
参与 Zookeeper 社区: 参与 Zookeeper 社区,例如贡献代码、提交 Bug 报告、参与讨论等,可以更深入地了解 Zookeeper 的发展动态,并与其他 Zookeeper 开发者交流学习。
希望本章内容能够帮助您开启 Zookeeper 源码分析之旅!