8.4 Session 管理源码分析 请注意,由于 ZooKeeper 源码版本众多,以下分析将基于相对通用的 ZooKeeper 3.x 版本源码,并侧重于核心概念和流程的讲解。实际代码可能因版本而略有差异,但核心原理是相通的。 8.4 Session 管理源码分析 8.4.1 Session 的重要性与概述 在分布式协调系统 ZooKeeper 中,Session (会话) 是一个至关重要的概念。客户端与 ZooKeeper 集群的每一次交互都建立在一个 Session 之上。Session 的有效性和稳定性直接影响到客户端能否可靠地与 ZooKeeper 集群进行通信,并享用 ZooKeeper 提供的各项服务,如数据节点的创建、读取、更新、删除,以及 Watcher 机制的触发等。
请注意,由于 ZooKeeper 源码版本众多,以下分析将基于相对通用的 ZooKeeper 3.x 版本源码,并侧重于核心概念和流程的讲解。实际代码可能因版本而略有差异,但核心原理是相通的。
在分布式协调系统 ZooKeeper 中,Session (会话) 是一个至关重要的概念。客户端与 ZooKeeper 集群的每一次交互都建立在一个 Session 之上。Session 的有效性和稳定性直接影响到客户端能否可靠地与 ZooKeeper 集群进行通信,并享用 ZooKeeper 提供的各项服务,如数据节点的创建、读取、更新、删除,以及 Watcher 机制的触发等。
Session 的核心作用可以概括为:
连接维持: Session 维持了客户端与 ZooKeeper 集群之间的长连接,确保客户端可以持续地进行操作。
状态同步: Session 的存在使得 ZooKeeper 可以跟踪客户端的状态,例如客户端是否在线,是否需要重新连接等。
授权与认证: Session 可以与权限控制机制结合,验证客户端的身份,并根据 Session 的权限进行操作限制。
临时节点管理: ZooKeeper 的临时节点 (Ephemeral Node) 与 Session 紧密关联,当 Session 失效时,与其关联的临时节点会被自动删除。
Session 生命周期:
一个 ZooKeeper Session 的生命周期大致可以分为以下几个阶段:
Session 关键参数:
Session ID (SessionId): ZooKeeper 集群为每个 Session 分配的唯一标识符,用于区分不同的客户端连接。
Session Timeout (SessionTimeout): Session 的超时时间,单位通常为毫秒。客户端在指定时间内未向 ZooKeeper 集群发送心跳包,Session 可能会被判定为超时失效。
Tick Time (TickTime): ZooKeeper 服务器内部的心跳时间间隔,SessionTimeout 通常是 TickTime 的整数倍。
Session 的创建过程主要发生在客户端发起连接请求时。我们从客户端和服务端两个角度来分析 Session 创建的源码流程。
1. 客户端 Session 创建流程 (ClientCnxn.java):
客户端的 ZooKeeper 构造函数会初始化 ClientCnxn 对象,ClientCnxn 负责与 ZooKeeper 服务器建立连接和管理 Session。
public ZooKeeper(String connectString, int sessionTimeout, Watcher watcher) throws IOException { this(connectString, sessionTimeout, watcher, false); } public ZooKeeper(String connectString, int sessionTimeout, Watcher watcher, boolean canBeReadOnly) throws IOException { // ... cnxn = new ClientCnxn(hostProvider, this, sessionTimeout, this, watchManager, defaultWatcher, getDataWatches, existsWatches, getChildrenWatches, canBeReadOnly); cnxn.start(); // 启动 ClientCnxn 连接线程 }
ClientCnxn.start() 方法会启动一个线程,负责建立与服务器的连接。连接建立的关键步骤在 ClientCnxn$SendThread.run() 方法中:
public void run() { // ... while (state.isAlive()) { try { if (!outgoingQueue.isEmpty()) { // ... 发送请求 } if (outgoingQueue.isEmpty()) { // ... if (readSelect(waitTimeout)) { // 使用 NIO 进行读操作 if (sock.getChannel().isOpen()) { lenBuffer.clear(); int rc = sock.getChannel().read(lenBuffer); // 读取数据长度 // ... 读取数据 process(incomingBuffer); // 处理接收到的数据 } } } } catch (InterruptedException e) { // ... } catch (IOException e) { // ... 处理 IO 异常,尝试重连 cleanup(); if (state.isAlive()) { state.set(States.CONNECTING); startConnect(); // 重新发起连接 } break; } } // ... }
在 ClientCnxn$SendThread.run() 方法中,startConnect() 方法负责发起连接请求:
void startConnect() throws IOException { // ... packetToSend(new ConnectRequest()); // 创建 ConnectRequest 请求 state.set(States.CONNECTING); selectorThread.wakeup(); // 唤醒 Selector 线程发送请求 }
ConnectRequest 是客户端发送给服务器的连接请求,包含了客户端的协议版本、最后看到的 ZXID、超时时间等信息。客户端将 ConnectRequest 放入 outgoingQueue 队列,等待 SelectorThread 线程将其发送到服务器。
2. 服务端 Session 创建流程 (NIOServerCnxn.java, SessionTrackerImpl.java):
服务端接收到客户端的 ConnectRequest 后,在 NIOServerCnxn.java 的 process() 方法中进行处理:
public void process(ByteBuffer bb) throws IOException { // ... ByteBufferInputStream bbis = new ByteBufferInputStream(bb); BinaryInputArchive bia = BinaryInputArchive.getArchive(bbis); // ... ConnectRequest connReq = new ConnectRequest(); connReq.deserialize(bia, "connect"); // 反序列化 ConnectRequest // ... ServerCnxn scxn = ServerCnxnFactory.this.createConnection(sock, incomingBuffer, NIOServerCnxn.this); // ... ZooKeeperServer zks = serverFactory.getZooKeeperServer(); long sessionId = zks.getSessionTracker().createSession(connReq.getTimeOut()); // 创建 Session // ... ConnectResponse connRsp = new ConnectResponse(0, zks.lastProcessedZxid(), sessionId, zks.getSessionTracker().getSessionTimeout(sessionId), serverFactory.isReadOnlyMode()); packetToSend(scxn, new ReplyHeader(-1, zxid, OpCode.connectResponse), connRsp); // 发送 ConnectResponse // ... }
服务端接收到 ConnectRequest 后,首先反序列化请求,然后调用 ZooKeeperServer.getSessionTracker().createSession() 方法创建 Session。SessionTracker 接口的默认实现是 SessionTrackerImpl.java。
SessionTrackerImpl.createSession() 方法负责生成 SessionId 并将 Session 信息添加到 Session 管理器中:
public long createSession(int sessionTimeout) { long sessionId = nextSessionId(); // 生成 SessionId sessionsById.put(sessionId, new Session(sessionTimeout)); // 存储 Session 信息 sessionExpiryQueue.add(new SessionSet(sessionTimeout, sessionId)); // 添加到 Session 过期队列 return sessionId; } private long nextSessionId() { long id = 0; while (id == 0) { id = ++nextSessionId; // 递增 SessionId } return id | serverId << 56; // 组合 ServerId 和递增 Id }
SessionTrackerImpl.createSession() 方法的关键步骤包括:
生成 SessionId: nextSessionId() 方法生成唯一的 SessionId,通常会结合服务器 ID 和递增的序列号,以保证全局唯一性。
存储 Session 信息: sessionsById 是一个 ConcurrentHashMap,用于存储 SessionId 和对应的 Session 对象,Session 对象包含了 SessionTimeout 等信息。
添加到过期队列: sessionExpiryQueue 是一个用于管理 Session 过期的队列,新创建的 Session 会被添加到该队列中,等待后续的过期检测。
创建 Session 后,服务端会构建 ConnectResponse 发送给客户端,其中包含了分配的 SessionId、SessionTimeout 等信息。客户端接收到 ConnectResponse 后,Session 创建过程完成,客户端进入 CONNECTED 状态。
Session 的有效性依赖于客户端和服务端之间的心跳机制。客户端定期向服务端发送心跳包 (PingRequest),服务端接收到心跳包后更新 Session 的最后活动时间,从而保持 Session 的有效性。
1. 客户端心跳发送 (ClientCnxn.java):
客户端在 ClientCnxn$SendThread.run() 方法中,会定期检查是否需要发送心跳包:
public void run() { // ... while (state.isAlive()) { try { // ... if (outgoingQueue.isEmpty()) { // ... if (readSelect(waitTimeout)) { // 使用 NIO 进行读操作 // ... } else { // readSelect 超时,可能是需要发送心跳 if (state.getState() == States.CONNECTED) { long now = System.currentTimeMillis(); if (now - lastPingSent >= negotiatedSessionTimeout / 3) { // 超过 SessionTimeout 的 1/3 时间未发送心跳 packetToSend(new PingRequest()); // 发送 PingRequest lastPingSent = now; } } } } } catch (InterruptedException e) { // ... } catch (IOException e) { // ... } } // ... }
客户端会定期检查距离上次发送心跳包的时间,如果超过了 negotiatedSessionTimeout / 3 (SessionTimeout 的三分之一),则会发送 PingRequest 心跳包。negotiatedSessionTimeout 是客户端和服务端协商后的实际 SessionTimeout。
2. 服务端心跳处理 (NIOServerCnxn.java, SessionTrackerImpl.java):
服务端接收到客户端的 PingRequest 后,在 NIOServerCnxn.java 的 process() 方法中进行处理:
public void process(ByteBuffer bb) throws IOException { // ... int op = hdr.getType(); switch (op) { case OpCode.ping: { // 接收到 PingRequest zks.getSessionTracker().touchSession(sessionId, sessionTimeout); // 更新 Session 最后活动时间 ReplyHeader rh = new ReplyHeader(hdr.getXid(), zks.getZxid(), 0); packetToSend(this, rh, pingResponse); // 发送 PingResponse return; } // ... } // ... }
服务端接收到 PingRequest 后,会调用 ZooKeeperServer.getSessionTracker().touchSession() 方法更新 Session 的最后活动时间。
SessionTrackerImpl.touchSession() 方法负责更新 Session 的 lastAccessTime:
public boolean touchSession(long sessionId, int sessionTimeout) { if (sessionsById.containsKey(sessionId)) { Session session = sessionsById.get(sessionId); session.lastAccessTime = System.currentTimeMillis(); // 更新最后访问时间 return true; } return false; }
SessionTrackerImpl.touchSession() 方法会更新指定 Session 的 lastAccessTime 为当前时间,表明该 Session 仍然活跃。服务端还会发送 PingResponse 给客户端,作为心跳回应。
心跳机制流程图:
如果客户端在 SessionTimeout 时间内未发送心跳包,或者网络出现异常导致心跳包丢失,服务端会判定 Session 超时失效。Session 过期处理主要由 SessionTrackerImpl.java 中的过期队列和过期检测线程负责。
1. Session 过期检测 (SessionTrackerImpl.java):
SessionTrackerImpl 内部维护了一个 SessionExpiryQueue,这是一个基于时间的优先级队列,用于存储待过期的 Session。SessionTrackerImpl 还有一个后台线程 ExpiryThread,负责定期从 SessionExpiryQueue 中取出待过期的 Session 并进行处理。
public class SessionTrackerImpl implements SessionTracker, Runnable { // ... private final SessionExpiryQueue sessionExpiryQueue = new SessionExpiryQueue(); private final ExpiryThread expiryThread = new ExpiryThread(); class ExpiryThread extends Thread { // ... public void run() { try { while (running) { long waitTime = sessionExpiryQueue.getWaitTime(); // 获取下次过期检测时间 if (waitTime > 0) { Thread.sleep(waitTime); // 等待 continue; } Set<Long> sessionsToExpire = sessionExpiryQueue.poll(); // 获取待过期 Session 集合 if (sessionsToExpire == null) { continue; } for (long sessionId : sessionsToExpire) { expire(sessionId); // 处理过期 Session } } } catch (InterruptedException e) { // ... } // ... } } // ... }
ExpiryThread.run() 方法的主要逻辑是:
获取下次过期检测时间: sessionExpiryQueue.getWaitTime() 方法返回距离队列中最早过期 Session 的剩余时间。如果返回正数,则表示队列中没有立即需要过期的 Session,线程可以休眠一段时间。
休眠等待: Thread.sleep(waitTime) 使线程休眠,减少 CPU 消耗。
获取待过期 Session: sessionExpiryQueue.poll() 方法从队列中取出所有已超时的 SessionId 集合。
处理过期 Session: expire(sessionId) 方法负责处理过期的 Session。
2. Session 过期处理 (SessionTrackerImpl.java):
SessionTrackerImpl.expire() 方法负责执行 Session 过期后的清理工作:
public void expire(long sessionId) { Session session = sessionsById.remove(sessionId); // 从 sessionsById 中移除 Session if (session != null) { closeSession(sessionId); // 关闭 Session server.getZKDatabase().removeSession(sessionId); // 从 ZKDatabase 中移除 Session 相关数据 } } private void closeSession(long sessionId) { NIOServerCnxn connection = serverCnxnFactory.getConnection(sessionId); if (connection != null) { connection.sendCloseSession(); // 通知客户端 Session 已过期 connection.close(); // 关闭连接 } // ... }
SessionTrackerImpl.expire() 方法的关键步骤包括:
移除 Session 信息: sessionsById.remove(sessionId) 从 sessionsById 移除已过期的 Session 信息。
关闭 Session: closeSession(sessionId) 方法负责关闭 Session,包括:
通知客户端: connection.sendCloseSession() 方法向客户端发送 SessionExpired 通知。
关闭连接: connection.close() 方法关闭与客户端的连接。
清理 ZKDatabase 数据: server.getZKDatabase().removeSession(sessionId) 方法从 ZKDatabase 中移除与该 Session 相关的临时节点等数据。
Session 过期流程图:
当客户端与 ZooKeeper 集群之间的连接断开 (例如网络故障) 时,客户端会尝试重新连接。ZooKeeper 客户端 SDK 通常会自动处理重连逻辑,对应用层透明。
客户端重连流程 (ClientCnxn.java):
客户端的 ClientCnxn$SendThread.run() 方法在捕获到 IOException 等异常时,会触发重连逻辑:
public void run() { // ... while (state.isAlive()) { try { // ... } catch (InterruptedException e) { // ... } catch (IOException e) { // 捕获 IO 异常 // ... 处理 IO 异常,尝试重连 cleanup(); if (state.isAlive()) { state.set(States.CONNECTING); // 设置状态为 CONNECTING startConnect(); // 重新发起连接 } break; } } // ... }
当 ClientCnxn$SendThread.run() 方法捕获到 IOException 时,会执行以下操作:
清理资源: cleanup() 方法清理连接相关的资源,例如关闭 SocketChannel 等。
设置状态为 CONNECTING: state.set(States.CONNECTING) 将客户端状态设置为 CONNECTING,表示正在尝试重新连接。
重新发起连接: startConnect() 方法重新发起连接请求,流程与首次连接类似,发送 ConnectRequest 给服务器。
Session 重连关键点:
SessionId 复用: 客户端在重连时,会尝试复用之前的 SessionId。在 ConnectRequest 中会包含上次会话的 SessionId 和 password。如果服务器仍然维护着该 Session,且 Session 未过期,则可以复用之前的 Session,从而保证 Session 的连续性。
会话迁移: 如果客户端连接到不同的 ZooKeeper 服务器,新的服务器会尝试将会话迁移过来。这涉及到 Leader 服务器的参与,以确保会话状态的一致性。
重连策略: 客户端 SDK 通常会采用一定的重连策略,例如指数退避算法,避免在网络故障时频繁重连导致集群压力过大。
Session 重连流程图:
代码实践:
为了更好地理解 Session 管理的源码,可以尝试以下实践:
断点调试: 在 ZooKeeper 客户端和服务器源码中设置断点,跟踪 Session 创建、心跳发送、过期检测等关键流程,观察变量变化和方法调用。
日志分析: 开启 ZooKeeper 客户端和服务器的 DEBUG 日志,分析日志中关于 Session 管理的信息,例如 SessionId 分配、心跳包记录、Session 过期事件等。
模拟网络故障: 在测试环境中,模拟网络故障 (例如断网、延迟) ,观察客户端的重连行为和 Session 状态变化。
修改 SessionTimeout: 尝试修改 ZooKeeper 的 SessionTimeout 参数,观察 Session 过期时间的变化,以及对客户端行为的影响。
总结:
ZooKeeper 的 Session 管理是其可靠性和稳定性的基石。通过源码分析,我们可以深入理解 Session 的创建、维护、过期和重连机制。掌握 Session 管理的原理,有助于我们更好地理解 ZooKeeper 的工作方式,并在实际应用中进行更有效的配置和优化。
关键要点回顾:
Session 是客户端与 ZooKeeper 集群交互的基础,负责连接维持、状态同步、授权认证和临时节点管理。
Session 生命周期包括 CONNECTING, CONNECTED, RECONNECTING, CLOSED 等状态。
Session 创建涉及客户端和服务端的 ConnectRequest 和 ConnectResponse 交互。
心跳机制通过 PingRequest 和 PingResponse 维持 Session 有效性。
Session 过期由 SessionTrackerImpl 的过期队列和过期检测线程负责。
客户端 SDK 自动处理 Session 重连,尝试复用 SessionId 或重建 Session。