6.824源码解析

源码解析

raft

主函数

一个Raft共识算法的实现中的Make函数,用于创建并初始化一个新的Raft节点。让我详细解析这段代码:

函数概述

Make是Raft节点的构造函数,负责初始化节点状态并启动后台goroutine。

代码结构分析

  1. 调试辅助
1
2
3
4
5
roler_string = map[int]string{
LEADER: "L",
CANDIDATE: "C",
FOLLOWER: "F",
}

用于调试时打印节点角色(领导者/候选人/跟随者)的简短表示。

  1. 节点状态初始化
1
2
3
4
5
6
7
8
9
10
11
12
13
14
rf := &Raft{
peers: peers, // 所有节点的RPC客户端
persister: persister, // 持久化存储
me: me, // 当前节点ID
currentTerm: 0, // 当前任期
votedFor: -1, // 当前任期投票给谁(-1表示未投票)
log: make([]LogEntry, 0), // 日志条目
roler: FOLLOWER, // 初始为跟随者
AppendExpireTime: make([]time.Time, num_servers), // 每个节点的追加超时
nextIndex: make([]int, num_servers), // 每个节点的下一条日志索引
matchIndex: make([]int, num_servers), // 每个节点已匹配的日志索引
commitQueue: make([]ApplyMsg, 0), // 待提交队列
}
rf.cv = sync.NewCond(&rf.mu) // 条件变量,用于等待状态变化
  1. 日志初始化
1
2
3
4
5
rf.log = append(rf.log, LogEntry{
Index: 0,
Term: 0,
Command: nil,
})

添加一个虚拟日志条目(索引0,任期0)。这是Raft实现中的常见做法,简化边界条件处理。

  1. 定时器初始化
1
2
3
4
rf.ResetElectionTimer()              // 重置选举超时定时器
for i := 0; i < num_servers; i++ {
rf.ResetAppendTimer(i, false) // 为每个节点重置追加日志定时器
}
  • 选举定时器:用于触发选举超时
  • 追加定时器:领导者定期向跟随者发送心跳/追加日志
  1. 持久化恢复
1
2
rf.readPersist(persister.ReadRaftState())
rf.commitIdx = rf.getFirstIndex()

从持久化存储中恢复状态(崩溃恢复场景),并设置提交索引为日志的第一个有效索引。

  1. 启动后台goroutine
1
2
3
go rf.election_ticker()   // 选举心跳:处理选举超时
go rf.append_ticker() // 追加日志心跳:领导者发送心跳
go rf.Applier(applyCh) // 异步应用已提交的日志到状态机

三个并发goroutine负责:

  • 选举:超时后发起选举
  • 日志复制:领导者定期发送心跳/日志
  • 应用:将已提交的日志应用到状态机

election_ticker

这是一个Raft算法中选举心跳循环election_ticker)的实现,负责检测和触发领导者选举。让我详细解析这段代码:

函数概述

election_ticker是一个后台goroutine,定期检查当前节点是否应该发起选举,并在超时时启动新的选举。

代码逐行解析

  1. 循环控制
1
for !rf.killed() {
  • 循环运行直到节点被”杀死”(killed()方法返回true)
  • 这是分布式系统中常见的优雅退出模式
  1. 睡眠等待
1
time.Sleep(ELECTION_TIMER_RESOLUTION * time.Millisecond)
  • 固定间隔检查(通常10-50ms)
  • 配合随机的选举超时时间,避免多个节点同时发起选举
  1. 加锁保护
1
2
3
rf.mu.Lock()
// ... 临界区代码 ...
rf.mu.Unlock()
  • 保护共享状态(ElectionExpireTime, roler, currentTerm等)
  • 确保并发安全
  1. 选举超时检查
1
if time.Now().After(rf.ElectionExpireTime) && (rf.roler == FOLLOWER || rf.roler == CANDIDATE) {

两个条件必须同时满足:

  • 时间超时:当前时间超过了选举超时时间
  • 角色限制:只有跟随者(FOLLOWER)或候选人(CANDIDATE)才能发起选举
    • 领导者(LEADER)不会发起选举
  1. 发起选举
1
2
3
4
5
6
7
8
9
10
11
// 转换角色为候选人并重置定时器
rf.changeToCandidate()
rf.ResetElectionTimer()

// 向其他节点请求投票
for i := range rf.peers {
if i == rf.me {
continue
}
go rf.CallForVote(i, rf.currentTerm, rf.getLastIndex(), rf.getLastTerm())
}

选举流程

  1. changeToCandidate()
    • 将角色切换为CANDIDATE
    • 增加currentTerm
    • 投票给自己(votedFor = rf.me
  2. ResetElectionTimer()
    • 重置选举超时时间(随机化,防止冲突)
  3. 并发请求投票
    • 为每个其他节点启动一个goroutine
    • 传递当前任期、最后日志索引和任期
    • 用于选举规则中的日志比较

关键设计特点

  1. 随机超时
1
2
3
4
// 在ResetElectionTimer中实现
rf.ElectionExpireTime = time.Now().Add(
time.Duration(ELECTION_TIMEOUT_MIN + rand.Intn(ELECTION_TIMEOUT_MAX - ELECTION_TIMEOUT_MIN)) * time.Millisecond
)
  • 每个节点有独立的随机超时时间
  • 减少选举冲突的可能性
  1. 角色限制

只有跟随者候选人能发起选举:

  • 跟随者:检测到领导者无响应时发起选举
  • 候选人:选举超时后重新发起选举(防止选举僵局)
  • 领导者:不应该发起选举,否则会导致任期不断增长
  1. 并发RPC
1
go rf.CallForVote(i, rf.currentTerm, rf.getLastIndex(), rf.getLastTerm())
  • 并行向所有节点请求投票
  • 提高选举效率
  • 每个RPC独立处理,互不影响

执行流程图

1
2
3
4
5
开始 → 睡眠固定间隔 → 检查超时 → 

未超时 → 继续睡眠

超时且角色符合 → 转为候选人 → 重置定时器 → 并行请求投票

CallForVote

函数概述

功能:向指定节点发送RequestVote RPC请求,并处理其响应。

调用时机:节点成为Candidate后,并发地向所有其他节点调用此函数。

逐段解析

  1. 构造RPC参数(1-7行)
1
2
3
4
5
6
args := RequestVoteArgs{}
reply := RequestVoteReply{}
args.Term = term
args.CandidateId = rf.me
args.LastLogIndex = lastLogIndex
args.LastLogTerm = lastLogTerm
  • 创建请求/回复对象
  • 填充参数:当前任期、候选人ID、最后日志信息
  • LastLogIndex/Term用于投票者判断候选者的日志是否足够新(Raft的日志完整性保证)
  1. 发送RPC(8行)
1
ok := rf.sendRequestVote(idx, &args, &reply)
  • 实际的网络调用
  • ok=true表示成功收到回复;false表示超时或网络故障
  • 失败时直接忽略,不处理
  1. 响应处理入口(10-12行)
1
2
3
if ok {
rf.mu.Lock()
defer rf.mu.Unlock()
  • 只有成功收到回复才处理
  • 加锁保护共享状态,确保并发安全
  • defer保证函数退出时解锁
  1. 第一层任期检查(14-17行)
1
2
3
if reply.Term < rf.currentTerm {
return
}
  • 场景:收到旧任期的回复(网络延迟导致)
  • 处理:直接丢弃,因为当前节点已进入更新的任期
  • 原因:旧任期的投票结果对当前状态无效
  1. 第二层任期检查(19-24行)
1
2
3
4
5
if reply.Term > rf.currentTerm {
rf.changeToFollower(reply.Term, -1)
rf.ResetElectionTimer()
return
}
  • 场景:发现更高的任期(集群中已有更新的任期)
  • 处理
    • 降级为Follower,更新任期
    • 重置选举计时器,避免立即重新选举
  • 原因:高任期意味着已有新Leader或更新候选者,当前节点应承认其权威
  • -1表示清空投票记录(未投票给任何人)
  1. 角色校验(26-30行)
1
2
3
if rf.roler != CANDIDATE {
return
}
  • 场景:在处理响应期间,节点角色已改变
  • 可能情况
    • 已收到多数票成为Leader
    • 被更高任期的节点降级为Follower
    • 选举超时重新发起新选举
  • 处理:忽略此响应,因为已不再需要
  1. 第三层任期检查(32-34行)
1
2
3
if reply.Term != args.Term {
return
}
  • 目的:确保回复任期与请求任期一致
  • 必要性:虽然检查了与rf.currentTerm的关系,但rf.currentTerm可能在处理期间被修改
  • 保证reply.Term == args.Term == rf.currentTerm三者一致
  1. 处理投票结果(36-49行)
1
2
3
4
5
6
7
8
9
10
11
if reply.VoteGranted {
DebugGetVote(rf.me, idx, term)
rf.receiveVoteNum += 1
if rf.receiveVoteNum > len(rf.peers)/2 {
rf.changeToLeader()
for i := 0; i < len(rf.peers); i++ {
if i == rf.me { continue }
rf.ResetAppendTimer(i, true)
}
}
}

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
2
3
4
for i := 0; i < len(rf.peers); i++ {
if i == rf.me { continue }
rf.ResetAppendTimer(i, true)
}
  • 遍历所有其他节点
  • ResetAppendTimer(i, true):立即触发AppendEntries RPC(心跳)
  • 目的:快速宣告领导权,防止其他节点发起新选举

关键设计思想

  1. 任期机制(Term)

任期是Raft的逻辑时钟,用于:

  • 识别过期信息(reply.Term < rf.currentTerm
  • 检测更新领导者(reply.Term > rf.currentTerm
  • 保证全局顺序一致性
  1. 多层防护

三层检查确保状态一致性:

  • 第一层:与当前任期比较
  • 第二层:处理更高任期
  • 第三层:与请求任期比较
  1. 多数派原则
  • 只有获得严格多数票才能成为Leader
  • 保证不会出现两个Leader(安全性)
  1. 快速收敛

成为Leader后立即发送心跳:

  • 减少集群选举时间
  • 避免其他候选者继续选举

append_ticker

函数概述

功能:Leader节点周期性检查并触发向所有Follower节点的日志复制或快照同步。

调用时机:节点成为Leader后,此goroutine持续运行,周期性执行。

逐段解析

  1. 主循环结构(2-3行)
1
2
for !rf.killed() {
time.Sleep(APPEND_TIMER_RESOLUTION * time.Millisecond)
  • 无限循环直到节点被杀死(killed()返回true)
  • 固定间隔唤醒(通常10-20ms),检查是否需要发送消息
  • 比心跳间隔(通常50-150ms)更频繁,确保及时响应
  1. 角色检查(6-10行)
1
2
3
4
5
rf.mu.Lock()
if rf.roler != LEADER {
rf.mu.Unlock()
continue
}
  • 加锁保护共享状态
  • 检查角色:只有Leader才执行心跳/日志复制
  • 如果不再是Leader(被降级),立即解锁并跳过本次循环
  • 关键:必须先加锁再检查,防止并发修改
  1. 遍历所有Follower(11-14行)
1
2
3
4
for i := 0; i < len(rf.peers); i++ {
if i == rf.me {
continue
}
  • 遍历所有其他节点(Follower)
  • 跳过自己(不需要给自己发送)
  1. 检查是否到达发送时间(15-17行)
1
2
3
if !time.Now().After(rf.AppendExpireTime[i]) {
continue
}
  • rf.AppendExpireTime[i]:每个节点的下一次发送时间
  • time.Now().After():检查当前时间是否已超过预定时间
  • 如果未到时间,跳过该节点(避免过于频繁发送)
  • 设计目的:每个节点独立定时,实现流控和负载均衡
  1. 判断发送内容:Snapshot还是Logs(18-27行)

5.1 Snapshot分支(18-21行)

1
2
3
4
5
if rf.nextIndex[i] <= rf.getFirstIndex() {
go rf.CallInstallSnapshot(i, rf.currentTerm,
rf.me, rf.getFirstIndex(),
rf.getFirstTerm(), rf.persister.ReadSnapshot())
}
  • 触发条件nextIndex[i] <= getFirstIndex()
    • nextIndex[i]:Leader认为应该发送给节点i的下一条日志索引
    • getFirstIndex():日志中第一条有效索引(快照覆盖了之前的日志)
    • 含义:节点需要的日志已被快照覆盖,无法通过日志复制同步
  • 处理:发送快照(InstallSnapshot RPC)
    • 包含当前任期、快照的最后一个索引/任期、快照数据
    • 使用go异步发送,不阻塞循环

5.2 Logs分支(22-27行)

1
2
3
4
5
6
7
else {
logs := make([]LogEntry, rf.getLastIndex()-rf.nextIndex[i]+1)
copy(logs, rf.log[rf.nextIndex[i]-rf.getFirstIndex():])
go rf.CallAppendEntries(i, rf.currentTerm, rf.me,
rf.nextIndex[i]-1, rf.getTermForIndex(rf.nextIndex[i]-1),
logs, rf.commitIdx)
}
  • 触发条件nextIndex[i] > getFirstIndex(),日志可用
  • 构造日志切片
    • nextIndex[i]getLastIndex()的所有日志条目
    • 计算偏移:rf.log[rf.nextIndex[i]-rf.getFirstIndex():]
  • 发送AppendEntries RPC
    • 参数:目标节点、任期、LeaderID、前一条日志索引/任期、日志条目、提交索引
    • 使用go异步发送
  1. 重置定时器(28行)
1
rf.ResetAppendTimer(i, false)
  • 重置该节点的发送定时器
  • false参数:表示不立即发送,而是设置新的到期时间(通常为心跳间隔)
  • 确保下次周期性发送
  1. 解锁(29行)
1
rf.mu.Unlock()
  • 释放锁,允许其他goroutine操作

关键数据结构

  1. nextIndex[]
  • Leader维护的每个Follower的下一条日志索引
  • 初始化为Leader的最后日志索引+1
  • 随着日志复制成功而递增,失败时递减(回退机制)
  1. getFirstIndex()getLastIndex()
  • getFirstIndex():快照后第一条日志的索引(日志截断后的起点)
  • getLastIndex():最后一条日志的索引
  • 日志存储为切片,但逻辑索引可能从非0开始
  1. AppendExpireTime[]
  • 每个节点的下一次发送时间
  • 避免所有节点同时发送,实现平滑流控

执行流程图

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
启动 → 循环睡眠10ms → 加锁 → 检查是否为Leader
↓ 是
遍历所有Follower

检查是否到期
↓ 是
判断发送内容
↙ ↘
nextIndex <= 开头 nextIndex > 开头
↓ ↓
发送快照(InstallSnapshot) 发送日志(AppendEntries)
↓ ↓
重置该节点定时器

解锁

设计思想

  1. 解耦定时与发送
  • append_ticker只负责触发,实际RPC在独立goroutine执行
  • 避免RPC阻塞定时循环
  • 提高并发性能
  1. 独立定时器
  • 每个Follower有独立的发送时间
  • 可以错峰发送,避免网络拥塞
  • 可以根据网络状况动态调整
  1. 快照优先
  • 当日志已被截断时,优先发送快照
  • 保证落后节点能快速追上
  • 符合Raft的快照机制
  1. 自适应处理
  • 根据nextIndex自动选择发送快照还是日志
  • 发送失败时,CallAppendEntries会递减nextIndex重试
  • 实现自动的日志一致性恢复

CallAppendEntries

函数概述

功能:Leader向Follower发送AppendEntries RPC后,处理其响应,更新复制状态和提交索引。

调用时机:Leader的append_ticker触发发送后,在独立goroutine中异步等待响应。

逐段解析

  1. 构造RPC参数(2-11行)
1
2
3
4
5
6
7
8
args := AppendEntriesArgs{}
reply := AppendEntriesReply{}
args.Term = term
args.LeaderId = me
args.PrevLogIndex = prevLogIndex
args.PrevLogTerm = prevLogTerm
args.Entries = logs
args.LeaderCommit = leaderCommit
  • 创建请求/回复对象
  • 填充参数:任期、LeaderID、前一条日志索引/任期、日志条目、Leader的提交索引
  • PrevLogIndex/Term:用于Follower进行日志一致性检查
  • Entries:待复制的日志条目(可能为空,此时为心跳)
  1. 发送RPC并检查结果(13-15行)
1
2
3
4
ok := rf.sendAppendEntries(idx, &args, &reply)
if ok {
rf.mu.Lock()
defer rf.mu.Unlock()
  • 实际网络调用
  • ok=true表示收到回复,进入处理逻辑
  • 加锁保护共享状态
  1. 任期检查(17-24行)

3.1 回复任期小于当前任期(17-19行)

1
2
3
if reply.Term < rf.currentTerm {
return
}
  • 旧任期的延迟响应,直接丢弃

3.2 回复任期大于当前任期(20-24行)

1
2
3
4
5
if reply.Term > rf.currentTerm {
rf.changeToFollower(reply.Term, -1)
rf.ResetElectionTimer()
return
}
  • 发现更高任期,说明已有新Leader
  • 立即降级为Follower,承认新Leader权威
  • 重置选举计时器
  1. 角色检查(26-28行)
1
2
3
if rf.roler != LEADER {
return
}
  • 如果节点不再是Leader,忽略此响应
  • 可能已被降级或收到更高任期RPC
  1. 请求任期检查(30-35行)
1
2
3
if args.Term != rf.currentTerm {
return
}
  • 确保请求任期与当前任期一致
  • 防止过期响应对应到错误任期
  • 注释说明:Leader可能在任期i发送请求,后来成为任期i+2的Leader,收到旧响应应忽略
  1. nextIndex一致性检查(37-39行)
1
2
3
if rf.nextIndex[idx] != args.PrevLogIndex+1 {
return
}
  • 关键检查:确保响应针对的是最新的发送请求
  • 如果nextIndex已变化(例如被其他响应更新),说明此响应已过期
  • 防止重复或乱序响应导致状态错误
  1. 处理成功响应(41-68行)

7.1 更新nextIndex和matchIndex(43-46行)

1
2
3
if reply.Success {
rf.nextIndex[idx] = rf.nextIndex[idx] + len(logs)
rf.matchIndex[idx] = rf.nextIndex[idx] - 1
  • 复制成功,推进复制进度
  • nextIndex前进到prevLogIndex + len(logs) + 1(即最后复制日志的索引+1)
  • matchIndex更新为最后复制的日志索引(nextIndex - 1
  • 注意:此时matchIndex是Follower已确认复制的最大日志索引

7.2 计算可提交的最大索引(47-59行)

1
2
3
4
5
6
7
8
9
10
11
12
13
diff := make([]int, rf.getLastIndex()+5)
for i := 0; i < len(rf.peers); i++ {
if i == rf.me { continue }
diff[0] += 1
diff[rf.matchIndex[i]+1] -= 1
}
ok_idx := 0
for i := 1; i < len(diff); i++ {
diff[i] += diff[i-1]
if diff[i]+1 > len(rf.peers)/2 {
ok_idx = i
}
}

这是统计多数派提交索引的高效算法

  • 使用差分数组统计每个索引被多少节点复制
  • 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
2
3
diff: [4, 0,0, -1,-1,-1, -1]
前缀和: [4,4,4,3,2,1,0]
ok_idx=3 (索引3有4个节点复制)

7.3 更新提交索引(60-67行)

1
2
3
4
5
6
7
if ok_idx > rf.commitIdx && rf.getTermForIndex(ok_idx) == rf.currentTerm {
for i := rf.commitIdx + 1; i <= ok_idx; i++ {
rf.commitQueue = append(rf.commitQueue, ApplyMsg{...})
}
rf.commitIdx = ok_idx
rf.cv.Broadcast()
}
  • 两个条件
    1. ok_idx > rf.commitIdx:有新的可提交索引
    2. getTermForIndex(ok_idx) == rf.currentTerm关键安全保证!只提交当前任期的日志
  • 原因:Raft不允许提交之前任期的日志(即使被多数派复制),必须通过当前任期日志提交来间接提交
  • 将新提交的日志添加到commitQueue,等待应用层消费
  • rf.cv.Broadcast():唤醒等待提交的goroutine(applier
  1. 处理失败响应(69-86行)

8.1 优化回退机制(70-84行)

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
else {
if reply.ConflictTerm == -1 {
rf.nextIndex[idx] = reply.ConflictIndex
} else {
findIdx := -1
for i := rf.getLastIndex() + 1; i > rf.getFirstIndex(); i-- {
if rf.getTermForIndex(i-1) == reply.ConflictTerm {
findIdx = i
break
}
}
if findIdx != -1 {
rf.nextIndex[idx] = findIdx
} else {
rf.nextIndex[idx] = reply.ConflictIndex
}
}
}

Follower返回的冲突信息

  • ConflictTerm:冲突日志的任期号(-1表示没有冲突)
  • ConflictIndex:冲突的第一个索引

优化回退策略

  1. 没有冲突任期ConflictTerm == -1):直接设置为ConflictIndex
  2. 有冲突任期
    • 在Leader日志中查找该任期号的日志
    • 如果找到:nextIndex设置为该任期最后一条日志的下一个索引
    • 如果没找到:设置为ConflictIndex
  • 优势:快速跳过整个冲突任期,避免逐条回退(比论文中简单的nextIndex--效率更高)
  1. 立即重试(88-91行)
1
2
3
if rf.nextIndex[idx] != rf.getLastIndex()+1 {
rf.ResetAppendTimer(idx, true)
}
  • 无论成功还是失败,只要nextIndex不是最后索引+1
  • 说明Follower日志还未完全同步
  • 立即发送下一次AppendEntries(true表示立即触发)
  • 目的:加速日志复制过程

关键数据结构

  1. nextIndex[]
  • Leader为每个Follower维护的下一条日志索引
  • 初始值:Leader最后日志索引+1
  • 成功时增加,失败时减小
  1. matchIndex[]
  • Leader为每个Follower维护的已复制最大日志索引
  • 初始值:0
  • 用于计算多数派提交索引
  1. commitIndex
  • Leader知道已被多数派复制的最大日志索引
  • 只有当前任期的日志才能被提交
  1. commitQueue
  • 待应用到状态机的日志队列
  • applier goroutine消费

执行流程图

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
发送AppendEntries → 收到响应 → 加锁

任期检查(<, >, =)

角色检查(仍为Leader)

请求任期检查(args.Term == curTerm)

nextIndex一致性检查

↙ 成功 ↘ 失败
更新nextIndex/matchIndex 优化回退nextIndex
↓ ↓
计算多数派提交索引 立即重试

更新commitIndex
添加到commitQueue
唤醒applier

设计思想

  1. 幂等性保证
  • 通过多重检查(任期、角色、nextIndex)确保只处理有效响应
  • 防止过期或重复响应破坏状态
  1. 快速恢复
  • 失败时使用优化回退,快速定位不一致位置
  • 成功后立即检查是否需要继续复制
  1. 提交安全性
  • 只提交当前任期的日志(Raft的安全性保证)
  • 通过commitQueue异步应用,解耦复制和应用
  1. 高效统计
  • 使用差分数组计算多数派索引,O(n)复杂度
  • 相比排序方法更高效

AppendEntries

函数概述

功能:Follower接收Leader发送的AppendEntries请求(心跳或日志复制),进行日志一致性检查并更新自身状态。

调用时机:Leader通过RPC调用Follower的AppendEntries方法。

逐段解析

  1. 加锁与调试日志(2-5行)
1
2
3
4
5
rf.mu.Lock()
defer rf.mu.Unlock()

DebugReceiveAppendEntries(rf, args)
defer DebugAfterReceiveAppendEntries(rf, args, reply)
  • 加锁保护所有状态修改
  • 调试日志记录接收前后状态(使用defer确保函数退出时记录)
  1. 任期检查(8-12行)
1
2
3
4
5
if args.Term < rf.currentTerm {
reply.Success = false
reply.Term = rf.currentTerm
return
}
  • 如果请求任期小于当前任期,拒绝请求
  • 回复中带上当前任期,让Leader知道自己过时了
  • 注意:这里没有处理args.Term > rf.currentTerm的情况,在函数末尾统一处理
  1. 获取当前日志边界(14-15行)
1
2
curLastLogIndex := rf.getLastIndex()
curFirstLogIndex := rf.getFirstIndex()
  • curLastLogIndex:当前最后一条日志的索引
  • curFirstLogIndex:当前第一条日志的索引(快照的最后一个索引)
  1. 日志一致性检查(17-35行)

4.1 快照边界检查(17-19行)

1
2
3
if args.PrevLogIndex < curFirstLogIndex {
Debug(dError, "S%d, PrevlogIndex %d is in the snapshot! %d", rf.me, args.PrevLogIndex, curFirstLogIndex)
}
  • 理论上这种情况应该由Leader避免(Leader会在发送日志前检查是否需要发快照)
  • 如果发生,记录错误日志(实际应该拒绝或要求快照)

4.2 日志不匹配处理(20-35行)

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
else if args.PrevLogIndex > curLastLogIndex || rf.getTermForIndex(args.PrevLogIndex) != args.PrevLogTerm {
reply.Success = false
reply.Term = args.Term

// 优化回退:返回冲突信息
if args.PrevLogIndex > curLastLogIndex {
reply.ConflictIndex = rf.getLastIndex() + 1
reply.ConflictTerm = -1
} else {
reply.ConflictTerm = rf.getTermForIndex(args.PrevLogIndex)
findIdx := args.PrevLogIndex
for i := args.PrevLogIndex; i > rf.getFirstIndex(); i-- {
if rf.getTermForIndex(i-1) != reply.ConflictTerm {
findIdx = i
break
}
}
reply.ConflictIndex = findIdx
}
}

触发条件

  • PrevLogIndex > curLastLogIndex:日志不够长
  • getTermForIndex(PrevLogIndex) != PrevLogTerm:日志任期不匹配

优化回退机制

  1. 日志不够长

    • ConflictIndex = getLastIndex() + 1(告诉Leader从哪里开始)
    • ConflictTerm = -1(表示没有冲突任期)
  2. 日志任期不匹配

    • ConflictTerm:冲突位置的任期
    • ConflictIndex:该任期第一条日志的索引
    • 例如:日志为[1,1,2,2,3],冲突在第4条(任期2),ConflictIndex=3
  3. 日志复制处理(37-57行)

5.1 成功响应(37-38行)

1
2
reply.Success = true
reply.Term = args.Term
  • 表示日志检查通过,可以开始复制

5.2 查找已匹配的日志(39-47行)

1
2
3
4
5
6
7
8
9
10
last_match_idx := args.PrevLogIndex
for i := 0; i < len(args.Entries); i++ {
if args.PrevLogIndex+1+i > curLastLogIndex {
break
}
if rf.getTermForIndex(args.PrevLogIndex+1+i) != args.Entries[i].Term {
break
}
last_match_idx = args.PrevLogIndex + 1 + i
}
  • PrevLogIndex+1开始,逐条比较日志
  • 直到发现不匹配或到达日志末尾
  • last_match_idx记录最后匹配的索引

5.3 处理部分匹配(48-53行)

1
2
3
4
5
6
if last_match_idx-args.PrevLogIndex != len(args.Entries) {
// 部分匹配
rf.log = rf.log[0 : last_match_idx-rf.getFirstIndex()+1]
rf.log = append(rf.log, args.Entries[last_match_idx-args.PrevLogIndex:]...)
rf.persist()
}
  • 如果last_match_idx不等于prevLogIndex + len(entries),说明有部分日志不匹配
  • 处理步骤
    1. 截断日志:保留到last_match_idx
    2. 追加新日志:从last_match_idx+1开始覆盖
    3. 持久化状态
  • 这正是Raft的日志覆盖机制:删除冲突后的所有日志,用Leader的日志替代
  1. 更新提交索引(55-65行)

6.1 更新commitIdx(55-57行)

1
2
3
4
old_commit_idx := rf.commitIdx
if args.LeaderCommit > rf.commitIdx {
rf.commitIdx = min(args.LeaderCommit, rf.getLastIndex())
}
  • 如果Leader的提交索引更大,更新自己的提交索引
  • min(args.LeaderCommit, getLastIndex()):不能提交超过自己拥有的日志

6.2 将新提交的日志加入队列(58-64行)

1
2
3
4
5
6
7
8
9
10
if rf.commitIdx > old_commit_idx {
for i := old_commit_idx + 1; i <= rf.commitIdx; i++ {
rf.commitQueue = append(rf.commitQueue, ApplyMsg{
CommandValid: true,
CommandIndex: i,
Command: rf.getCommand(i),
})
}
rf.cv.Broadcast()
}
  • 如果有新提交的日志(commitIdx > old_commit_idx
  • 将新提交的日志加入commitQueue,等待应用层消费
  • cv.Broadcast():唤醒等待的applier goroutine
  1. 角色调整与计时器重置(68-72行)
1
2
3
4
if args.Term > rf.currentTerm || rf.roler != FOLLOWER {
rf.changeToFollower(args.Term, -1)
}
rf.ResetElectionTimer()
  • 角色调整
    • 如果请求任期更大,或当前不是Follower(可能是Candidate)
    • 降级为Follower,更新任期
  • 重置选举计时器:收到合法RPC,说明Leader存活,防止发起新选举

核心数据结构

  1. 回复信息
1
2
3
4
5
6
type AppendEntriesReply struct {
Term int // 当前任期
Success bool // 是否成功
ConflictTerm int // 冲突的任期
ConflictIndex int // 冲突的第一个索引
}
  1. 日志存储
  • rf.log[0]可能是占位日志(对应快照的最后一条)
  • rf.log[i].Index = i + getFirstIndex()
  • 实际索引 = 切片索引 + getFirstIndex()

执行流程图

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
收到AppendEntries → 加锁

任期检查(<)

↙ ↘
日志不匹配 日志匹配
(返回失败) (开始复制)
↓ ↓
返回冲突信息 查找已匹配日志
ConflictIndex ↓
ConflictTerm 部分匹配 → 截断+覆盖

更新commitIdx

加入commitQueue

唤醒applier

降级为Follower
重置选举计时器

设计思想

  1. 日志一致性检查
  • 通过PrevLogIndexPrevLogTerm检查
  • 确保Follower的日志与Leader在相同位置有相同任期
  • 这是Raft保证日志一致性的核心
  1. 优化回退机制
1
2
3
4
5
6
7
if args.PrevLogIndex > curLastLogIndex {
ConflictIndex = getLastIndex() + 1
ConflictTerm = -1
} else {
ConflictTerm = getTermForIndex(PrevLogIndex)
ConflictIndex = 该任期第一条日志
}
  • 相比于论文中简单的nextIndex--,这是性能优化
  • 一次性跳过整个冲突任期,减少往返次数
  1. 日志覆盖
1
2
rf.log = rf.log[0 : last_match_idx - getFirstIndex() + 1]
rf.log = append(rf.log, args.Entries[last_match_idx - args.PrevLogIndex:]...)
  • 删除冲突后的日志
  • 用Leader的日志替代
  • 这是Raft的日志强制覆盖机制
  1. 提交索引更新
  • 不能提交超过自己拥有的日志
  • min(args.LeaderCommit, getLastIndex())保证安全性
  • 这符合Raft的”只有Leader能决定哪些日志被提交”原则

CallInstallSnapshot

函数概述

功能:Leader向Follower发送快照后,处理其响应,更新复制状态。

调用时机:Leader在append_ticker中检测到nextIndex[i] <= getFirstIndex()时,异步发送InstallSnapshot RPC。

逐段解析

  1. 构造RPC参数(2-9行)
1
2
3
4
5
6
7
8
args := InstallSnapshotArgs{
Term: term,
LeaderId: me,
LastIncludedIndex: lastIncludedIndex,
LastIncludedTerm: lastIncludedTerm,
Snapshot: data,
}
reply := InstallSnapshotReply{}
  • 创建快照安装请求参数
  • LastIncludedIndex/Term:快照包含的最后一条日志的索引和任期
  • Snapshot:实际的快照数据(二进制)
  • 回复对象用于接收响应
  1. 发送RPC(10-11行)
1
2
3
4
ok := rf.sendInstallSnapshot(idx, &args, &reply)
if ok {
rf.mu.Lock()
defer rf.mu.Unlock()
  • 发送快照RPC(可能数据量大,耗时较长)
  • 成功收到回复后加锁处理
  1. 任期检查(14-21行)

3.1 回复任期小于当前任期(14-16行)

1
2
3
if reply.Term < rf.currentTerm {
return
}
  • 过期响应,直接忽略

3.2 回复任期大于当前任期(17-21行)

1
2
3
4
5
if reply.Term > rf.currentTerm {
rf.changeToFollower(reply.Term, -1)
rf.ResetElectionTimer()
return
}
  • 发现更高任期,降级为Follower
  • 重置选举计时器
  1. 请求任期检查(23-27行)
1
2
3
if args.Term != rf.currentTerm {
return
}
  • 确保请求时的任期与当前任期一致
  • 过滤过期响应(Leader可能已变更任期)
  1. 更新复制状态(29-35行)
1
2
3
4
if lastIncludedIndex > rf.matchIndex[idx] {
rf.matchIndex[idx] = lastIncludedIndex
rf.nextIndex[idx] = lastIncludedIndex + 1
}
  • 条件:只有快照比当前记录的matchIndex更新时才更新
  • 更新matchIndex:设置为快照的最后索引(表示该节点已复制到此处)
  • 更新nextIndex:设置为lastIncludedIndex + 1(下一条要发送的日志索引)

关键设计思想

  1. 快照作为日志压缩机制
  • 当日志过长或节点落后太多时,发送快照而非逐条日志
  • 减少网络传输和存储开销
  • nextIndex <= getFirstIndex()是触发条件
  1. 幂等性保证
1
2
3
if lastIncludedIndex > rf.matchIndex[idx] {
// 只有新快照才更新
}
  • 防止重复或过期的快照响应覆盖更新的状态
  • 如果matchIndex已经更大,说明有更新的复制进度
  1. 状态同步

更新matchIndexnextIndex后:

  • Leader知道该Follower已经同步到快照位置
  • 后续可以从nextIndex开始继续复制日志
  • 保持复制进度的一致性

执行流程图

1
2
3
4
5
6
7
8
发送快照 → 收到响应 → 加锁

任期检查(<, >, =)

请求任期检查

更新复制状态
(matchIndex/nextIndex)

与AppendEntries的对比

特性 AppendEntries InstallSnapshot
传输内容 增量日志 完整快照
触发条件 nextIndex > getFirstIndex() nextIndex <= getFirstIndex()
数据量 小(增量) 大(完整状态)
更新方式 递增推进 跳跃式更新
提交计算 影响commitIndex 不影响commitIndex

边界情况处理

情况1:重复快照响应

1
2
3
if lastIncludedIndex > rf.matchIndex[idx] {
// 只有更大的索引才更新
}
  • 防止旧快照覆盖新进度

情况2:快照发送后节点变为Leader

  • 函数开始时会检查rf.roler != LEADER(但在这段代码中检查在append_ticker里)
  • 如果不再是Leader,RPC函数会通过任期检查降级

情况3:节点重启

  • 收到快照后,Follower会截断日志并应用快照
  • Leader更新状态后继续复制

installSnapshot

函数概述

功能:Follower接收Leader发送的快照,更新自己的状态机和日志。

调用时机:Leader检测到Follower需要快照同步时,发送InstallSnapshot RPC。

逐段解析

  1. 加锁与任期检查(2-9行)
1
2
3
4
5
6
7
rf.mu.Lock()
defer rf.mu.Unlock()

if args.Term < rf.currentTerm {
reply.Term = rf.currentTerm
return
}
  • 加锁:保护共享状态修改
  • 任期检查:如果请求任期小于当前任期,拒绝并返回当前任期
  • 回复任期设置为rf.currentTerm,让Leader知道自己的真实任期
  1. 调试日志(11-16行)
1
2
3
4
5
Debug(dSnap, "S%d Receive Snapshot From S%d T%d, LII: %d, LIT:%d, Snap:%v",
rf.me, args.LeaderId, args.Term,
args.LastIncludedIndex, args.LastIncludedTerm,
args.Snapshot)
Debug(dSnap, "S%d Before Process, Log is: %v", rf.me, rf.log)
  • 记录接收到的快照信息
  • 记录处理前的日志状态,便于调试
  1. 获取当前状态(18-19行)
1
2
curSnapLastIndex := rf.getFirstIndex()
curLogLastIndex := rf.getLastIndex()
  • curSnapLastIndex:当前快照的最后索引(日志起始索引)
  • curLogLastIndex:当前最后一条日志的索引
  1. 快照有效性检查(20-23行)
1
2
3
4
if args.LastIncludedIndex <= curSnapLastIndex {
reply.Term = args.Term
return
}
  • 如果快照不比当前更新:忽略此RPC
  • <=表示快照已经应用过或更旧
  • 设置回复任期后返回
  1. 主要处理逻辑(24-56行)

5.1 将快照加入应用队列(25-31行)

1
2
3
4
5
6
7
rf.commitQueue = append(rf.commitQueue, ApplyMsg{
CommandValid: false,
SnapshotValid: true,
SnapshotIndex: args.LastIncludedIndex,
SnapshotTerm: args.LastIncludedTerm,
Snapshot: args.Snapshot,
})
  • 创建一个快照类型的ApplyMsg
  • CommandValid: false表示这不是普通日志条目
  • SnapshotValid: true表示这是一个快照
  • 放入commitQueue等待应用层消费

5.2 情况A:快照包含部分日志(32-47行)

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
if args.LastIncludedIndex < curLogLastIndex {
// 日志有额外条目
old_log := rf.log
rf.log = make([]LogEntry, curLogLastIndex-args.LastIncludedIndex+1)
copy(rf.log, old_log[args.LastIncludedIndex-curSnapLastIndex:])
rf.log[0].Command = nil

if args.LastIncludedIndex < rf.commitIdx {
// 重新提交快照之后的日志
for i := args.LastIncludedIndex + 1; i <= rf.commitIdx; i++ {
rf.commitQueue = append(rf.commitQueue, ApplyMsg{
CommandValid: true,
CommandIndex: i,
Command: rf.getCommand(i),
})
}
} else {
rf.commitIdx = args.LastIncludedIndex
}
}

触发条件:快照包含的日志少于当前日志(LastIncludedIndex < curLogLastIndex

处理步骤

  1. 截断日志
    • 保留LastIncludedIndex之后的日志
    • 计算新日志长度:curLogLastIndex - args.LastIncludedIndex + 1
    • 从旧日志中复制[args.LastIncludedIndex-curSnapLastIndex:]部分
  2. 清空第一个日志的命令
    • rf.log[0].Command = nil:占位日志不包含实际命令
    • 这个日志代表快照的最后一条记录
  3. 处理提交索引
    • 如果LastIncludedIndex < rf.commitIdx:有已提交的日志在快照之后
      • 重新将这些已提交日志加入commitQueue
      • 确保应用层能处理这些命令
    • 否则:commitIdx更新为快照的最后索引

5.3 情况B:快照比当前日志更新(48-54行)

1
2
3
4
5
6
7
8
9
else {
rf.log = make([]LogEntry, 0)
rf.log = append(rf.log, LogEntry{
Index: args.LastIncludedIndex,
Term: args.LastIncludedTerm,
Command: nil,
})
rf.commitIdx = args.LastIncludedIndex
}

触发条件:快照包含的日志比当前日志更新(LastIncludedIndex >= curLogLastIndex

处理步骤

  1. 清空日志:创建空日志切片

  2. 添加占位日志:只保留快照的最后一条记录(无命令)

  3. 更新提交索引commitIdx设置为快照的最后索引

  4. 持久化(55-56行)

1
2
state_seri_result := SerilizeState(rf)
rf.persister.SaveStateAndSnapshot(state_seri_result, args.Snapshot)
  • 序列化当前状态(任期、投票、日志等)
  • 同时保存状态和快照到持久化存储
  • 崩溃恢复时可以从这里恢复
  1. 角色调整(57-59行)
1
2
3
if args.Term > rf.currentTerm || rf.roler != FOLLOWER {
rf.changeToFollower(args.Term, -1)
}
  • 如果请求任期更大,或当前不是Follower
  • 降级为Follower,更新任期
  • 保证集群一致性
  1. 重置选举计时器(60行)
1
rf.ResetElectionTimer()
  • 收到合法RPC,重置选举计时器
  • 防止在同步快照期间发起不必要的选举
  1. 唤醒等待的goroutine(56行)
1
rf.cv.Broadcast()
  • 唤醒等待提交的applier goroutine
  • 通知有新的提交条目(快照)需要应用

关键数据结构

  1. 日志存储结构
1
2
3
[Entry0(占位)] [Entry1] [Entry2] ... [EntryN]

LastIncludedIndex(快照的最后索引)
  • 第一个条目是占位条目,Command=nil
  • 逻辑索引从LastIncludedIndex开始
  1. ApplyMsg类型
1
2
3
4
5
6
7
8
9
10
type ApplyMsg struct {
CommandValid bool // 普通日志
Command interface{}
CommandIndex int

SnapshotValid bool // 快照
Snapshot []byte
SnapshotIndex int
SnapshotTerm int
}

执行流程图

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
收到快照 → 加锁 → 任期检查

忽略旧快照(<= current)

加入commitQueue(快照)

↙ ↘
快照在日志中间 快照在日志末尾
(部分保留) (完全覆盖)
↓ ↓
截断日志 清空日志
重新提交已提交日志 添加占位日志
↓ ↓
持久化状态和快照

降级为Follower

重置选举计时器

设计思想

  1. 日志压缩与快照
  • 快照替代了日志中较早的部分
  • 通过截断日志减少存储开销
  • 保留快照之后的新日志
  1. 提交索引处理
1
2
3
if args.LastIncludedIndex < rf.commitIdx {
// 重新提交已提交的日志
}
  • 快照可能覆盖部分已提交的日志
  • 需要重新提交快照之后的已提交日志
  • 确保应用层状态正确
  1. 占位日志
1
rf.log[0].Command = nil
  • 第一个日志作为占位,代表快照的最后状态
  • 索引对齐:log[0].Index = LastIncludedIndex
  • 简化索引计算
  1. 幂等性
1
2
3
if args.LastIncludedIndex <= curSnapLastIndex {
return
}
  • 重复接收已应用的快照直接忽略
  • 保证处理的幂等性

边界情况

情况1:快照覆盖部分已提交日志

1
2
3
原有日志: [1,2,3,4,5] commitIdx=3
快照: LastIncludedIndex=2
处理: 保留[3,4,5],重新提交[3]
  • 保证commitIdx=3的日志仍然被应用

情况2:快照覆盖所有日志

1
2
3
4
原有日志: [1,2,3] commitIdx=3
快照: LastIncludedIndex=5
处理: 清空日志,添加[5]
commitIdx=5
  • 所有旧日志被快照替代

情况3:快照比当前旧

1
2
3
当前: LastIncludedIndex=5
快照: LastIncludedIndex=3
处理: 忽略

Applier

这段代码是Raft算法中应用层日志提交器(Applier)的核心实现。我来全面解析:

函数概述

功能:从提交队列中取出已提交的日志和快照,通过channel传递给上层应用(状态机)。

调用时机:节点启动时创建独立的goroutine运行,持续运行直到节点被杀死。

逐段解析

  1. 主循环(2-3行)
1
2
3
for !rf.killed() {
// ...
}
  • 无限循环,直到节点被杀死(rf.killed()返回true)
  • 这是一个长期运行的goroutine
  1. 加锁与等待(4-6行)
1
2
3
4
rf.mu.Lock()
for len(rf.commitQueue) == 0 {
rf.cv.Wait()
}
  • 加锁:保护commitQueue的读取
  • 条件等待:如果队列为空,阻塞等待
  • rf.cv.Wait():释放锁并挂起,直到被cv.Broadcast()唤醒
  • 唤醒时机
    • AppendEntries中提交新日志时
    • InstallSnapshot中收到快照时
    • CallAppendEntries中Leader提交日志时(Leader的applier也会使用)
  1. 取出所有消息(7-9行)
1
2
3
msgs := rf.commitQueue
rf.commitQueue = make([]ApplyMsg, 0)
rf.mu.Unlock()
  • 取出队列:保存引用到局部变量msgs
  • 清空队列:重置为空切片,准备接收新消息
  • 解锁:尽快释放锁,允许其他goroutine继续提交
  • 这是一个批量处理设计:一次取出所有已提交的消息
  1. 处理消息(10-20行)
1
2
3
4
5
6
7
8
9
10
for _, msg := range msgs {
if msg.CommandValid {
Debug(dLog2, "S%d Apply Command IDX%d CMD: %v", rf.me, msg.CommandIndex, msg.Command)
} else if msg.SnapshotValid {
Debug(dLog2, "S%d Apply Snapshot. LII: %d, LIT: %d, snapShot: %v", rf.me, msg.SnapshotIndex, msg.SnapshotTerm, msg.Snapshot)
} else {
Debug(dError, "S%d, Apply unknown Command!", rf.me)
}
applyCh <- msg
}
  • 遍历所有消息
  • 日志和快照分别打印不同的调试信息
  • 核心操作applyCh <- msg将消息发送给上层应用
    • 这是阻塞操作:如果应用层处理慢,会阻塞在此
    • 但这是必要的:确保消息按顺序应用

关键数据结构

  1. ApplyMsg
1
2
3
4
5
6
7
8
9
10
11
12
type ApplyMsg struct {
// 普通日志
CommandValid bool
Command interface{}
CommandIndex int

// 快照
SnapshotValid bool
Snapshot []byte
SnapshotIndex int
SnapshotTerm int
}
  1. commitQueue
1
2
3
4
5
type Raft struct {
commitQueue []ApplyMsg // 待应用的提交队列
cv *sync.Cond // 条件变量
// ...
}

执行流程图

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
启动 → 无限循环

加锁 → 检查队列
↓ ↓
为空 非空
↓ ↓
cv.Wait() 取出所有消息
等待唤醒 清空队列
↓ 解锁
唤醒 ↓
循环 遍历消息

applyCh <- msg

继续循环

设计思想

  1. 生产者-消费者模式
  • 生产者
    • AppendEntries:收到Leader的AppendEntries时,将新提交的日志放入队列
    • InstallSnapshot:收到快照时,将快照放入队列
    • CallAppendEntries:Leader自己提交日志时
  • 消费者Applier goroutine消费队列
  1. 条件变量同步
1
2
rf.cv.Wait()     // 等待
rf.cv.Broadcast() // 唤醒
  • 避免忙等(busy-waiting),节省CPU
  • 生产者在放入消息后立即唤醒消费者
  1. 批量处理
1
2
msgs := rf.commitQueue
rf.commitQueue = make([]ApplyMsg, 0)
  • 一次处理所有累积的消息
  • 减少锁的获取次数
  • 提高吞吐量
  1. 解耦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
2
3
4
5
// 发送:向channel中放入数据
applyCh <- msg // 箭头指向channel

// 接收:从channel中取出数据
msg := <-applyCh // 箭头从channel指出来
在Raft中的角色架构图
1
2
3
4
5
┌─────────────┐          ┌─────────────┐          ┌─────────────┐
│ Raft层 │ │ Channel │ │ 应用层 │
│ (共识协议) │ ──────> │ (applyCh) │ ──────> │ (状态机) │
│ │ 发送 │ 缓冲/阻塞 │ 接收 │ │
└─────────────┘ └─────────────┘ └─────────────┘

数据流转

  1. Raft层(生产者):
    • 决定哪些日志可以提交(commitIdx)
    • 将提交的日志封装成ApplyMsg
    • 放入commitQueue队列
  2. Applier(中转者):
    • commitQueue取出消息
    • 通过applyCh <- msg发送给上层
  3. 应用层(消费者):
    • applyCh接收消息
    • 执行命令(如PUT/GET)
    • 更新状态机状态
为什么使用Channel?
  1. 解耦
1
2
// Raft不需要知道应用层如何执行命令
applyCh <- msg // 只负责发送,不管后续处理
  1. 顺序保证
1
2
3
4
// Channel保证消息按发送顺序接收
applyCh <- msg1 // 先发送
applyCh <- msg2 // 后发送
// 应用层先收到msg1,再收到msg2
  1. 线程安全
1
2
// 多个goroutine可以安全地发送/接收
// Channel内部处理同步
  1. 阻塞同步
1
2
3
4
5
// 发送阻塞:如果channel满了,发送者会等待
applyCh <- msg // 等待应用层接收

// 接收阻塞:如果channel为空,接收者会等待
msg := <-applyCh // 等待Raft发送
具体示例

创建Channel(在测试代码中)

1
2
3
applyCh := make(chan ApplyMsg)  // 无缓冲channel
// 或
applyCh := make(chan ApplyMsg, 100) // 有缓冲channel

消费者(应用层)

1
2
3
4
5
6
7
8
9
10
func (kv *KVServer) applier() {
for msg := range kv.applyCh { // 持续接收
if msg.CommandValid {
// 执行命令
kv.execute(msg.Command)
}
// 通知客户端结果
kv.notify(msg.CommandIndex)
}
}

生产者(Raft层)

1
2
// 在Applier goroutine中
applyCh <- msg // 发送给上层

KvRaft

调用时序图

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
客户端                    KVServer                   Raft层                  状态机
| | | |
|--PutAppend RPC--------->| | |
| | | |
| |--去重检查 | |
| | | |
| |--rf.Start(op)---------->| |
| | | |
| | |--追加日志 |
| | | |
| | |--复制到Follower |
| | | |
| | |--等待多数派确认 |
| | | |
| | |--提交日志 |
| | | |
| | |--applyCh-------------->|
| | | |
| | | |--执行PUT/APPEND
| | | |
| | | |--update去重表
| | | |
| | | |--replyChan
| | | |
| |<--rec_chan接收结果------| |
| | | |
|<--PutAppendReply-------| | |
| | | |