源码解析
raft
主函数
一个Raft共识算法的实现中的Make函数,用于创建并初始化一个新的Raft节点。让我详细解析这段代码:
函数概述
Make是Raft节点的构造函数,负责初始化节点状态并启动后台goroutine。
代码结构分析
- 调试辅助
1 | roler_string = map[int]string{ |
用于调试时打印节点角色(领导者/候选人/跟随者)的简短表示。
- 节点状态初始化
1 | rf := &Raft{ |
- 日志初始化
1 | rf.log = append(rf.log, LogEntry{ |
添加一个虚拟日志条目(索引0,任期0)。这是Raft实现中的常见做法,简化边界条件处理。
- 定时器初始化
1 | rf.ResetElectionTimer() // 重置选举超时定时器 |
- 选举定时器:用于触发选举超时
- 追加定时器:领导者定期向跟随者发送心跳/追加日志
- 持久化恢复
1 | rf.readPersist(persister.ReadRaftState()) |
从持久化存储中恢复状态(崩溃恢复场景),并设置提交索引为日志的第一个有效索引。
- 启动后台goroutine
1 | go rf.election_ticker() // 选举心跳:处理选举超时 |
三个并发goroutine负责:
- 选举:超时后发起选举
- 日志复制:领导者定期发送心跳/日志
- 应用:将已提交的日志应用到状态机
election_ticker
这是一个Raft算法中选举心跳循环(election_ticker)的实现,负责检测和触发领导者选举。让我详细解析这段代码:
函数概述
election_ticker是一个后台goroutine,定期检查当前节点是否应该发起选举,并在超时时启动新的选举。
代码逐行解析
- 循环控制
1 | for !rf.killed() { |
- 循环运行直到节点被”杀死”(
killed()方法返回true) - 这是分布式系统中常见的优雅退出模式
- 睡眠等待
1 | time.Sleep(ELECTION_TIMER_RESOLUTION * time.Millisecond) |
- 固定间隔检查(通常10-50ms)
- 配合随机的选举超时时间,避免多个节点同时发起选举
- 加锁保护
1 | rf.mu.Lock() |
- 保护共享状态(
ElectionExpireTime,roler,currentTerm等) - 确保并发安全
- 选举超时检查
1 | if time.Now().After(rf.ElectionExpireTime) && (rf.roler == FOLLOWER || rf.roler == CANDIDATE) { |
两个条件必须同时满足:
- 时间超时:当前时间超过了选举超时时间
- 角色限制:只有跟随者(FOLLOWER)或候选人(CANDIDATE)才能发起选举
- 领导者(LEADER)不会发起选举
- 发起选举
1 | // 转换角色为候选人并重置定时器 |
选举流程:
changeToCandidate():- 将角色切换为
CANDIDATE - 增加
currentTerm - 投票给自己(
votedFor = rf.me)
- 将角色切换为
ResetElectionTimer():- 重置选举超时时间(随机化,防止冲突)
- 并发请求投票:
- 为每个其他节点启动一个goroutine
- 传递当前任期、最后日志索引和任期
- 用于选举规则中的日志比较
关键设计特点
- 随机超时
1 | // 在ResetElectionTimer中实现 |
- 每个节点有独立的随机超时时间
- 减少选举冲突的可能性
- 角色限制
只有跟随者和候选人能发起选举:
- 跟随者:检测到领导者无响应时发起选举
- 候选人:选举超时后重新发起选举(防止选举僵局)
- 领导者:不应该发起选举,否则会导致任期不断增长
- 并发RPC
1 | go rf.CallForVote(i, rf.currentTerm, rf.getLastIndex(), rf.getLastTerm()) |
- 并行向所有节点请求投票
- 提高选举效率
- 每个RPC独立处理,互不影响
执行流程图
1 | 开始 → 睡眠固定间隔 → 检查超时 → |
CallForVote
函数概述
功能:向指定节点发送RequestVote RPC请求,并处理其响应。
调用时机:节点成为Candidate后,并发地向所有其他节点调用此函数。
逐段解析
- 构造RPC参数(1-7行)
1 | args := RequestVoteArgs{} |
- 创建请求/回复对象
- 填充参数:当前任期、候选人ID、最后日志信息
LastLogIndex/Term用于投票者判断候选者的日志是否足够新(Raft的日志完整性保证)
- 发送RPC(8行)
1 | ok := rf.sendRequestVote(idx, &args, &reply) |
- 实际的网络调用
ok=true表示成功收到回复;false表示超时或网络故障- 失败时直接忽略,不处理
- 响应处理入口(10-12行)
1 | if ok { |
- 只有成功收到回复才处理
- 加锁保护共享状态,确保并发安全
defer保证函数退出时解锁
- 第一层任期检查(14-17行)
1 | if reply.Term < rf.currentTerm { |
- 场景:收到旧任期的回复(网络延迟导致)
- 处理:直接丢弃,因为当前节点已进入更新的任期
- 原因:旧任期的投票结果对当前状态无效
- 第二层任期检查(19-24行)
1 | if reply.Term > rf.currentTerm { |
- 场景:发现更高的任期(集群中已有更新的任期)
- 处理:
- 降级为Follower,更新任期
- 重置选举计时器,避免立即重新选举
- 原因:高任期意味着已有新Leader或更新候选者,当前节点应承认其权威
-1表示清空投票记录(未投票给任何人)
- 角色校验(26-30行)
1 | if rf.roler != CANDIDATE { |
- 场景:在处理响应期间,节点角色已改变
- 可能情况:
- 已收到多数票成为Leader
- 被更高任期的节点降级为Follower
- 选举超时重新发起新选举
- 处理:忽略此响应,因为已不再需要
- 第三层任期检查(32-34行)
1 | if reply.Term != args.Term { |
- 目的:确保回复任期与请求任期一致
- 必要性:虽然检查了与
rf.currentTerm的关系,但rf.currentTerm可能在处理期间被修改 - 保证:
reply.Term == args.Term == rf.currentTerm三者一致
- 处理投票结果(36-49行)
1 | if reply.VoteGranted { |
8.1 获得投票(36-37行)
VoteGranted=true表示目标节点投票给当前候选人- 增加获得票数计数
8.2 多数派判断(38行)
1 | if rf.receiveVoteNum > len(rf.peers)/2 |
- 判断是否超过半数节点投票
- 例如:5节点需要3票(3>2),3节点需要2票(2>1)
- 使用
>而非>=,确保严格多数
8.3 成为Leader(39行)
1 | rf.changeToLeader() |
- 切换角色为Leader
- 初始化Leader相关状态(nextIndex、matchIndex等)
8.4 立即发送心跳(40-47行)
1 | for i := 0; i < len(rf.peers); i++ { |
- 遍历所有其他节点
ResetAppendTimer(i, true):立即触发AppendEntries RPC(心跳)- 目的:快速宣告领导权,防止其他节点发起新选举
关键设计思想
- 任期机制(Term)
任期是Raft的逻辑时钟,用于:
- 识别过期信息(
reply.Term < rf.currentTerm) - 检测更新领导者(
reply.Term > rf.currentTerm) - 保证全局顺序一致性
- 多层防护
三层检查确保状态一致性:
- 第一层:与当前任期比较
- 第二层:处理更高任期
- 第三层:与请求任期比较
- 多数派原则
- 只有获得严格多数票才能成为Leader
- 保证不会出现两个Leader(安全性)
- 快速收敛
成为Leader后立即发送心跳:
- 减少集群选举时间
- 避免其他候选者继续选举
append_ticker
函数概述
功能:Leader节点周期性检查并触发向所有Follower节点的日志复制或快照同步。
调用时机:节点成为Leader后,此goroutine持续运行,周期性执行。
逐段解析
- 主循环结构(2-3行)
1 | for !rf.killed() { |
- 无限循环直到节点被杀死(
killed()返回true) - 固定间隔唤醒(通常10-20ms),检查是否需要发送消息
- 比心跳间隔(通常50-150ms)更频繁,确保及时响应
- 角色检查(6-10行)
1 | rf.mu.Lock() |
- 加锁保护共享状态
- 检查角色:只有Leader才执行心跳/日志复制
- 如果不再是Leader(被降级),立即解锁并跳过本次循环
- 关键:必须先加锁再检查,防止并发修改
- 遍历所有Follower(11-14行)
1 | for i := 0; i < len(rf.peers); i++ { |
- 遍历所有其他节点(Follower)
- 跳过自己(不需要给自己发送)
- 检查是否到达发送时间(15-17行)
1 | if !time.Now().After(rf.AppendExpireTime[i]) { |
rf.AppendExpireTime[i]:每个节点的下一次发送时间time.Now().After():检查当前时间是否已超过预定时间- 如果未到时间,跳过该节点(避免过于频繁发送)
- 设计目的:每个节点独立定时,实现流控和负载均衡
- 判断发送内容:Snapshot还是Logs(18-27行)
5.1 Snapshot分支(18-21行)
1 | if rf.nextIndex[i] <= rf.getFirstIndex() { |
- 触发条件:
nextIndex[i] <= getFirstIndex()nextIndex[i]:Leader认为应该发送给节点i的下一条日志索引getFirstIndex():日志中第一条有效索引(快照覆盖了之前的日志)- 含义:节点需要的日志已被快照覆盖,无法通过日志复制同步
- 处理:发送快照(InstallSnapshot RPC)
- 包含当前任期、快照的最后一个索引/任期、快照数据
- 使用
go异步发送,不阻塞循环
5.2 Logs分支(22-27行)
1 | else { |
- 触发条件:
nextIndex[i] > getFirstIndex(),日志可用 - 构造日志切片:
- 从
nextIndex[i]到getLastIndex()的所有日志条目 - 计算偏移:
rf.log[rf.nextIndex[i]-rf.getFirstIndex():]
- 从
- 发送AppendEntries RPC:
- 参数:目标节点、任期、LeaderID、前一条日志索引/任期、日志条目、提交索引
- 使用
go异步发送
- 重置定时器(28行)
1 | rf.ResetAppendTimer(i, false) |
- 重置该节点的发送定时器
false参数:表示不立即发送,而是设置新的到期时间(通常为心跳间隔)- 确保下次周期性发送
- 解锁(29行)
1 | rf.mu.Unlock() |
- 释放锁,允许其他goroutine操作
关键数据结构
nextIndex[]
- Leader维护的每个Follower的下一条日志索引
- 初始化为Leader的最后日志索引+1
- 随着日志复制成功而递增,失败时递减(回退机制)
getFirstIndex()和getLastIndex()
getFirstIndex():快照后第一条日志的索引(日志截断后的起点)getLastIndex():最后一条日志的索引- 日志存储为切片,但逻辑索引可能从非0开始
AppendExpireTime[]
- 每个节点的下一次发送时间
- 避免所有节点同时发送,实现平滑流控
执行流程图
1 | 启动 → 循环睡眠10ms → 加锁 → 检查是否为Leader |
设计思想
- 解耦定时与发送
append_ticker只负责触发,实际RPC在独立goroutine执行- 避免RPC阻塞定时循环
- 提高并发性能
- 独立定时器
- 每个Follower有独立的发送时间
- 可以错峰发送,避免网络拥塞
- 可以根据网络状况动态调整
- 快照优先
- 当日志已被截断时,优先发送快照
- 保证落后节点能快速追上
- 符合Raft的快照机制
- 自适应处理
- 根据
nextIndex自动选择发送快照还是日志 - 发送失败时,
CallAppendEntries会递减nextIndex重试 - 实现自动的日志一致性恢复
CallAppendEntries
函数概述
功能:Leader向Follower发送AppendEntries RPC后,处理其响应,更新复制状态和提交索引。
调用时机:Leader的append_ticker触发发送后,在独立goroutine中异步等待响应。
逐段解析
- 构造RPC参数(2-11行)
1 | args := AppendEntriesArgs{} |
- 创建请求/回复对象
- 填充参数:任期、LeaderID、前一条日志索引/任期、日志条目、Leader的提交索引
PrevLogIndex/Term:用于Follower进行日志一致性检查Entries:待复制的日志条目(可能为空,此时为心跳)
- 发送RPC并检查结果(13-15行)
1 | ok := rf.sendAppendEntries(idx, &args, &reply) |
- 实际网络调用
ok=true表示收到回复,进入处理逻辑- 加锁保护共享状态
- 任期检查(17-24行)
3.1 回复任期小于当前任期(17-19行)
1 | if reply.Term < rf.currentTerm { |
- 旧任期的延迟响应,直接丢弃
3.2 回复任期大于当前任期(20-24行)
1 | if reply.Term > rf.currentTerm { |
- 发现更高任期,说明已有新Leader
- 立即降级为Follower,承认新Leader权威
- 重置选举计时器
- 角色检查(26-28行)
1 | if rf.roler != LEADER { |
- 如果节点不再是Leader,忽略此响应
- 可能已被降级或收到更高任期RPC
- 请求任期检查(30-35行)
1 | if args.Term != rf.currentTerm { |
- 确保请求任期与当前任期一致
- 防止过期响应对应到错误任期
- 注释说明:Leader可能在任期i发送请求,后来成为任期i+2的Leader,收到旧响应应忽略
- nextIndex一致性检查(37-39行)
1 | if rf.nextIndex[idx] != args.PrevLogIndex+1 { |
- 关键检查:确保响应针对的是最新的发送请求
- 如果
nextIndex已变化(例如被其他响应更新),说明此响应已过期 - 防止重复或乱序响应导致状态错误
- 处理成功响应(41-68行)
7.1 更新nextIndex和matchIndex(43-46行)
1 | if reply.Success { |
- 复制成功,推进复制进度
nextIndex前进到prevLogIndex + len(logs) + 1(即最后复制日志的索引+1)matchIndex更新为最后复制的日志索引(nextIndex - 1)- 注意:此时
matchIndex是Follower已确认复制的最大日志索引
7.2 计算可提交的最大索引(47-59行)
1 | diff := make([]int, rf.getLastIndex()+5) |
这是统计多数派提交索引的高效算法:
- 使用差分数组统计每个索引被多少节点复制
diff[0] += 1:所有节点至少从索引0开始diff[matchIndex+1] -= 1:在matchIndex+1处结束计数- 前缀和得到每个索引的复制节点数
diff[i]+1 > len(rf.peers)/2:判断是否超过半数(+1包括自己)ok_idx:最后一个满足多数派的索引
示例:5节点,matchIndex分别为[3,2,4,5]
1 | diff: [4, 0,0, -1,-1,-1, -1] |
7.3 更新提交索引(60-67行)
1 | if ok_idx > rf.commitIdx && rf.getTermForIndex(ok_idx) == rf.currentTerm { |
- 两个条件:
ok_idx > rf.commitIdx:有新的可提交索引getTermForIndex(ok_idx) == rf.currentTerm:关键安全保证!只提交当前任期的日志
- 原因:Raft不允许提交之前任期的日志(即使被多数派复制),必须通过当前任期日志提交来间接提交
- 将新提交的日志添加到
commitQueue,等待应用层消费 rf.cv.Broadcast():唤醒等待提交的goroutine(applier)
- 处理失败响应(69-86行)
8.1 优化回退机制(70-84行)
1 | else { |
Follower返回的冲突信息:
ConflictTerm:冲突日志的任期号(-1表示没有冲突)ConflictIndex:冲突的第一个索引
优化回退策略:
- 没有冲突任期(
ConflictTerm == -1):直接设置为ConflictIndex - 有冲突任期:
- 在Leader日志中查找该任期号的日志
- 如果找到:
nextIndex设置为该任期最后一条日志的下一个索引 - 如果没找到:设置为
ConflictIndex
- 优势:快速跳过整个冲突任期,避免逐条回退(比论文中简单的
nextIndex--效率更高)
- 立即重试(88-91行)
1 | if rf.nextIndex[idx] != rf.getLastIndex()+1 { |
- 无论成功还是失败,只要
nextIndex不是最后索引+1 - 说明Follower日志还未完全同步
- 立即发送下一次AppendEntries(
true表示立即触发) - 目的:加速日志复制过程
关键数据结构
nextIndex[]
- Leader为每个Follower维护的下一条日志索引
- 初始值:Leader最后日志索引+1
- 成功时增加,失败时减小
matchIndex[]
- Leader为每个Follower维护的已复制最大日志索引
- 初始值:0
- 用于计算多数派提交索引
commitIndex
- Leader知道已被多数派复制的最大日志索引
- 只有当前任期的日志才能被提交
commitQueue
- 待应用到状态机的日志队列
- 由
appliergoroutine消费
执行流程图
1 | 发送AppendEntries → 收到响应 → 加锁 |
设计思想
- 幂等性保证
- 通过多重检查(任期、角色、nextIndex)确保只处理有效响应
- 防止过期或重复响应破坏状态
- 快速恢复
- 失败时使用优化回退,快速定位不一致位置
- 成功后立即检查是否需要继续复制
- 提交安全性
- 只提交当前任期的日志(Raft的安全性保证)
- 通过commitQueue异步应用,解耦复制和应用
- 高效统计
- 使用差分数组计算多数派索引,O(n)复杂度
- 相比排序方法更高效
AppendEntries
函数概述
功能:Follower接收Leader发送的AppendEntries请求(心跳或日志复制),进行日志一致性检查并更新自身状态。
调用时机:Leader通过RPC调用Follower的AppendEntries方法。
逐段解析
- 加锁与调试日志(2-5行)
1 | rf.mu.Lock() |
- 加锁保护所有状态修改
- 调试日志记录接收前后状态(使用defer确保函数退出时记录)
- 任期检查(8-12行)
1 | if args.Term < rf.currentTerm { |
- 如果请求任期小于当前任期,拒绝请求
- 回复中带上当前任期,让Leader知道自己过时了
- 注意:这里没有处理
args.Term > rf.currentTerm的情况,在函数末尾统一处理
- 获取当前日志边界(14-15行)
1 | curLastLogIndex := rf.getLastIndex() |
curLastLogIndex:当前最后一条日志的索引curFirstLogIndex:当前第一条日志的索引(快照的最后一个索引)
- 日志一致性检查(17-35行)
4.1 快照边界检查(17-19行)
1 | if args.PrevLogIndex < curFirstLogIndex { |
- 理论上这种情况应该由Leader避免(Leader会在发送日志前检查是否需要发快照)
- 如果发生,记录错误日志(实际应该拒绝或要求快照)
4.2 日志不匹配处理(20-35行)
1 | else if args.PrevLogIndex > curLastLogIndex || rf.getTermForIndex(args.PrevLogIndex) != args.PrevLogTerm { |
触发条件:
PrevLogIndex > curLastLogIndex:日志不够长- 或
getTermForIndex(PrevLogIndex) != PrevLogTerm:日志任期不匹配
优化回退机制:
日志不够长:
ConflictIndex = getLastIndex() + 1(告诉Leader从哪里开始)ConflictTerm = -1(表示没有冲突任期)
日志任期不匹配:
ConflictTerm:冲突位置的任期ConflictIndex:该任期第一条日志的索引- 例如:日志为[1,1,2,2,3],冲突在第4条(任期2),ConflictIndex=3
日志复制处理(37-57行)
5.1 成功响应(37-38行)
1 | reply.Success = true |
- 表示日志检查通过,可以开始复制
5.2 查找已匹配的日志(39-47行)
1 | last_match_idx := args.PrevLogIndex |
- 从
PrevLogIndex+1开始,逐条比较日志 - 直到发现不匹配或到达日志末尾
last_match_idx记录最后匹配的索引
5.3 处理部分匹配(48-53行)
1 | if last_match_idx-args.PrevLogIndex != len(args.Entries) { |
- 如果
last_match_idx不等于prevLogIndex + len(entries),说明有部分日志不匹配 - 处理步骤:
- 截断日志:保留到
last_match_idx - 追加新日志:从
last_match_idx+1开始覆盖 - 持久化状态
- 截断日志:保留到
- 这正是Raft的日志覆盖机制:删除冲突后的所有日志,用Leader的日志替代
- 更新提交索引(55-65行)
6.1 更新commitIdx(55-57行)
1 | old_commit_idx := rf.commitIdx |
- 如果Leader的提交索引更大,更新自己的提交索引
min(args.LeaderCommit, getLastIndex()):不能提交超过自己拥有的日志
6.2 将新提交的日志加入队列(58-64行)
1 | if rf.commitIdx > old_commit_idx { |
- 如果有新提交的日志(
commitIdx > old_commit_idx) - 将新提交的日志加入
commitQueue,等待应用层消费 cv.Broadcast():唤醒等待的applier goroutine
- 角色调整与计时器重置(68-72行)
1 | if args.Term > rf.currentTerm || rf.roler != FOLLOWER { |
- 角色调整:
- 如果请求任期更大,或当前不是Follower(可能是Candidate)
- 降级为Follower,更新任期
- 重置选举计时器:收到合法RPC,说明Leader存活,防止发起新选举
核心数据结构
- 回复信息
1 | type AppendEntriesReply struct { |
- 日志存储
rf.log[0]可能是占位日志(对应快照的最后一条)rf.log[i].Index = i + getFirstIndex()- 实际索引 = 切片索引 +
getFirstIndex()
执行流程图
1 | 收到AppendEntries → 加锁 |
设计思想
- 日志一致性检查
- 通过
PrevLogIndex和PrevLogTerm检查 - 确保Follower的日志与Leader在相同位置有相同任期
- 这是Raft保证日志一致性的核心
- 优化回退机制
1 | if args.PrevLogIndex > curLastLogIndex { |
- 相比于论文中简单的
nextIndex--,这是性能优化 - 一次性跳过整个冲突任期,减少往返次数
- 日志覆盖
1 | rf.log = rf.log[0 : last_match_idx - getFirstIndex() + 1] |
- 删除冲突后的日志
- 用Leader的日志替代
- 这是Raft的日志强制覆盖机制
- 提交索引更新
- 不能提交超过自己拥有的日志
- 用
min(args.LeaderCommit, getLastIndex())保证安全性 - 这符合Raft的”只有Leader能决定哪些日志被提交”原则
CallInstallSnapshot
函数概述
功能:Leader向Follower发送快照后,处理其响应,更新复制状态。
调用时机:Leader在append_ticker中检测到nextIndex[i] <= getFirstIndex()时,异步发送InstallSnapshot RPC。
逐段解析
- 构造RPC参数(2-9行)
1 | args := InstallSnapshotArgs{ |
- 创建快照安装请求参数
LastIncludedIndex/Term:快照包含的最后一条日志的索引和任期Snapshot:实际的快照数据(二进制)- 回复对象用于接收响应
- 发送RPC(10-11行)
1 | ok := rf.sendInstallSnapshot(idx, &args, &reply) |
- 发送快照RPC(可能数据量大,耗时较长)
- 成功收到回复后加锁处理
- 任期检查(14-21行)
3.1 回复任期小于当前任期(14-16行)
1 | if reply.Term < rf.currentTerm { |
- 过期响应,直接忽略
3.2 回复任期大于当前任期(17-21行)
1 | if reply.Term > rf.currentTerm { |
- 发现更高任期,降级为Follower
- 重置选举计时器
- 请求任期检查(23-27行)
1 | if args.Term != rf.currentTerm { |
- 确保请求时的任期与当前任期一致
- 过滤过期响应(Leader可能已变更任期)
- 更新复制状态(29-35行)
1 | if lastIncludedIndex > rf.matchIndex[idx] { |
- 条件:只有快照比当前记录的
matchIndex更新时才更新 - 更新
matchIndex:设置为快照的最后索引(表示该节点已复制到此处) - 更新
nextIndex:设置为lastIncludedIndex + 1(下一条要发送的日志索引)
关键设计思想
- 快照作为日志压缩机制
- 当日志过长或节点落后太多时,发送快照而非逐条日志
- 减少网络传输和存储开销
nextIndex <= getFirstIndex()是触发条件
- 幂等性保证
1 | if lastIncludedIndex > rf.matchIndex[idx] { |
- 防止重复或过期的快照响应覆盖更新的状态
- 如果
matchIndex已经更大,说明有更新的复制进度
- 状态同步
更新matchIndex和nextIndex后:
- Leader知道该Follower已经同步到快照位置
- 后续可以从
nextIndex开始继续复制日志 - 保持复制进度的一致性
执行流程图
1 | 发送快照 → 收到响应 → 加锁 |
与AppendEntries的对比
| 特性 | AppendEntries | InstallSnapshot |
|---|---|---|
| 传输内容 | 增量日志 | 完整快照 |
| 触发条件 | nextIndex > getFirstIndex() |
nextIndex <= getFirstIndex() |
| 数据量 | 小(增量) | 大(完整状态) |
| 更新方式 | 递增推进 | 跳跃式更新 |
| 提交计算 | 影响commitIndex | 不影响commitIndex |
边界情况处理
情况1:重复快照响应
1 | if lastIncludedIndex > rf.matchIndex[idx] { |
- 防止旧快照覆盖新进度
情况2:快照发送后节点变为Leader
- 函数开始时会检查
rf.roler != LEADER(但在这段代码中检查在append_ticker里) - 如果不再是Leader,RPC函数会通过任期检查降级
情况3:节点重启
- 收到快照后,Follower会截断日志并应用快照
- Leader更新状态后继续复制
installSnapshot
函数概述
功能:Follower接收Leader发送的快照,更新自己的状态机和日志。
调用时机:Leader检测到Follower需要快照同步时,发送InstallSnapshot RPC。
逐段解析
- 加锁与任期检查(2-9行)
1 | rf.mu.Lock() |
- 加锁:保护共享状态修改
- 任期检查:如果请求任期小于当前任期,拒绝并返回当前任期
- 回复任期设置为
rf.currentTerm,让Leader知道自己的真实任期
- 调试日志(11-16行)
1 | Debug(dSnap, "S%d Receive Snapshot From S%d T%d, LII: %d, LIT:%d, Snap:%v", |
- 记录接收到的快照信息
- 记录处理前的日志状态,便于调试
- 获取当前状态(18-19行)
1 | curSnapLastIndex := rf.getFirstIndex() |
curSnapLastIndex:当前快照的最后索引(日志起始索引)curLogLastIndex:当前最后一条日志的索引
- 快照有效性检查(20-23行)
1 | if args.LastIncludedIndex <= curSnapLastIndex { |
- 如果快照不比当前更新:忽略此RPC
<=表示快照已经应用过或更旧- 设置回复任期后返回
- 主要处理逻辑(24-56行)
5.1 将快照加入应用队列(25-31行)
1 | rf.commitQueue = append(rf.commitQueue, ApplyMsg{ |
- 创建一个快照类型的ApplyMsg
CommandValid: false表示这不是普通日志条目SnapshotValid: true表示这是一个快照- 放入
commitQueue等待应用层消费
5.2 情况A:快照包含部分日志(32-47行)
1 | if args.LastIncludedIndex < curLogLastIndex { |
触发条件:快照包含的日志少于当前日志(LastIncludedIndex < curLogLastIndex)
处理步骤:
- 截断日志:
- 保留
LastIncludedIndex之后的日志 - 计算新日志长度:
curLogLastIndex - args.LastIncludedIndex + 1 - 从旧日志中复制
[args.LastIncludedIndex-curSnapLastIndex:]部分
- 保留
- 清空第一个日志的命令:
rf.log[0].Command = nil:占位日志不包含实际命令- 这个日志代表快照的最后一条记录
- 处理提交索引:
- 如果
LastIncludedIndex < rf.commitIdx:有已提交的日志在快照之后- 重新将这些已提交日志加入
commitQueue - 确保应用层能处理这些命令
- 重新将这些已提交日志加入
- 否则:
commitIdx更新为快照的最后索引
- 如果
5.3 情况B:快照比当前日志更新(48-54行)
1 | else { |
触发条件:快照包含的日志比当前日志更新(LastIncludedIndex >= curLogLastIndex)
处理步骤:
清空日志:创建空日志切片
添加占位日志:只保留快照的最后一条记录(无命令)
更新提交索引:
commitIdx设置为快照的最后索引持久化(55-56行)
1 | state_seri_result := SerilizeState(rf) |
- 序列化当前状态(任期、投票、日志等)
- 同时保存状态和快照到持久化存储
- 崩溃恢复时可以从这里恢复
- 角色调整(57-59行)
1 | if args.Term > rf.currentTerm || rf.roler != FOLLOWER { |
- 如果请求任期更大,或当前不是Follower
- 降级为Follower,更新任期
- 保证集群一致性
- 重置选举计时器(60行)
1 | rf.ResetElectionTimer() |
- 收到合法RPC,重置选举计时器
- 防止在同步快照期间发起不必要的选举
- 唤醒等待的goroutine(56行)
1 | rf.cv.Broadcast() |
- 唤醒等待提交的
appliergoroutine - 通知有新的提交条目(快照)需要应用
关键数据结构
- 日志存储结构
1 | [Entry0(占位)] [Entry1] [Entry2] ... [EntryN] |
- 第一个条目是占位条目,
Command=nil - 逻辑索引从
LastIncludedIndex开始
- ApplyMsg类型
1 | type ApplyMsg struct { |
执行流程图
1 | 收到快照 → 加锁 → 任期检查 |
设计思想
- 日志压缩与快照
- 快照替代了日志中较早的部分
- 通过截断日志减少存储开销
- 保留快照之后的新日志
- 提交索引处理
1 | if args.LastIncludedIndex < rf.commitIdx { |
- 快照可能覆盖部分已提交的日志
- 需要重新提交快照之后的已提交日志
- 确保应用层状态正确
- 占位日志
1 | rf.log[0].Command = nil |
- 第一个日志作为占位,代表快照的最后状态
- 索引对齐:
log[0].Index = LastIncludedIndex - 简化索引计算
- 幂等性
1 | if args.LastIncludedIndex <= curSnapLastIndex { |
- 重复接收已应用的快照直接忽略
- 保证处理的幂等性
边界情况
情况1:快照覆盖部分已提交日志
1 | 原有日志: [1,2,3,4,5] commitIdx=3 |
- 保证commitIdx=3的日志仍然被应用
情况2:快照覆盖所有日志
1 | 原有日志: [1,2,3] commitIdx=3 |
- 所有旧日志被快照替代
情况3:快照比当前旧
1 | 当前: LastIncludedIndex=5 |
Applier
这段代码是Raft算法中应用层日志提交器(Applier)的核心实现。我来全面解析:
函数概述
功能:从提交队列中取出已提交的日志和快照,通过channel传递给上层应用(状态机)。
调用时机:节点启动时创建独立的goroutine运行,持续运行直到节点被杀死。
逐段解析
- 主循环(2-3行)
1 | for !rf.killed() { |
- 无限循环,直到节点被杀死(
rf.killed()返回true) - 这是一个长期运行的goroutine
- 加锁与等待(4-6行)
1 | rf.mu.Lock() |
- 加锁:保护
commitQueue的读取 - 条件等待:如果队列为空,阻塞等待
rf.cv.Wait():释放锁并挂起,直到被cv.Broadcast()唤醒- 唤醒时机:
AppendEntries中提交新日志时InstallSnapshot中收到快照时CallAppendEntries中Leader提交日志时(Leader的applier也会使用)
- 取出所有消息(7-9行)
1 | msgs := rf.commitQueue |
- 取出队列:保存引用到局部变量
msgs - 清空队列:重置为空切片,准备接收新消息
- 解锁:尽快释放锁,允许其他goroutine继续提交
- 这是一个批量处理设计:一次取出所有已提交的消息
- 处理消息(10-20行)
1 | for _, msg := range msgs { |
- 遍历所有消息
- 日志和快照分别打印不同的调试信息
- 核心操作:
applyCh <- msg将消息发送给上层应用- 这是阻塞操作:如果应用层处理慢,会阻塞在此
- 但这是必要的:确保消息按顺序应用
关键数据结构
- ApplyMsg
1 | type ApplyMsg struct { |
- commitQueue
1 | type Raft struct { |
执行流程图
1 | 启动 → 无限循环 |
设计思想
- 生产者-消费者模式
- 生产者:
AppendEntries:收到Leader的AppendEntries时,将新提交的日志放入队列InstallSnapshot:收到快照时,将快照放入队列CallAppendEntries:Leader自己提交日志时
- 消费者:
Appliergoroutine消费队列
- 条件变量同步
1 | rf.cv.Wait() // 等待 |
- 避免忙等(busy-waiting),节省CPU
- 生产者在放入消息后立即唤醒消费者
- 批量处理
1 | msgs := rf.commitQueue |
- 一次处理所有累积的消息
- 减少锁的获取次数
- 提高吞吐量
- 解耦Raft与状态机
- Raft通过
applyCh将消息传递给上层 - 上层可以是KV服务器、数据库等
- 实现了Raft协议与应用逻辑的解耦
channel说明
applyCh <- msg 是Go语言中向channel发送数据的操作。在这段代码中,它表示将已提交的日志条目或快照通过channel传递给上层应用(如KV状态机)。
Channel基本概念
什么是Channel?
1 | applyCh chan ApplyMsg // 声明一个channel,传递ApplyMsg类型的数据 |
- Channel是Go语言中用于goroutine间通信的管道
- 遵循FIFO(先进先出)原则
- 可以理解为线程安全的队列
发送和接收操作
1 | // 发送:向channel中放入数据 |
在Raft中的角色架构图
1 | ┌─────────────┐ ┌─────────────┐ ┌─────────────┐ |
数据流转
- Raft层(生产者):
- 决定哪些日志可以提交(commitIdx)
- 将提交的日志封装成ApplyMsg
- 放入
commitQueue队列
- Applier(中转者):
- 从
commitQueue取出消息 - 通过
applyCh <- msg发送给上层
- 从
- 应用层(消费者):
- 从
applyCh接收消息 - 执行命令(如PUT/GET)
- 更新状态机状态
- 从
为什么使用Channel?
- 解耦
1 | // Raft不需要知道应用层如何执行命令 |
- 顺序保证
1 | // Channel保证消息按发送顺序接收 |
- 线程安全
1 | // 多个goroutine可以安全地发送/接收 |
- 阻塞同步
1 | // 发送阻塞:如果channel满了,发送者会等待 |
具体示例
创建Channel(在测试代码中)
1 | applyCh := make(chan ApplyMsg) // 无缓冲channel |
消费者(应用层)
1 | func (kv *KVServer) applier() { |
生产者(Raft层)
1 | // 在Applier goroutine中 |
KvRaft
调用时序图
1 | 客户端 KVServer Raft层 状态机 |