
回顾

Leader,
Candidate和
Follower。当集群启动时,每个节点的状态都是Follower,因为集群中没有Leader,所以在一个随机的Election Timeout之后,某一节点转换成Candidate状态并发起投票,如果获得多数票,则此节点成为Leader。这里“随机”的Election Timeout非常重要,因为如果每个节点都使用一个固定的Timeout的话,他们都会几乎在同一时刻发起自己的Election,导致谁也当不上Leader,这种情况称之为
SplitVote。当Candidate成功当选为Leader时,发起会定期给所有Follower发送Hearbeat信号以刷新它们的Election Timer。当Leader通过收到的RPC请求或者回复发现有Term比自己高的Leader时,马上转变成Follower。
目标


定义Raft节点
Raft类,除了论文中提到的几个状态,笔者还定义了几个信号量,
killSig用来通知各个goroutine节点已经被kill掉了,
hearbeatSig用来刷新Follower的Election Timer,
stopListingHeartbeatSig用来通知专属Follower的goroutine退出,
stopSendingHeartbeatSig用来通知专属Leader的goroutine退出。
type Raft struct {mu sync.Mutex // Lock to protect shared access to this peer's statepeers []*labrpc.ClientEnd // RPC end points of all peersme int // this peer's Index into peers[]// Paper-defined Raft statescurrentTerm intvotedFor intlog []*LogEntrycommitIndex intlastApplied intnextIndex []intmatchIndex []int// Self-defined Raft statesisLeader bool// None raft relatedkillSig chan struct{}heartbeatSig chan struct{}stopListeningHeartbeatSig chan struct{}stopSendingHeartbeatSig chan struct{}}
节点启动
Make函数将启动并返回一个Raft对象。在初始化节点状态时,
rf.isLeader为
false,对应了状态图中所有节点初始状态为
Follower的情况。在函数的最后启动了
rf.listenHearbeatLoop来监听来自可能的Leader的心跳信号,同时在Election Timeout之内没有收到心跳的话,节点将发起投票。
Make函数有一个参数是
peers,可以看成其他所有节点的地址,用来发送和接受RPC请求。
func Make(peers []*labrpc.ClientEnd, me int,persister *Persister, applyCh chan ApplyMsg) *Raft {rf := &Raft{}rf.peers = peersrf.me = me// Your initialization code here (2A, 2B, 2C).// Initialize Raft statesrf.currentTerm = -1rf.votedFor = -1rf.log = make([]*LogEntry, 0)rf.commitIndex = -1rf.lastApplied = -1rf.nextIndex = make([]int, len(peers))rf.matchIndex = make([]int, len(peers))rf.isLeader = false// Initialize signal channelsrf.killSig = make(chan struct{}, 1)rf.heartbeatSig = make(chan struct{}, 1)rf.stopListeningHeartbeatSig = make(chan struct{}, 1)rf.stopSendingHeartbeatSig = make(chan struct{}, 1)go rf.listenHeartbeatLoop()return rf}
发起投票
listenHeartbeatLoop中,Follower节点在每个循环中创建一个随机的timer来实现Election Timeout,然后同时监听来自
rf.hearbeatSig,
timer.C,
rf.killSig和
rf.stopListenHeartbeatSig的信号,当在Election Timeout之内(timer.C 触发)没有收到心跳信号,节点将通过
rf.startElection()发起投票。
func (rf *Raft) listenHeartbeatLoop() {defer DPrintf("raft %d exiting listenHeartbeatLoop\n", rf.me)for {timer := time.NewTimer(time.Duration(rand.Intn(maxElectionTimeout-minElectionTimeout)+minElectionTimeout) * time.Millisecond)select {case <-rf.heartbeatSig:continuecase <-timer.C:// Turn to candidate and request for votesDPrintf("raft %d turning Candidate\n", rf.me)go rf.startElection()case <-rf.stopListeningHeartbeatSig:returncase <-rf.killSig:return}}}
startElection方法开始先锁住节点,然后将
rf.currentTerm加一,更新
rf.votedFor为自己,并记录下
electionTerm。随后并发地向除自己外所有节点发送
rf.requestVote()请求,然后等待所有请求完成并收集票数。当票数收集完毕时,节点首先判断自己的
rf.currentTerm是否为当初的
electionTerm,这一点非常重要,因为节点在等候票数的时候可能收到了其他节点的
RequestVote请求并且投出了自己的一票,导致当前term比electionTerm要高。如果这种情况出现的话,本轮投票立即终止,否则节点将会通过错误的任期成为Leader。在本文中,所有RPC请求结束后都会做一次这种验证。最后当节点确认自己收到的票数过半时,通过
rf.turnLeader()成为Leader。
rf.turnLeader()将
rf.isLeader设为True,并发送
rf.stopListeningHearbeatSig关闭
rf.listenHeartbeatLoop。最后启动
rf.sendHeartbeatLoop()给所有folloer发送心跳。
rf.mu解锁,这会导致当两个节点同时投票时出现死锁的bug。
func (rf *Raft) startElection() {rf.mu.Lock()if rf.isLeader {rf.mu.Unlock()return}rf.currentTerm += 1rf.votedFor = rf.meelectionTerm := rf.currentTermrf.mu.Unlock()votes := make(chan bool, len(rf.peers)-1)for server := range rf.peers {if server == rf.me {continue}go rf.requestVote(electionTerm, server, votes)}numGranted := 1for i := 0; i < len(rf.peers)-1; i++ {if vote := <-votes; vote {numGranted++}if numGranted > len(rf.peers)/2 {break}}rf.mu.Lock()defer rf.mu.Unlock()if electionTerm != rf.currentTerm {return}if numGranted > len(rf.peers)/2 {if !rf.isLeader {rf.turnLeader()}}return}func (rf *Raft) turnLeader() {DPrintf("raft %d turning to leader\n", rf.me)rf.isLeader = trueselect {case rf.stopListeningHeartbeatSig <- struct{}{}:default:}go rf.sendHeartbeatLoop()}
rf.requestVote对单个其他节点发送求票请求。一开始在请求中设置term与本节点的id。LastLogIndex与LastLogTerm暂时忽略,因为这属于第二篇文章(Log replication)的范畴。当收到回复时,节点首先判断回复的Term是否大于
rf.currentTerm,一旦条件为真,说明了比自己term高的Leader已经产生了,此时需要更新自己的
rf.currentTerm并且放弃本轮票选。如果没有发现比自己term高的Leader,节点判断
reply.VoteGranted是否为真,如果为真,说明得到了对面节点的投票。
func (rf *Raft) requestVote(electionTerm int, server int, votes chan bool) {args := &RequestVoteArgs{Term: electionTerm,CandidateId: rf.me,// TODO: fill these twoLastLogIndex: 0,LastLogTerm: 0,}reply := &RequestVoteReply{}DPrintf("raft %d sending RequestVote %v to raft %d \n", rf.me, args, server)if success := rf.sendRequestVote(server, args, reply); !success {DPrintf("raft %d sendRequestVote RPC %v to raft %d failure\n", rf.me, args, server)votes <- falsereturn}DPrintf("raft %d got RequestVote reply %v from raft %d\n", rf.me, reply, server)rf.mu.Lock()defer rf.mu.Unlock()if args.Term != rf.currentTerm {votes <- falsereturn}if reply.Term > rf.currentTerm {if rf.isLeader {rf.turnFollower()}rf.currentTerm = reply.Termrf.votedFor = -1votes <- falsereturn}if reply.VoteGranted {votes <- truereturn}votes <- falsereturn}
处理投票请求
rf.RequestVote方法定义了一个节点收到其他节点RequestVote RPC时的行为。当节点收到其他节点的求票请求时,首先判断对方的Term是否小于自己的
rf.currentTerm,如果小于,那么告知对方最新的Term并且拒绝请求。如果对方的Term等于自己的Term,并且对方的id不等于自己的
rf.VotedFor,说明节点在这个Term已经给其他节点(可以包括自己)投票了,也拒绝请求。否则将票投给对方并更新自己的Term。在方法最后最好刷新一下节点的hearbeat信号,防止在投票之后和收到新Leader心跳信号之前自己发起新一轮投票。
func (rf *Raft) RequestVote(args *RequestVoteArgs, reply *RequestVoteReply) {DPrintf("raft %d received RequestVoteRequest %v from raft %d\n", rf.me, args, args.CandidateId)defer DPrintf("raft %d returning RequestVoteResponse %v to raft %d\n", rf.me, reply, args.CandidateId)// Your code here (2A, 2B).rf.mu.Lock()defer rf.mu.Unlock()if args.Term < rf.currentTerm {reply.Term = rf.currentTermreply.VoteGranted = falsereturn}if args.Term == rf.currentTerm && rf.votedFor >= 0 && rf.votedFor != args.CandidateId {reply.Term = args.Termreply.VoteGranted = falsereturn}// TODO: implement log checking belowif rf.isLeader {rf.turnFollower()}// update current Term and grant voterf.currentTerm = args.Termrf.votedFor = args.CandidateIdreply.Term = args.Termreply.VoteGranted = truerf.refreshHeartbeat()}
心跳发送
发起投票部分中我们提到了当节点成功当选Leader后会启动
rf.sendHearbeatLoop()来持续地给所有节点发送心跳以维持自己Leader的地位。这个循环会每隔一段时间调用
rf.sendHeatbeats()来广播一次心跳。
rf.sendHeartbeats()会并发地向所有Follower节点调用
rf.sendHeartbeat()以提高性能。
rf.sendHeartbeat()通过
AppendEntriesRPC向对方节点发送自己的Term与id,当收到回复后,判断回复的term是否大于自己的
rf.currentTerm,如果大于,更新自己的term并马上step down成为follower。
func (rf *Raft) sendHeartbeatLoop() {rf.sendHeartbeats()defer DPrintf("raft %d exiting sendHeartbeatLoop\n", rf.me)ticker := time.NewTicker(heartbeatPeriod * time.Millisecond)for {select {case <-ticker.C:go rf.sendHeartbeats()case <-rf.stopSendingHeartbeatSig:returncase <-rf.killSig:return}}}func (rf *Raft) sendHeartbeats() {rf.mu.Lock()if !rf.isLeader {rf.mu.Unlock()return}term := rf.currentTermrf.mu.Unlock()replies := make(chan bool, len(rf.peers)-1)for server := range rf.peers {if server == rf.me {continue}go rf.sendHeartbeat(term, server, replies)}numSuccess := 1for i := 0; i < len(rf.peers)-1; i++ {if reply := <-replies; reply {numSuccess++}}if numSuccess <= len(rf.peers)/2 {rf.mu.Lock()defer rf.mu.Unlock()if rf.isLeader {rf.turnFollower()}}}func (rf *Raft) sendHeartbeat(term int, server int, replies chan bool) {args := &AppendEntriesArgs{Term: term,LeaderId: rf.me,// TODO: fill the field belowPrevLogIndex: 0,PrevLogTerm: 0,Entries: nil,LeaderCommit: 0,}reply := &AppendEntriesReply{}if success := rf.sendAppendEntries(server, args, reply); !success {DPrintf("raft %d sendAppendEntries RPC %v to raft %d failure\n", rf.me, args, server)replies <- falsereturn}rf.mu.Lock()defer rf.mu.Unlock()if args.Term != rf.currentTerm {replies <- falsereturn}if !reply.Success {if reply.Term > rf.currentTerm {if rf.isLeader {rf.turnFollower()}rf.currentTerm = reply.Termrf.votedFor = -1}replies <- falsereturn}replies <- true}
心跳接收
AppendEntries(心跳)请求时,首先判断对方的term是否小于自己的
rf.currentTerm,如果小于,那么告知对方最新的term并拒绝请求。否则接受请求并刷新自己的heartbeat信号。
func (rf *Raft) AppendEntries(args *AppendEntriesArgs, reply *AppendEntriesReply) {DPrintf("raft %d received AppendEntriesRequest %v from raft %d\n", rf.me, args, args.LeaderId)defer DPrintf("raft %d return AppendEntiresResponse %v to raft %d\n", rf.me, reply, args.LeaderId)rf.mu.Lock()defer rf.mu.Unlock()if args.Term < rf.currentTerm {reply.Term = rf.currentTermreply.Success = falsereturn}reply.Term = args.Termreply.Success = truerf.currentTerm = args.Termif rf.isLeader {rf.votedFor = -1rf.turnFollower()}rf.refreshHeartbeat()// TODO: implement append entriesreturn}func (rf *Raft) refreshHeartbeat() {select {case rf.heartbeatSig <- struct{}{}:default:}}
总结
LeaderElection,代码量也就三百来行,但是其对并发编程和Raft论文的的理解都有着比较高的要求。相信通过细读这篇文章,读者能够对Raft有第一步的了解,在第二篇文章中笔者将带大家实现一遍
Logreplication。
exportGOPATH=<some_path_prefix>/6.824git checkout9b5e2f5f2b7128c3360c957c60f8d006ccc83382实战Apache Kafka
CSP范式编程

文章转载自Techdemic,如果涉嫌侵权,请发送邮件至:contact@modb.pro进行举报,并提供相关证据,一经查实,墨天轮将立刻删除相关内容。




