作者:京东物流 梁吉超
zookeeper是一个分布式服务框架,主要解决分布式应用中常见的多种数据问题,例如集群管理,状态同步等。为解决这些问题zookeeper需要Leader选举进行保障数据的强一致性机制和稳定性。本文通过集群的配置,对leader选举源进行解析,让读者们了解如何利用BIO通信机制,多线程多层队列实现高性能架构。****
01Leader选举机制
Leader选举机制采用半数选举算法。
每一个zookeeper服务端称之为一个节点,每
个节点都有投票权,把其选票投向每一个有选举权的节点,当其中一个节点选举出票数过半,这个节点就会成为Leader,其它节点成为Follower。
02Leader选举集群配置
-
重命名zoo_sample.cfg文件为zoo1.cfg ,zoo2.cfg,zoo3.cfg,zoo4.cfg
-
修改zoo.cfg文件,修改值如下:
1【plain】 2zoo1.cfg文件内容: 3dataDir=/export/data/zookeeper-1 4clientPort=2181 5server.1=127.0.0.1:2001:3001 6server.2=127.0.0.1:2002:3002:participant 7server.3=127.0.0.1:2003:3003:participant 8server.4=127.0.0.1:2004:3004:observer 9 10 11zoo2.cfg文件内容: 12dataDir=/export/data/zookeeper-2 13clientPort=2182 14server.1=127.0.0.1:2001:3001 15server.2=127.0.0.1:2002:3002:participant 16server.3=127.0.0.1:2003:3003:participant 17server.4=127.0.0.1:2004:3004:observer 18 19 20zoo3.cfg文件内容: 21dataDir=/export/data/zookeeper-3 22clientPort=2183 23server.1=127.0.0.1:2001:3001 24server.2=127.0.0.1:2002:3002:participant 25server.3=127.0.0.1:2003:3003:participant 26server.4=127.0.0.1:2004:3004:observer 27 28 29zoo4.cfg文件内容: 30dataDir=/export/data/zookeeper-4 31clientPort=2184 32server.1=127.0.0.1:2001:3001 33server.2=127.0.0.1:2002:3002:participant 34server.3=127.0.0.1:2003:3003:participant 35server.4=127.0.0.1:2004:3004:observer 36 37 38 39 40 41
- server.第几号服务器(对应myid文件内容)=ip:数据同步端口:选举端口:选举标识
- participant默认参与选举标识,可不写. observer不参与选举
4.在/export/data/zookeeper-1,/export/data/zookeeper-2,/export/data/zookeeper-3,/export/data/zookeeper-4目录下创建myid文件,文件内容分别写1 ,2,3,4,用于标识sid(全称:Server ID)赋值。
- 启动三个zookeeper实例:
- bin/zkServer.sh start conf/zoo1.cfg
- bin/zkServer.sh start conf/zoo2.cfg
- bin/zkServer.sh start conf/zoo3.cfg
- 每启动一个实例,都会读取启动参数配置zoo.cfg文件,这样实例就可以知道其作为服务端身份信息sid以及集群中有多少个实例参与选举。
03Leader选举流程
图1 第一轮到第二轮投票流程
前提:
设定票据数据格式vote(sid,zxid,epoch)
- sid是Server ID每台服务的唯一标识,是myid文件内容;
- zxid是数据事务id号;
- epoch为选举周期,为方便理解下面讲解内容暂定为1初次选举,不写入下面内容里。
按照顺序启动sid=1,sid=2节点
第一轮投票:
-
sid=1节点:初始选票为自己,将选票vote(1,0)发送给sid=2节点;
-
sid=2节点:初始选票为自己,将选票vote(2,0)发送给sid=1节点;
-
sid=1节点:收到sid=2节点选票vote(2,0)和当前自己的选票vote(1,0),首先比对zxid值,zxid越大代表数据最新,优先选择zxid最大的选票,如果zxid相同,选举最大sid。当前投票选举结果为vote(2,0),sid=1节点的选票变为vote(2,0);
-
sid=2节点:收到sid=1节点选票vote(1,0)和当前自己的选票vote(2,0),参照上述选举方式,选举结果为vote(2,0),sid=2节点的选票不变;
-
第一轮投票选举结束。
第二轮投票:
-
sid=1节点:当前自己的选票为vote(2,0),将选票vote(2,0)发送给sid=2节点;
-
sid=2节点:当前自己的选票为vote(2,0),将选票vote(2,0)发送给sid=1节点;
-
sid=1节点:收到sid=2节点选票vote(2,0)和自己的选票vote(2,0), 按照半数选举算法,总共3个节点参与选举,已有2个节点选举出相同选票,推举sid=2节点为Leader,自己角色变为Follower;
-
sid=2节点:收到sid=1节点选票vote(2,0)和自己的选票vote(2,0),按照半数选举算法推举sid=2节点为Leader,自己角色变为Leader。
这时启动sid=3节点后,集群里已经选举出leader,sid=1和sid=2节点会将自己的leader选票发回给sid=3节点,通过半数选举结果还是sid=2节点为leader。
3.1 Leader选举采用多层队列架构
zookeeper选举底层主要分为选举应用层和消息传输队列层,第一层应用层队列统一接收和发送选票,而第二层传输层队列,是按照服务端sid分成了多个队列,是为了避免给每台服务端发送消息互相影响。比如对某台机器发送不成功不会影响正常服务端的发送。
图2 多层队列上下关系交互流程图
04解析代码入口类
通过查看zkServer.sh文件内容找到服务启动类:
org.apache.zookeeper.server.quorum.QuorumPeerMain
05选举流程代码解析
图3 选举代码实现流程图
- 加载配置文件QuorumPeerConfig.parse(path);
针对 Leader选举关键配置信息如下:
- 读取dataDir目录找到myid文件内容,设置当前应用sid标识,做为投票人身份信息。下面遇到myid变量为当前节点自己sid标识。
-
- 设置peerType当前应用是否参与选举
- new QuorumMaj()解析server.前缀加载集群成员信息,加载allMembers所有成员,votingMembers参与选举成员,observingMembers观察者成员,设置half值votingMembers.size()/2.
1【Java】 2public QuorumMaj(Properties props) throws ConfigException { 3 for (Entry<Object, Object> entry : props.entrySet()) { 4 String key = entry.getKey().toString(); 5 String value = entry.getValue().toString(); 6 //读取集群配置文件中的server.开头的应用实例配置信息 7 if (key.startsWith("server.")) { 8 int dot = key.indexOf('.'); 9 long sid = Long.parseLong(key.substring(dot + 1)); 10 QuorumServer qs = new QuorumServer(sid, value); 11 allMembers.put(Long.valueOf(sid), qs); 12 if (qs.type == LearnerType.PARTICIPANT) 13//应用实例绑定的角色为PARTICIPANT意为参与选举 14 votingMembers.put(Long.valueOf(sid), qs); 15 else { 16 //观察者成员 17 observingMembers.put(Long.valueOf(sid), qs); 18 } 19 } else if (key.equals("version")) { 20 version = Long.parseLong(value, 16); 21 } 22 } 23 //过半基数 24 half = votingMembers.size() / 2; 25 } 26 27 28 29 30 31
-
QuorumPeerMain.runFromConfig(config) 启动服务;
-
QuorumPeer.startLeaderElection() 开启选举服务;
- 设置当前选票new Vote(sid,zxid,epoch)
1【plain】 2synchronized public void startLeaderElection(){ 3try { 4 if (getPeerState() == ServerState.LOOKING) { 5 //首轮:当前节点默认投票对象为自己 6 currentVote = new Vote(myid, getLastLoggedZxid(), getCurrentEpoch()); 7 } 8 } catch(IOException e) { 9 RuntimeException re = new RuntimeException(e.getMessage()); 10 re.setStackTrace(e.getStackTrace()); 11 throw re; 12 } 13//........ 14} 15 16 17 18 19 20
- 创建选举管理类:QuorumCnxnManager;
- 初始化recvQueue<Message(sid,ByteBuffer)>接收投票队列(第二层传输队列);
- 初始化queueSendMap<sid,queue>按sid发送投票队列(第二层传输队列);
- 初始化senderWorkerMap<sid,SendWorker>发送投票工作线程容器,表示着与sid投票节点已连接;
- 初始化选举监听线程类QuorumCnxnManager.Listener。
1【Java】 2//QuorumPeer.createCnxnManager() 3public QuorumCnxManager(QuorumPeer self, 4 final long mySid, 5 Map<Long,QuorumPeer.QuorumServer> view, 6 QuorumAuthServer authServer, 7 QuorumAuthLearner authLearner, 8 int socketTimeout, 9 boolean listenOnAllIPs, 10 int quorumCnxnThreadsSize, 11 boolean quorumSaslAuthEnabled) { 12 //接收投票队列(第二层传输队列) 13 this.recvQueue = new ArrayBlockingQueue<Message>(RECV_CAPACITY); 14 //按sid发送投票队列(第二层传输队列) 15 this.queueSendMap = new ConcurrentHashMap<Long, ArrayBlockingQueue<ByteBuffer>>(); 16 //发送投票工作线程容器,表示着与sid投票节点已连接 17 this.senderWorkerMap = new ConcurrentHashMap<Long, SendWorker>(); 18 this.lastMessageSent = new ConcurrentHashMap<Long, ByteBuffer>(); 19 20 21 String cnxToValue = System.getProperty("zookeeper.cnxTimeout"); 22 if(cnxToValue != null){ 23 this.cnxTO = Integer.parseInt(cnxToValue); 24 } 25 26 27 this.self = self; 28 29 30 this.mySid = mySid; 31 this.socketTimeout = socketTimeout; 32 this.view = view; 33 this.listenOnAllIPs = listenOnAllIPs; 34 35 36 initializeAuth(mySid, authServer, authLearner, quorumCnxnThreadsSize, 37 quorumSaslAuthEnabled); 38 // Starts listener thread that waits for connection requests 39 //创建选举监听线程 接收选举投票请求 40 listener = new Listener(); 41 listener.setName("QuorumPeerListener"); 42} 43//QuorumPeer.createElectionAlgorithm 44protected Election createElectionAlgorithm(int electionAlgorithm){ 45 Election le=null; 46 //TODO: use a factory rather than a switch 47 switch (electionAlgorithm) { 48 case 0: 49 le = new LeaderElection(this); 50 break; 51 case 1: 52 le = new AuthFastLeaderElection(this); 53 break; 54 case 2: 55 le = new AuthFastLeaderElection(this, true); 56 break; 57 case 3: 58 qcm = createCnxnManager();// new QuorumCnxManager(... new Listener()) 59 QuorumCnxManager.Listener listener = qcm.listener; 60 if(listener != null){ 61 listener.start();//启动选举监听线程 62 FastLeaderElection fle = new FastLeaderElection(this, qcm); 63 fle.start(); 64 le = fle; 65 } else { 66 LOG.error("Null listener when initializing cnx manager"); 67 } 68 break; 69 default: 70 assert false; 71 } 72return le;} 73 74 75 76 77 78
- 开启选举监听线程QuorumCnxnManager.Listener;
- 创建ServerSockket等待大于自己sid节点连接,连接信息存储到senderWorkerMap<sid,SendWorker>;
- sid>self.sid才可以连接过来。
1【Java】 2//上面的listener.start()执行后,选择此方法 3public void run() { 4 int numRetries = 0; 5 InetSocketAddress addr; 6 Socket client = null; 7 while((!shutdown) && (numRetries < 3)){ 8 try { 9 ss = new ServerSocket(); 10 ss.setReuseAddress(true); 11 if (self.getQuorumListenOnAllIPs()) { 12 int port = self.getElectionAddress().getPort(); 13 addr = new InetSocketAddress(port); 14 } else { 15 // Resolve hostname for this server in case the 16 // underlying ip address has changed. 17 self.recreateSocketAddresses(self.getId()); 18 addr = self.getElectionAddress(); 19 } 20 LOG.info("My election bind port: " + addr.toString()); 21 setName(addr.toString()); 22 ss.bind(addr); 23 while (!shutdown) { 24 client = ss.accept(); 25 setSockOpts(client); 26 LOG.info("Received connection request " 27 + client.getRemoteSocketAddress()); 28 // Receive and handle the connection request 29 // asynchronously if the quorum sasl authentication is 30 // enabled. This is required because sasl server 31 // authentication process may take few seconds to finish, 32 // this may delay next peer connection requests. 33 if (quorumSaslAuthEnabled) { 34 receiveConnectionAsync(client); 35 } else { 36//接收连接信息 37 receiveConnection(client); 38 } 39 numRetries = 0; 40 } 41 } catch (IOException e) { 42 if (shutdown) { 43 break; 44 } 45 LOG.error("Exception while listening", e); 46 numRetries++; 47 try { 48 ss.close(); 49 Thread.sleep(1000); 50 } catch (IOException ie) { 51 LOG.error("Error closing server socket", ie); 52 } catch (InterruptedException ie) { 53 LOG.error("Interrupted while sleeping. " + 54 "Ignoring exception", ie); 55 } 56 closeSocket(client); 57 } 58 } 59 LOG.info("Leaving listener"); 60 if (!shutdown) { 61 LOG.error("As I'm leaving the listener thread, " 62 + "I won't be able to participate in leader " 63 + "election any longer: " 64 + self.getElectionAddress()); 65 } else if (ss != null) { 66 // Clean up for shutdown. 67 try { 68 ss.close(); 69 } catch (IOException ie) { 70 // Don't log an error for shutdown. 71 LOG.debug("Error closing server socket", ie); 72 } 73 } 74} 75 76 77//代码执行路径:receiveConnection()->handleConnection(...) 78private void handleConnection(Socket sock, DataInputStream din) 79 throws IOException { 80//...省略 81 if (sid < self.getId()) { 82 /* 83 * This replica might still believe that the connection to sid is 84 * up, so we have to shut down the workers before trying to open a 85 * new connection. 86 */ 87 SendWorker sw = senderWorkerMap.get(sid); 88 if (sw != null) { 89 sw.finish(); 90 } 91 92 93 /* 94 * Now we start a new connection 95 */ 96 LOG.debug("Create new connection to server: {}", sid); 97 closeSocket(sock); 98 99 100 if (electionAddr != null) { 101 connectOne(sid, electionAddr); 102 } else { 103 connectOne(sid); 104 } 105 106 107 } else { // Otherwise start worker threads to receive data. 108 SendWorker sw = new SendWorker(sock, sid); 109 RecvWorker rw = new RecvWorker(sock, din, sid, sw); 110 sw.setRecv(rw); 111 112 113 SendWorker vsw = senderWorkerMap.get(sid); 114 115 116 if (vsw != null) { 117 vsw.finish(); 118 } 119 //存储连接信息<sid,SendWorker> 120 senderWorkerMap.put(sid, sw); 121 122 123 queueSendMap.putIfAbsent(sid, 124 new ArrayBlockingQueue<ByteBuffer>(SEND_CAPACITY)); 125 126 127 sw.start(); 128 rw.start(); 129 } 130} 131 132 133 134 135 136
- 创建FastLeaderElection快速选举服务;
- 初始选票发送队列sendqueue(第一层队列)
- 初始选票接收队列recvqueue(第一层队列)
- 创建线程WorkerSender
- 创建线程WorkerReceiver
1【Java】 2//FastLeaderElection.starter 3private void starter(QuorumPeer self, QuorumCnxManager manager) { 4 this.self = self; 5 proposedLeader = -1; 6 proposedZxid = -1; 7 //发送队列sendqueue(第一层队列) 8 sendqueue = new LinkedBlockingQueue<ToSend>(); 9 //接收队列recvqueue(第一层队列) 10 recvqueue = new LinkedBlockingQueue<Notification>(); 11 this.messenger = new Messenger(manager); 12} 13//new Messenger(manager) 14Messenger(QuorumCnxManager manager) { 15 //创建线程WorkerSender 16 this.ws = new WorkerSender(manager); 17 18 19 this.wsThread = new Thread(this.ws, 20 "WorkerSender[myid=" + self.getId() + "]"); 21 this.wsThread.setDaemon(true); 22 //创建线程WorkerReceiver 23 this.wr = new WorkerReceiver(manager); 24 25 26 this.wrThread = new Thread(this.wr, 27 "WorkerReceiver[myid=" + self.getId() + "]"); 28 this.wrThread.setDaemon(true); 29} 30 31 32 33 34 35
- 开启WorkerSender和WorkerReceiver线程。
WorkerSender线程自旋获取sendqueue第一层队列元素
- sendqueue队列元素内容为相关选票信息详见ToSend类;
- 首先判断选票sid是否和自己sid值相同,相等直接放入到recvQueue队列中;
- 不相同将sendqueue队列元素转储到queueSendMap<sid,queue>第二层传输队列中。
1【Java】//FastLeaderElection.Messenger.WorkerSenderclass WorkerSender extends ZooKeeperThread{ 2//... 3 public void run() { 4 while (!stop) { 5 try { 6 ToSend m = sendqueue.poll(3000, TimeUnit.MILLISECONDS); 7 if(m == null) continue; 8 //将投票信息发送出去 9 process(m); 10 } catch (InterruptedException e) { 11 break; 12 } 13 } 14 LOG.info("WorkerSender is down"); 15 } 16} 17//QuorumCnxManager#toSend 18public void toSend(Long sid, ByteBuffer b) { 19 /* 20 * If sending message to myself, then simply enqueue it (loopback). 21 */ 22 if (this.mySid == sid) { 23 b.position(0); 24 addToRecvQueue(new Message(b.duplicate(), sid)); 25 /* 26 * Otherwise send to the corresponding thread to send. 27 */ 28 } else { 29 /* 30 * Start a new connection if doesn't have one already. 31 */ 32 ArrayBlockingQueue<ByteBuffer> bq = new ArrayBlockingQueue<ByteBuffer>( 33 SEND_CAPACITY); 34 ArrayBlockingQueue<ByteBuffer> oldq = queueSendMap.putIfAbsent(sid, bq); 35 //转储到queueSendMap<sid,queue>第二层传输队列中 36 if (oldq != null) { 37 addToSendQueue(oldq, b); 38 } else { 39 addToSendQueue(bq, b); 40 } 41 connectOne(sid); 42 } 43} 44 45 46 47 48 49
WorkerReceiver线程自旋获取recvQueue第二层传输队列元素转存到recvqueue第一层队列中。
1【Java】 2//WorkerReceiver 3public void run() { 4 Message response; 5 while (!stop) { 6 // Sleeps on receive 7 try { 8 //自旋获取recvQueue第二层传输队列元素 9 response = manager.pollRecvQueue(3000, TimeUnit.MILLISECONDS); 10 if(response == null) continue; 11 // The current protocol and two previous generations all send at least 28 bytes 12 if (response.buffer.capacity() < 28) { 13 LOG.error("Got a short response: " + response.buffer.capacity()); 14 continue; 15 } 16 //... 17 if(self.getPeerState() == QuorumPeer.ServerState.LOOKING){ 18 //第二层传输队列元素转存到recvqueue第一层队列中 19 recvqueue.offer(n); 20 //... 21 } 22 } 23//... 24} 25 26 27 28 29 30
06选举核心逻辑
- 启动线程QuorumPeer
开始Leader选举投票makeLEStrategy().lookForLeader();
sendNotifications()向其它节点发送选票信息,选票信息存储到sendqueue队列中。sendqueue队列由WorkerSender线程处理。
1【plain】 2//QuorunPeer.run 3//... 4try { 5 reconfigFlagClear(); 6 if (shuttingDownLE) { 7 shuttingDownLE = false; 8 startLeaderElection(); 9 } 10 //makeLEStrategy().lookForLeader() 发送投票 11 setCurrentVote(makeLEStrategy().lookForLeader()); 12} catch (Exception e) { 13 LOG.warn("Unexpected exception", e); 14 setPeerState(ServerState.LOOKING); 15} 16//... 17//FastLeaderElection.lookLeader 18public Vote lookForLeader() throws InterruptedException { 19//... 20 //向其他应用发送投票 21sendNotifications(); 22//... 23} 24 25 26private void sendNotifications() { 27 //获取应用节点 28 for (long sid : self.getCurrentAndNextConfigVoters()) { 29 QuorumVerifier qv = self.getQuorumVerifier(); 30 ToSend notmsg = new ToSend(ToSend.mType.notification, 31 proposedLeader, 32 proposedZxid, 33 logicalclock.get(), 34 QuorumPeer.ServerState.LOOKING, 35 sid, 36 proposedEpoch, qv.toString().getBytes()); 37 if(LOG.isDebugEnabled()){ 38 LOG.debug("Sending Notification: " + proposedLeader + " (n.leader), 0x" + 39 Long.toHexString(proposedZxid) + " (n.zxid), 0x" + Long.toHexString(logicalclock.get()) + 40 " (n.round), " + sid + " (recipient), " + self.getId() + 41 " (myid), 0x" + Long.toHexString(proposedEpoch) + " (n.peerEpoch)"); 42 } 43 //储存投票信息 44 sendqueue.offer(notmsg); 45 } 46} 47 48 49class WorkerSender extends ZooKeeperThread { 50 //... 51 public void run() { 52 while (!stop) { 53 try { 54//提取已储存的投票信息 55 ToSend m = sendqueue.poll(3000, TimeUnit.MILLISECONDS); 56 if(m == null) continue; 57 58 59 process(m); 60 } catch (InterruptedException e) { 61 break; 62 } 63 } 64 LOG.info("WorkerSender is down"); 65 } 66//... 67} 68 69 70 71 72 73
自旋recvqueue队列元素获取投票过来的选票信息:
1【Java】 2public Vote lookForLeader() throws InterruptedException { 3//... 4/* 5 * Loop in which we exchange notifications until we find a leader 6 */ 7while ((self.getPeerState() == ServerState.LOOKING) && 8 (!stop)){ 9 /* 10 * Remove next notification from queue, times out after 2 times 11 * the termination time 12 */ 13 //提取投递过来的选票信息 14 Notification n = recvqueue.poll(notTimeout, 15 TimeUnit.MILLISECONDS); 16/* 17 * Sends more notifications if haven't received enough. 18 * Otherwise processes new notification. 19 */ 20if(n == null){ 21 if(manager.haveDelivered()){ 22 //已全部连接成功,并且前一轮投票都完成,需要再次发起投票 23 sendNotifications(); 24 } else { 25 //如果未收到选票信息,manager.contentAll()自动连接其它socket节点 26 manager.connectAll(); 27 } 28 /* 29 * Exponential backoff 30 */ 31 int tmpTimeOut = notTimeout*2; 32 notTimeout = (tmpTimeOut < maxNotificationInterval? 33 tmpTimeOut : maxNotificationInterval); 34 LOG.info("Notification time out: " + notTimeout); 35 } 36 //.... 37 } 38 //... 39} 40 41 42 43 44 45
1【Java】 2//manager.connectAll()->connectOne(sid)->initiateConnection(...)->startConnection(...) 3 4 5private boolean startConnection(Socket sock, Long sid) 6 throws IOException { 7 DataOutputStream dout = null; 8 DataInputStream din = null; 9 try { 10 // Use BufferedOutputStream to reduce the number of IP packets. This is 11 // important for x-DC scenarios. 12 BufferedOutputStream buf = new BufferedOutputStream(sock.getOutputStream()); 13 dout = new DataOutputStream(buf); 14 15 16 // Sending id and challenge 17 // represents protocol version (in other words - message type) 18 dout.writeLong(PROTOCOL_VERSION); 19 dout.writeLong(self.getId()); 20 String addr = self.getElectionAddress().getHostString() + ":" + self.getElectionAddress().getPort(); 21 byte[] addr_bytes = addr.getBytes(); 22 dout.writeInt(addr_bytes.length); 23 dout.write(addr_bytes); 24 dout.flush(); 25 26 27 din = new DataInputStream( 28 new BufferedInputStream(sock.getInputStream())); 29 } catch (IOException e) { 30 LOG.warn("Ignoring exception reading or writing challenge: ", e); 31 closeSocket(sock); 32 return false; 33 } 34 35 36 // authenticate learner 37 QuorumPeer.QuorumServer qps = self.getVotingView().get(sid); 38 if (qps != null) { 39 // TODO - investigate why reconfig makes qps null. 40 authLearner.authenticate(sock, qps.hostname); 41 } 42 43 44 // If lost the challenge, then drop the new connection 45 //保证集群中所有节点之间只有一个通道连接 46 if (sid > self.getId()) { 47 LOG.info("Have smaller server identifier, so dropping the " + 48 "connection: (" + sid + ", " + self.getId() + ")"); 49 closeSocket(sock); 50 // Otherwise proceed with the connection 51 } else { 52 SendWorker sw = new SendWorker(sock, sid); 53 RecvWorker rw = new RecvWorker(sock, din, sid, sw); 54 sw.setRecv(rw); 55 56 57 SendWorker vsw = senderWorkerMap.get(sid); 58 59 60 if(vsw != null) 61 vsw.finish(); 62 63 64 senderWorkerMap.put(sid, sw); 65 queueSendMap.putIfAbsent(sid, new ArrayBlockingQueue<ByteBuffer>( 66 SEND_CAPACITY)); 67 68 69 sw.start(); 70 rw.start(); 71 72 73 return true; 74 75 76 } 77 return false; 78} 79 80 81 82 83 84
如上述代码中所示,sid>self.sid才可以创建连接Socket和SendWorker,RecvWorker线程,存储到senderWorkerMap<sid,SendWorker>中。对应第2步中的sid<self.sid逻辑,保证集群中所有节点之间只有一个通道连接。
图4 节点之间连接方式
1【Java】 2 3 4public Vote lookForLeader() throws InterruptedException { 5//... 6 if (n.electionEpoch > logicalclock.get()) { 7 //当前选举周期小于选票周期,重置recvset选票池 8 //大于当前周期更新当前选票信息,再次发送投票 9 logicalclock.set(n.electionEpoch); 10 recvset.clear(); 11 if(totalOrderPredicate(n.leader, n.zxid, n.peerEpoch, 12 getInitId(), getInitLastLoggedZxid(), getPeerEpoch())) { 13 updateProposal(n.leader, n.zxid, n.peerEpoch); 14 } else { 15 updateProposal(getInitId(), 16 getInitLastLoggedZxid(), 17 getPeerEpoch()); 18 } 19 sendNotifications(); 20 } else if (n.electionEpoch < logicalclock.get()) { 21 if(LOG.isDebugEnabled()){ 22 LOG.debug("Notification election epoch is smaller than logicalclock. n.electionEpoch = 0x" 23 + Long.toHexString(n.electionEpoch) 24 + ", logicalclock=0x" + Long.toHexString(logicalclock.get())); 25 } 26 break; 27 } else if (totalOrderPredicate(n.leader, n.zxid, n.peerEpoch, 28 proposedLeader, proposedZxid, proposedEpoch)) {//相同选举周期 29 //接收的选票与当前选票PK成功后,替换当前选票 30 updateProposal(n.leader, n.zxid, n.peerEpoch); 31 sendNotifications(); 32 } 33//... 34 35 36} 37 38 39 40 41 42
在上代码中,自旋从recvqueue队列中获取到选票信息。开始进行选举:
- 判断当前选票和接收过来的选票周期是否一致
- 大于当前周期更新当前选票信息,再次发送投票
- 周期相等:当前选票信息和接收的选票信息进行PK
1【Java】 2//接收的选票与当前选票PK 3protected boolean totalOrderPredicate(long newId, long newZxid, long newEpoch, long curId, long curZxid, long curEpoch) { 4 LOG.debug("id: " + newId + ", proposed id: " + curId + ", zxid: 0x" + 5 Long.toHexString(newZxid) + ", proposed zxid: 0x" + Long.toHexString(curZxid)); 6 if(self.getQuorumVerifier().getWeight(newId) == 0){ 7 return false; 8 } 9 10 11 /* 12 * We return true if one of the following three cases hold: 13 * 1- New epoch is higher 14 * 2- New epoch is the same as current epoch, but new zxid is higher 15 * 3- New epoch is the same as current epoch, new zxid is the same 16 * as current zxid, but server id is higher. 17 */ 18 return ((newEpoch > curEpoch) || 19 ((newEpoch == curEpoch) && 20 ((newZxid > curZxid) || ((newZxid == curZxid) && (newId > curId)))));wId > curId))))); 21 } 22 23 24 25 26 27
在上述代码中的totalOrderPredicate方法逻辑如下:
- 竞选周期大于当前周期为true
- 竞选周期相等,竞选zxid大于当前zxid为true
- 竞选周期相等,竞选zxid等于当前zxid,竞选sid大于当前sid为true
- 经过上述条件判断为true将当前选票信息替换为竞选成功的选票,同时再次将新的选票投出去。
1【Java】 2public Vote lookForLeader() throws InterruptedException { 3//... 4 //存储节点对应的选票信息 5 // key:选票来源sid value:选票推举的Leader sid 6 recvset.put(n.sid, new Vote(n.leader, n.zxid, n.electionEpoch, n.peerEpoch)); 7 8 9 //半数选举开始 10 if (termPredicate(recvset, 11 new Vote(proposedLeader, proposedZxid, 12 logicalclock.get(), proposedEpoch))) { 13 // Verify if there is any change in the proposed leader 14 while((n = recvqueue.poll(finalizeWait, 15 TimeUnit.MILLISECONDS)) != null){ 16 if(totalOrderPredicate(n.leader, n.zxid, n.peerEpoch, 17 proposedLeader, proposedZxid, proposedEpoch)){ 18 recvqueue.put(n); 19 break; 20 } 21 } 22 /*WorkerSender 23 * This predicate is true once we don't read any new 24 * relevant message from the reception queue 25 */ 26 if (n == null) { 27 //已选举出leader 更新当前节点是否为leader 28 self.setPeerState((proposedLeader == self.getId()) ? 29 ServerState.LEADING: learningState()); 30 31 32 Vote endVote = new Vote(proposedLeader, 33 proposedZxid, proposedEpoch); 34 leaveInstance(endVote); 35 return endVote; 36 } 37 } 38//... 39} 40/** 41 * Termination predicate. Given a set of votes, determines if have 42 * sufficient to declare the end of the election round. 43 * 44 * @param votes 45 * Set of votes 46 * @param vote 47 * Identifier of the vote received last PK后的选票 48 */ 49private boolean termPredicate(HashMap<Long, Vote> votes, Vote vote) { 50 SyncedLearnerTracker voteSet = new SyncedLearnerTracker(); 51 voteSet.addQuorumVerifier(self.getQuorumVerifier()); 52 if (self.getLastSeenQuorumVerifier() != null 53 && self.getLastSeenQuorumVerifier().getVersion() > self 54 .getQuorumVerifier().getVersion()) { 55 voteSet.addQuorumVerifier(self.getLastSeenQuorumVerifier()); 56 } 57 /* 58 * First make the views consistent. Sometimes peers will have different 59 * zxids for a server depending on timing. 60 */ 61 //votes 来源于recvset 存储各个节点推举出来的选票信息 62 for (Map.Entry<Long, Vote> entry : votes.entrySet()) { 63//选举出的sid和其它节点选择的sid相同存储到voteSet变量中。 64 if (vote.equals(entry.getValue())) { 65//保存推举出来的sid 66 voteSet.addAck(entry.getKey()); 67 } 68 } 69 //判断选举出来的选票数量是否过半 70 return voteSet.hasAllQuorums(); 71} 72//QuorumMaj#containsQuorum 73public boolean containsQuorum(Set<Long> ackSet) { 74 return (ackSet.size() > half); 75 } 76 77 78 79 80 81
在上述代码中:recvset是存储每个sid推举的选票信息。
第一轮 sid1:vote(1,0,1) ,sid2:vote(2,0,1);
第二轮 sid1:vote(2,0,1) ,sid2:vote(2,0,1)。
最终经过选举信息vote(2,0,1)为推荐leader,并用推荐leader在recvset选票池里比对持相同票数量为2个。因为总共有3个节点参与选举,sid1和sid2都选举sid2为leader,满足票数过半要求,故确认sid2为leader。
- setPeerState更新当前节点角色;
- proposedLeader选举出来的sid和自己sid相等,设置为Leader;
- 上述条件不相等,设置为Follower或Observing;
- 更新currentVote当前选票为Leader的选票vote(2,0,1)。
07总结
通过对Leader选举源码的解析,可以了解到:
-
多个应用节点之间网络通信采用BIO方式进行相互投票,同时保证每个节点之间只使用一个通道,减少网络资源的消耗,足以见得在BIO分布式中间件开发中的技术重要性。
-
基于BIO的基础上,灵活运用多线程和内存消息队列完好实现多层队列架构,每层队列由不同的线程分工协作,提高快速选举性能目的。
-
为BIO在多线程技术上的实践带来了宝贵的经验。
