8.1 Leader 选举源码分析


文档摘要

8.1 Leader 选举源码分析 8.1 Leader 选举源码分析 8.1.1 引言 在分布式系统领域中,一致性是一个至关重要的概念。ZooKeeper 作为一款高性能的分布式协调服务,其核心功能之一便是保证集群内数据的一致性。而 Leader 选举机制正是 ZooKeeper 实现数据一致性的基石。在 ZooKeeper 集群中,必须选举出一个 Leader 服务器来负责处理客户端的写请求,并协调 Follower 服务器进行数据同步,从而保证整个集群的数据一致性。 8.1.2 Leader 选举概述 ZooKeeper 使用基于 Fast Paxos 算法的 Leader 选举机制,也被称为 Fast Leader Election (FLE)。

8.1 Leader 选举源码分析

8.1 Leader 选举源码分析

8.1.1 引言

在分布式系统领域中,一致性是一个至关重要的概念。ZooKeeper 作为一款高性能的分布式协调服务,其核心功能之一便是保证集群内数据的一致性。而 Leader 选举机制正是 ZooKeeper 实现数据一致性的基石。在 ZooKeeper 集群中,必须选举出一个 Leader 服务器来负责处理客户端的写请求,并协调 Follower 服务器进行数据同步,从而保证整个集群的数据一致性。

8.1.2 Leader 选举概述

ZooKeeper 使用基于 Fast Paxos 算法的 Leader 选举机制,也被称为 Fast Leader Election (FLE)。当 ZooKeeper 集群启动或 Leader 服务器发生故障时,集群中的服务器会进入 Leader 选举过程,最终选举出一个 Leader 服务器来对外提供服务。

Leader 选举的目标:

  • 选出一个唯一的 Leader: 在任何时刻,集群中只能有一个 Leader 服务器。

  • 确保 Leader 的数据是最新的: 选举出的 Leader 必须拥有集群中最新的数据。

  • 快速完成选举: 在集群启动或故障恢复时,尽快完成 Leader 选举,减少服务不可用时间。

  • 容错性: 即使部分服务器发生故障,集群也能正常进行 Leader 选举。

Leader 选举的基本流程:

  1. 服务器启动或检测到 Leader 丢失: 当 ZooKeeper 服务器启动或检测到当前集群没有 Leader 时,会进入 Leader 选举状态。

  2. 发起投票: 每个服务器都会将自己的投票信息广播给集群中的其他服务器。投票信息通常包含服务器的 ID、事务 ID (zxid) 等。

  3. 接收投票并更新自身投票: 服务器接收到其他服务器的投票后,会根据一定的规则更新自己的投票。

  4. 统计投票结果: 每个服务器都会统计收到的投票,判断是否已经选举出 Leader。

  5. 确定 Leader 或进入下一轮选举: 如果某个服务器收到了超过半数服务器的相同投票,则该投票对应的服务器成为 Leader。否则,进入下一轮选举。

8.1.3 Leader 选举源码详解

ZooKeeper Leader 选举的核心代码主要位于 org.apache.zookeeper.server.quorum 包下,关键类包括:

  • QuorumPeer: 代表一个 ZooKeeper 服务器节点,负责服务器的启动、运行、状态维护以及参与 Leader 选举。

  • FastLeaderElection: 实现 Fast Leader Election 算法的核心类,负责选举的具体逻辑。

  • Vote: 封装投票信息的类,包含服务器 ID、事务 ID 等信息。

  • ToSend: 用于封装待发送的消息,包括投票信息等。

  • Recv: 用于封装接收到的消息。

下面我们将结合源码,逐步分析 Leader 选举的流程。

8.1.3.1 选举入口:QuorumPeer.run()QuorumPeer.startLeaderElection()

ZooKeeper 服务器启动后,QuorumPeer.run() 方法会被调用,该方法是服务器主循环。在 run() 方法中,服务器会根据自身状态判断是否需要进行 Leader 选举。如果服务器处于 LOOKING 状态 (表示正在进行 Leader 选举),则会调用 startLeaderElection() 方法启动选举过程。

// QuorumPeer.java public void run() { // ... try { while (running) { switch (getPeerState()) { case LOOKING: if (tick.get() > 0) { tick.set(0); } setPeerState(startLeaderElection()); // 启动 Leader 选举 break; // ... 其他状态处理 } } } finally { // ... } } protected PeerState startLeaderElection() throws InterruptedException, KeeperException { try { electionAlg = createElectionAlgorithm(electionType); // 创建选举算法实例 } catch (UnknownElectionTypeException e) { // ... } electionAlg.shutdown(); return electionAlg.lookForLeader(); // 执行选举并返回选举结果状态 }

startLeaderElection() 方法首先会根据配置的选举类型 (electionType) 创建对应的选举算法实例,默认使用 FastLeaderElection。然后调用 electionAlg.lookForLeader() 方法启动选举过程。

8.1.3.2 选举算法核心:FastLeaderElection.lookForLeader()

FastLeaderElection.lookForLeader() 方法是 Leader 选举的核心方法,负责实现 Fast Leader Election 算法的逻辑。

// FastLeaderElection.java public Vote lookForLeader() throws InterruptedException { try { self.set选举状态(ServerState.LOOKING); // 设置自身状态为 LOOKING long start_time = System.currentTimeMillis(); final HashMap<Long, Vote> recvset = new HashMap<Long, Vote>(); // 接收到的投票集合 final HashMap<Long, Vote> outofelection = new HashMap<Long, Vote>(); // 已经完成选举的服务器集合 int electionRound = 0; while ((self.get选举状态() == ServerState.LOOKING)) { electionRound++; recvset.clear(); // 清空接收到的投票集合 // 1. 发送投票 sendInitial投票(); while ((self.get选举状态() == ServerState.LOOKING)) { // 2. 接收投票 Recv 工作单元 = recvQueue.poll(pollTimeout, TimeUnit.MILLISECONDS); if (工作单元 == null) { // 超时,重新发送投票 sendInitial投票(); continue; } // 3. 处理接收到的投票 Vote v = 工作单元.vote; if (v != null) { if (isValidVote(v)) { recvset.put(v.getId(), v); // 将投票添加到接收集合 if (should处理(v)) { switch (v.getState()) { case LOOKING: if (处理投票(v, recvset, outofelection)) { // 处理投票,判断是否选出 Leader return 决定选举结果(); // 决定选举结果并返回 } break; case OBSERVING: case FOLLOWING: case LEADING: // ... 处理已完成选举的服务器投票 处理已完成选举服务器投票(v, outofelection); if (outofelection.size() > self.getQuorumSize() / 2) { // 超过半数服务器已完成选举 return 决定选举结果(); } break; default: unexpected状态(v); } } } else { // ... 处理无效投票 } } else { // ... 处理空投票 } } } return null; // Should never happen } finally { // ... 清理资源 } }

lookForLeader() 方法的核心逻辑可以概括为以下几个步骤,并结合 Mermaid 图进行展示:

1. 发送初始投票 (sendInitial投票()):

每个服务器启动选举时,都会向集群中的其他服务器广播自己的初始投票。初始投票信息通常包含:

  • id: 服务器的唯一标识符 (serverId)。

  • zxid: 服务器当前最大的事务 ID (zxid),代表服务器数据的最新程度。

  • electionEpoch: 选举轮次,用于区分不同轮次的选举。

  • state: 服务器当前状态,初始状态为 LOOKING

初始投票的生成逻辑在 FastLeaderElection.总投票() 方法中:

// FastLeaderElection.java protected Vote 总投票() { return new Vote(self.getId(), self.getLastLoggedZxid(), self.getElectionEpoch()); } protected void sendInitial投票() { Vote v = 总投票(); // 生成初始投票 send广播(v); // 广播投票信息 }

send广播(v) 方法负责将投票信息广播给集群中的其他服务器,具体的网络通信实现由 QuorumPeer 类负责。

2. 接收投票 (recvQueue.poll()):

每个服务器都会维护一个接收队列 recvQueue,用于接收来自其他服务器的投票信息。lookForLeader() 方法使用 recvQueue.poll() 方法从队列中获取接收到的投票,并设置超时时间 pollTimeout,避免长时间阻塞。

3. 处理投票 (处理投票(v, recvset, outofelection)):

处理投票(v, recvset, outofelection) 方法是选举算法的核心逻辑,负责处理接收到的投票,并根据投票信息更新自身状态和判断是否选出 Leader。

// FastLeaderElection.java protected boolean 处理投票(Vote v, HashMap<Long, Vote> recvset, HashMap<Long, Vote> outofelection) { boolean 返回值 = false; // 投票比较规则:优先比较 zxid,zxid 大的优先;zxid 相同,serverId 大的优先 Vote currentVote = 总投票(); // 获取自身当前投票 if (v.zxid > currentVote.zxid || (v.zxid == currentVote.zxid && v.getId() > currentVote.getId())) { if (LOG.isDebugEnabled()) { LOG.debug("New vote received " + v + " from server " + v.getId() + " (n=" + self.getId() + ") , New 投票: " + v); } set总投票(v); // 更新自身投票为收到的更优投票 } else { if (LOG.isDebugEnabled()) { LOG.debug("Current vote is better: " + currentVote + " , New vote received " + v + " from server " + v.getId() + " (n=" + self.getId() + ")"); } } if (recvset.containsKey(self.getId())) { recvset.remove(self.getId()); // 移除自身投票,避免重复计数 } recvset.put(v.getId(), v); // 将接收到的投票添加到接收集合 // 统计投票结果 if (quorum认证(recvset)) { // 判断是否达到 Quorum 数量 返回值 = true; try { setLeader(决定Leader(recvset)); // 决定 Leader } catch (Exception e) { // ... } set选举状态(ServerState.FOLLOWING); // 设置自身状态为 FOLLOWING (如果自身不是 Leader) } return 返回值; }

处理投票() 方法的核心逻辑如下:

  • 投票比较: 将接收到的投票 v 与自身当前的投票进行比较。比较规则是 优先比较 zxidzxid 大的投票更优;如果 zxid 相同,则比较 serverIdserverId 大的投票更优。如果接收到的投票更优,则更新自身的投票 (set总投票(v))。

  • 统计投票: 将接收到的投票添加到接收集合 recvset 中,并调用 quorum认证(recvset) 方法判断是否达到 Quorum 数量。

  • Quorum 认证 (quorum认证(recvset)): 判断接收到的投票集合中,是否有超过半数服务器投票给同一个 Leader。

// FastLeaderElection.java protected boolean quorum认证(HashMap<Long, Vote> 投票集合) { if (投票集合.size() < self.getQuorumSize()) { // 投票数量不足 Quorum return false; } HashMap<Vote, Integer> countmap = new HashMap<Vote, Integer>(); for (Vote v : 投票集合.values()) { Integer count = countmap.get(v); if (count == null) { count = 1; } else { count++; } countmap.put(v, count); } for (Entry<Vote, Integer> entry : countmap.entrySet()) { if (entry.getValue() > self.getQuorumSize() / 2) { // 超过半数服务器投票给同一个 Leader return true; } } return false; }

quorum认证() 方法首先判断投票数量是否达到 Quorum (超过半数服务器)。然后统计每个投票的出现次数,如果某个投票的出现次数超过半数,则认为达到 Quorum,选举成功。

  • 决定 Leader (决定Leader(recvset)): 如果达到 Quorum,则调用 决定Leader(recvset) 方法从接收到的投票中选出 Leader。Leader 的选择规则与投票比较规则相同:优先选择 zxid 最大的投票,如果 zxid 相同,则选择 serverId 最大的投票
// FastLeaderElection.java protected Vote 决定Leader(HashMap<Long, Vote> 投票集合) { Vote 返回值 = null; long maxZxid = -1; long maxId = -1; for (Vote v : 投票集合.values()) { if ((v.zxid > maxZxid) || ((v.zxid == maxZxid) && (v.getId() > maxId))) { 返回值 = v; maxZxid = v.zxid; maxId = v.getId(); } } return 返回值; }
  • 设置状态: 如果当前服务器被选举为 Leader,则设置自身状态为 LEADING;否则,设置自身状态为 FOLLOWING

4. 决定选举结果 (决定选举结果()):

决定选举结果() 方法负责根据自身状态返回选举结果。如果自身被选举为 Leader,则返回 Leader 的投票信息;否则,返回当前 Leader 的投票信息 (从 总投票() 方法获取)。

// FastLeaderElection.java private Vote 决定选举结果() { Vote 返回值 = 总投票(); if (self.getId() == 返回值.getId()) { // 自身被选举为 Leader self.set选举状态(ServerState.LEADING); // 设置自身状态为 LEADING } else { self.set选举状态(ServerState.FOLLOWING); // 设置自身状态为 FOLLOWING } return 返回值; }

8.1.3.3 状态转换

在 Leader 选举过程中,ZooKeeper 服务器会经历不同的状态:

  • LOOKING: 选举状态,表示服务器正在参与 Leader 选举。

  • LEADING: Leader 状态,表示服务器被选举为 Leader。

  • FOLLOWING: Follower 状态,表示服务器是 Follower,跟随 Leader。

  • OBSERVING: Observer 状态,表示服务器是 Observer,只同步数据,不参与投票。

状态转换图如下:

8.1.4 Leader 选举代码实践

为了更好地理解 Leader 选举过程,我们可以通过代码实践来模拟一个简化的 Leader 选举场景。以下是一个使用 Java 模拟的简化版 Leader 选举示例,注意:这只是一个简化的模拟,并非 ZooKeeper 源码的完整实现。

import java.util.*; import java.util.concurrent.*; class Vote { int serverId; long zxid; public Vote(int serverId, long zxid) { this.serverId = serverId; this.zxid = zxid; } public int getServerId() { return serverId; } public long getZxid() { return zxid; } @Override public boolean equals(Object o) { if (this == o) return true; if (o == null || getClass() != o.getClass()) return false; Vote vote = (Vote) o; return serverId == vote.serverId && zxid == vote.zxid; } @Override public int hashCode() { return Objects.hash(serverId, zxid); } @Override public String toString() { return "Vote{" + "serverId=" + serverId + ", zxid=" + zxid + '}'; } } class Server { int serverId; long currentZxid = 0; Map<Integer, BlockingQueue<Vote>> receiveQueues = new HashMap<>(); Map<Integer, Server> serverMap; Vote currentVote; public Server(int serverId, Map<Integer, Server> serverMap) { this.serverId = serverId; this.serverMap = serverMap; this.currentVote = new Vote(serverId, currentZxid); for (int id : serverMap.keySet()) { if (id != serverId) { receiveQueues.put(id, new LinkedBlockingQueue<>()); } } } public int getServerId() { return serverId; } public long getCurrentZxid() { return currentZxid; } public Vote getCurrentVote() { return currentVote; } public void setCurrentVote(Vote vote) { this.currentVote = vote; } public void receiveVote(Vote vote, int senderId) { if (receiveQueues.containsKey(senderId)) { receiveQueues.get(senderId).offer(vote); } } public void sendVote(Vote vote) { for (Server server : serverMap.values()) { if (server.getServerId() != serverId) { server.receiveVote(vote, serverId); } } } public Vote lookForLeader() throws InterruptedException { System.out.println("Server " + serverId + " starting leader election."); sendVote(currentVote); // 发送初始投票 Map<Integer, Vote> receivedVotes = new HashMap<>(); receivedVotes.put(serverId, currentVote); while (true) { for (BlockingQueue<Vote> queue : receiveQueues.values()) { Vote receivedVote = queue.poll(1, TimeUnit.SECONDS); // 接收投票,设置超时 if (receivedVote != null) { System.out.println("Server " + serverId + " received vote: " + receivedVote + " from server " + receivedVote.getServerId()); receivedVotes.put(receivedVote.getServerId(), receivedVote); if (isQuorum(receivedVotes)) { Vote leaderVote = decideLeader(receivedVotes); System.out.println("Server " + serverId + " decided leader: " + leaderVote); return leaderVote; } currentVote = updateCurrentVote(currentVote, receivedVote); // 更新自身投票 sendVote(currentVote); // 重新发送投票 } } } } private Vote updateCurrentVote(Vote currentVote, Vote receivedVote) { if (receivedVote.getZxid() > currentVote.getZxid() || (receivedVote.getZxid() == currentVote.getZxid() && receivedVote.getServerId() > currentVote.getServerId())) { return receivedVote; } return currentVote; } private boolean isQuorum(Map<Integer, Vote> receivedVotes) { if (receivedVotes.size() < (serverMap.size() / 2 + 1)) { return false; } Map<Vote, Integer> voteCounts = new HashMap<>(); for (Vote vote : receivedVotes.values()) { voteCounts.put(vote, voteCounts.getOrDefault(vote, 0) + 1); } for (int count : voteCounts.values()) { if (count >= (serverMap.size() / 2 + 1)) { return true; } } return false; } private Vote decideLeader(Map<Integer, Vote> receivedVotes) { Vote leaderVote = null; long maxZxid = -1; int maxServerId = -1; for (Vote vote : receivedVotes.values()) { if (vote.getZxid() > maxZxid || (vote.getZxid() == maxZxid && vote.getServerId() > maxServerId)) { leaderVote = vote; maxZxid = vote.getZxid(); maxServerId = vote.getServerId(); } } return leaderVote; } } public class SimpleLeaderElection { public static void main(String[] args) throws InterruptedException, ExecutionException { int serverCount = 3; Map<Integer, Server> serverMap = new HashMap<>(); List<Future<Vote>> futures = new ArrayList<>(); ExecutorService executorService = Executors.newFixedThreadPool(serverCount); for (int i = 1; i <= serverCount; i++) { serverMap.put(i, new Server(i, serverMap)); // 注意:serverMap 在循环中被更新 } for (int i = 1; i <= serverCount; i++) { final Server server = serverMap.get(i); futures.add(executorService.submit(() -> server.lookForLeader())); } Vote leader = futures.get(0).get(); // 获取任意一个服务器的选举结果,集群最终会选出同一个 Leader System.out.println("Elected Leader is Server " + leader.getServerId()); executorService.shutdown(); } }

代码说明:

  1. Vote 类: 模拟投票信息,包含 serverIdzxid

  2. Server 类: 模拟 ZooKeeper 服务器,包含服务器 ID、当前 zxid、接收队列、服务器映射表、当前投票等。

  3. lookForLeader() 方法: 模拟 Leader 选举过程,包括发送初始投票、接收投票、更新投票、判断 Quorum、决定 Leader 等逻辑。

  4. isQuorum() 方法: 判断是否达到 Quorum 数量。

  5. decideLeader() 方法: 根据投票信息决定 Leader。

  6. SimpleLeaderElection 类: 主程序,创建多个 Server 实例,并启动选举。

运行代码: 运行 SimpleLeaderElection 类,可以看到模拟的 Leader 选举过程,最终会选举出一个 Leader 服务器。

注意: 此代码仅为简化演示 Leader 选举原理,实际 ZooKeeper 的 Leader 选举实现要复杂得多,包括网络通信、状态管理、异常处理等。

8.1.5 总结

Leader 选举是 ZooKeeper 集群正常运行的基础,理解其源码实现对于深入掌握 ZooKeeper 的工作原理至关重要。通过本章节的学习,读者可以对 ZooKeeper Leader 选举的内部机制有更清晰的认识,为后续深入学习 ZooKeeper 的其他核心功能打下坚实的基础。


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