raft
# 1.简介
# 1.1.功能
etcd的raft模块实现了分布式一致性算法,基于Raft协议保证系统的数据一致性。作为etcd核心,raft模块独立于存储模块,提供节点选举、日志复制、心跳计时等诸多功能。--- leader选举 基于raft协议实现leader选举机制,分布式系统内只有选举为leader的节点才能处理请求 --- 日志复制 客户端提交请求时,leader本地记录请求日志,同时将记录复制到follower节点 --- 安全性保障 限制leader的commit操作和follower的vote操作,确保半数节点同意才会提交数据并改变状态 --- 节点变更 支持动态加入或删除节点,通过重新投票和重新分配副本实现新节点的加入和旧节点的删除 --- 日志压缩 etcd支持将旧日志压缩,减少存储时间1
2
3
4
5
6
7
8
9
10
11
12
13
14
# 1.2.模块
etcd中实现的raft本质上提供的是一个基于raft协议的sdk,通过sdk方式接入etcd解耦于storage模块。基于最小原则,raft模块设计时实现了基本功能,包括leader选举、日志处理、状态变更等逻辑,上层的网络传输和数据存储模块定义接口移交给应用层实现。--- 存储层 定义storage接口管理raft log,具体的持久化由应用层实现,提供基于内存数组的非持久化实现MemoryStorage --- 网络层 节点间的数据通信实现没有任何约束,通过channel和应用层交互,由应用层自定义实现处理接收到的消息 --- 模块实现 raft ├── confchange ├── ··· // 成员变更处理,实现joint consensus算法 ├── quorum ├── ··· // 集群多数投票处理,包括仲裁计算和投票结果判断 ├── raftpb ├── ··· // raft协议相关的protocol buffer消息(message、entry、snapshot) ├── tracker ├── ··· // 跟踪集群中所有节点状态,管理复制进度 ├── bootstrap.go // 初始化集群状态的工具函数 ├── doc.go // 文档和基础类型 ├── log.go // 日志存储接口storage,提供日志条目、快照和状态的读写方法 ├── log_unstable.go // 管理未提交日志条目,处理内存中的临时日志 ├── logger.go // 日志接口,允许应用层注入自定义的日志实现 ├── node.go // node接口,作为应用层与raft状态机的交互桥梁,处理消息循环和状态转换 ├── raft.go // 核心实现,涵盖选举、日志复制、状态转换 ├── rawnode.go // 提供底层raft节点接口,允许直接操作raft状态 ├── read_only.go // 实现线性一致读,基于readIndex或租约机制保证读一致性 ├── status.go // 提供raft状态统计信息,用于监控和调试 ├── storage.go // 抽象化存储层接口,支持持久化日志和快照,应用层实现,例如WAL模块 └── util.go // 工具函数,包括ID生成、类型转换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
# 1.3.宏观架构
etcd实现raft时,根据职责边界拆分为应用层和算法层两个模块。算法层作为大脑,负责根据共识机制进行请求内容的正确性校验、预写日志的状态的维护管理、可提交日志的进度推进。应用层扮演协调者的角色,是raft节点中承上启下的模块,对外与客户端交互,对内与各模块交互,串联raft节点运行的整体流程。节点启动时,应用层和算法层的处理作为两个独立的goroutine,依据几个channel进行模块间的异步通信。
注意
算法层的处理逻辑本质上类似
拼图,应用层提交的日志具备乱序、不一定正确的特征,算法层基于共识依次纠正碎片的顺序,纠正日志的内容,然后归还给应用层持久化。
# 2.算法数据
# 2.1.Entry
Entry是raft协议中的一笔预写日志,主要分为普通类型和配置变更两种类型。Entry中包含索引、任期、类型、数据,term+index共同构成Entry的全局唯一索引。// ./raft/raftpb/raft.pb.go type EntryType int32 const ( EntryNormal EntryType = 0 EntryConfChange EntryType = 1 EntryConfChangeV2 EntryType = 2 ) type Entry struct { Term uint64 // 任期 Index uint64 // 预写日志最后一条日志的索引 Type EntryType // 预写日志的类型 Data []byte // 预写日志的数据 XXX_unrecognized []byte // 未识别的协议缓冲区数据,用于protobuf兼容性处理 }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
# 2.2.Message
Message是raft协议中一条消息,用于驱动raft节点的状态变更。Message作为大而全的通用结构体,耦合了日志同步、leader选举、心跳等请求及响应,这种方式虽然简化了协议,却不可避免地破坏职责单一原则,每次作用时只会用到部分字段。// ./raft/raftpb/raft.pb.go type Message struct { Type MessageType // 消息类型 To uint64 // 接收消息的节点ID From uint64 // 发送消息的节点ID Term uint64 // 当前节点任期 LogTerm uint64 // 待同步日志的上一笔日志的任期 Index uint64 // 待同步日志的上一笔日志的索引 Entries []Entry // 待完成同步的预写日志 Commit uint64 // leader已提交的日志索引 Snapshot Snapshot // 快照 Reject bool // 标识响应结果为拒绝或赞同(响应投票或日志同步) RejectHint uint64 // 拒绝同步日志请求时返回的follower的最后一条日志索引,便于leader下一次补发缺失日志 Context []byte // 上下文信息,用于携带额外数据,例如集群配置 XXX_unrecognized []byte // 未识别的协议缓冲区数据,用于protobuf兼容性处理 }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16MessageType定义了raft协议支持的所有消息类型,根据宏观流程可以分为日志同步、leader选举、心跳广播、读请求四类。type MessageType int32 const ( MsgHup MessageType = 0 // 领导者发起选举消息 MsgBeat MessageType = 1 // 心跳,维持leader和follower联系 MsgProp MessageType = 2 // 向领导者发送提案(leader处理) MsgApp MessageType = 3 // 向跟随者发送的日志同步 MsgAppResp MessageType = 4 // 日志同步请求的响应 MsgVote MessageType = 5 // 投票消息 MsgVoteResp MessageType = 6 // 回复投票消息 ... MsgHeartbeat MessageType = 8 // 领导者向跟随着发送心跳,用于确认自己仍然活着,告诉跟随者自己任期号 MsgHeartbeatResp MessageType = 9 // 回复心跳消息 MsgUnreachable MessageType = 10 // 通知领导者某个跟随者不可达 ... MsgCheckQuorum MessageType = 12 // 检查当前集群是否有活跃成员 MsgTransferLeader MessageType = 13 // 请求领导权转移 MsgTimeoutNow MessageType = 14 // 转化为候选者请求 MsgReadIndex MessageType = 15 // 客户端发起的只读请求,用于读取最新的状态 MsgReadIndexResp MessageType = 16 // 回复只读请求 MsgPreVote MessageType = 17 // 预投票,类似探测,用于防止网络分区影响无限自增term MsgPreVoteResp MessageType = 18 // 预投票响应 )1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
# 2.3.raftLog
预写日志产生时,需要经历未持久化——>已持久化的过程,前者在
raftLog.unstable完成,后者在应用层完成,通过raftLog.storage为算法层提供已持久化预写日志的查询能力。// ./raft/log.go type raftLog struct { storage Storage // 保存最后一次快照后提交的数据 unstable unstable // 保存未持久化的数据和快照,最终转存到storage committed uint64 // 已提交的日志索引 applied uint64 // 已应用到状态机的日志索引(已提交-->应用) ... }1
2
3
4
5
6
7
8
9
10
11
unstable提供未持久化预写日志代理能力,可读可写,entries是还未进行持久化的预写日志列表,offset是首笔预写日志在持久化预写日志中的偏移量。storage是持久化日志存储接口,提供日志的查询能力,作为抽象接口由上层应用自定义实现。// 提供未持久化预写日志代理能力,可读可写 type unstable struct { snapshot *pb.Snapshot // 未持久化的快照 entries []pb.Entry // 还未持久化的预写日志 offset uint64 // 预先日志列表中起始日志相对持久化预写日志的偏移量 ... } // 持久化日志存储接口,提供日志的查询能力,storage定义抽象接口,由应用层实现 type Storage interface { InitialState() (pb.HardState, pb.ConfState, error) // 返回保存的初始化状态 Entries(lo, hi, maxSize uint64) ([]pb.Entry, error) // 返回索引范围[lo,hi)内不大于maxSize的数据 Term(i uint64) (uint64, error) // 根据索引获取对应任期 LastIndex() (uint64, error) // 获取最后一条数据的索引 FirstIndex() (uint64, error) // 获取第一条数据的索引 Snapshot() (pb.Snapshot, error) // 获取最近快照 }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
# 2.4.Ready
Ready是算法层和应用层交互的数据格式,每当算法层执行完一轮逻辑处理会向readyc channel传入一个Ready结构体,其中封装了算法层处理好的结果,应用层可以直接消费。./raft/node.go type Ready struct { *SoftState // 软状态,集群的leader和当前节点的状态 pb.HardState // 硬状态,节点当前的Term、Vote和Commit ReadStates []ReadState // 应用索引大于ReadStates中索引时本地服务线性化读请求 Entries []pb.Entry // 待持久化的预写日志 Snapshot pb.Snapshot // 需要持久化的快照,启用异步存储写入时不需要立即处理 CommittedEntries []pb.Entry // 本轮算法层已经提交的预写日志,由上层应用到状态机(已提交待应用的预写日志) Messages []pb.Message // 本轮算法层需要发送的消息,由应用层调用通信模块发送 MustSync bool // 限制hardState和entries是否必须持久化到磁盘 }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18raft节点的软状态包含节点的Leader和当前节点的角色状态,由于容易发生改变且通过通信可重新恢复,因此该部分信息无需进行持久化存储。硬状态记录的是节点的当前任期、投票归属和已提交的日志索引,这些信息丢失后无法恢复,因此需要持久化存储。// 节点角色可划分为follower、candidate和leader // candidate由竞选和预竞选情况可以拆分为candidate和preCandidate type StateType uint64 const ( StateFollower StateType = iota StateCandidate StateLeader StatePreCandidate ) // raft节点的软状态 type SoftState struct { Lead uint64 // 集群leader RaftState StateType // 节点当前状态 } // raft节点的硬状态 type HardState struct { Term // 当前任期 Vote // 投票归属 Commit // 已提交的日志索引 XXX_unrecognized []byte // 未识别的protobuf数据,用于兼容性处理 }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
# 2.5.Node
Node定义raft节点的基本操作,是算法层raft节点的抽象,也是应用层和算法层交互的唯一入口,应用层持有Node作为算法层raft节点引用,通过Node API完成和算法层的channel交互。type Node interface { Tick() // 驱动节点内部逻辑时钟向前一步,用于计算选举超时和心跳超时 Campaign(ctx context.Context) error // 转换节点为候选人状态及开始竞选成为leader Propose(ctx context.Context, data []byte) error // 用于向日志追加数据,可能被丢弃或失败,上层需要保证重试 ProposeConfChange(ctx context.Context, cc pb.ConfChangeI) error // 提议一个配置变更,可能被丢弃或失败 Step(ctx context.Context, msg pb.Message) error // 接收及处理其他节点发送的消息,比如投票、心跳 Ready() <-chan Ready // 返回通道,用于和上层交互节点状态 Advance() // 通知节点应用处理完上一个ready返回的状态,准确获取下一个可用状态 ApplyConfChange(cc pb.ConfChangeI) *pb.ConfState // 提议的配置变更应用到节点 TransferLeadership(ctx context.Context, lead, transferee uint64) // leader权力移交给其他节点 ReadIndex(ctx context.Context, rctx []byte) error // 请求读取 Status() Status // 返回当前节点状态信息 ReportUnreachable(id uint64) // 报告某个节点最近的通信中不可达 ReportSnapshot(id uint64, status SnapshotStatus) // 报告快照的发送状态,包括发送快照的节点ID和发送状态 Stop() // 终止节点运行 }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
29
Tick会传送定时驱动信号,每次调用Tick方法的时间间隔都是固定的,称为一个tick,是raft节点的最小计时单位,leader节点的心跳计时和leader/candidate的选举计时都以tick作为单位。实现上,应用层向node.tickc channel中传入一个信号,算法层时刻监听channel状态,获取到信号后会驱动raft节点进行定时函数处理。func (n *node) Tick() { select { case n.tickc <- struct{}{}: ... } }1
2
3
4
5
6Propose向日志追加数据,应用层调用Node.Propose向算法层发起一笔写数据请求。调用过程中,会向node.proc channel中传入一条类型为MspProp的消息,算法层获取消息后会驱动raft节点进入写请求提议流程。func (n *node) Propose(ctx context.Context, data []byte) error { return n.stepWait(ctx, pb.Message{Type: pb.MsgProp, Entries: []pb.Entry{{Data: data}}}) } func (n *node) stepWait(ctx context.Context, m pb.Message) error { return n.stepWithWaitOption(ctx, m, true) } func (n *node) stepWithWaitOption(ctx context.Context, m pb.Message, wait bool) error { ... ch := n.propc pm := msgWithResult{m: m} ... select { // mspProp消息写入proc channel case ch <- pm: ... } ... return nil }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21ProposeConfigChange用于配置变更,应用层向node.proc channel发送消息类型为MspProp、日志类型为EntryConfChange的消息,推动raft节点进入配置变更流程。func confChangeToMsg(c pb.ConfChangeI) (pb.Message, error) { // EntryConfChange日志类型 typ, data, err := pb.MarshalConfChange(c) if err != nil { return pb.Message{}, err } return pb.Message{Type: pb.MsgProp, Entries: []pb.Entry{{Type: typ, Data: data}}}, nil } func (n *node) ProposeConfChange(ctx context.Context, cc pb.ConfChangeI) error { msg, err := confChangeToMsg(cc) if err != nil { return err } return n.Step(ctx, msg) }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16ApplyConfChange用于提交配置变更,应用层通过node.confc向算法层发送配置变更内容,算法层接收到后直接应用,确保变更即时生效。func (n *node) ApplyConfChange(cc pb.ConfChangeI) *pb.ConfState { var cs pb.ConfState select { // 向配置更新channel写入需更新的配置 case n.confc <- cc.AsV2(): ... } ... return &cs }1
2
3
4
5
6
7
8
9
10ReadIndex用于发起读请求,应用层通过node.revc channel传入类型为MsgReadIndex的消息向算法层发起读数据请求。func (n *node) ReadIndex(ctx context.Context, rctx []byte) error { // 向node.revc channel发送读请求消息 return n.step(ctx, pb.Message{Type: pb.MsgReadIndex, Entries: []pb.Entry{{Data: rctx}}}) }1
2
3
4Ready和Advance方法用于应用层和算法层的状态交换,应用层调用Node.Ready方法返回node.readyc channel用于监听算法层处理结果,应用层调用Node.Advance通过node.advancec channel向算法层发送信号,驱动算法层进入下一轮调度循环。因此,Node.Ready和Node.Advance方法是成对使用的,两个方法每配合一次意味着应用层和算法层完成一轮交互。func (n *node) Ready() <-chan Ready { return n.readyc } func (n *node) Advance() { select { case n.advancec <- struct{}{}: case <-n.done: } }1
2
3
4
5
6
7
8
node是Node接口的算法层实现,定义了一系列channel传递消息,通过监听对应channel主席那个不同处理。// Node实现 ./raft/read_only.go type node struct { propc chan msgWithResult // 接受客户端请求返回结果的管道 recvc chan pb.Message // 接收其他节点raft消息的管道 confc chan pb.ConfChangeV2 // 用于接收配置变更消息的管道 confstatec chan pb.ConfState // 接收当前配置状态的管道 readyc chan Ready // 接收raft日志提交信号的管道 advancec chan struct{} // 接收驱动raft算法状态机的信号管道 tickc chan struct{} // 接收逻辑时钟的信号管道 done chan struct{} // 接收完成raft操作的信号管道 stop chan struct{} // 接收节点停止信号的管道 status chan chan Status // 接收节点状态信息的管道 rn *RawNode // raft节点 }1
2
3
4
5
6
7
8
9
10
11
12
13
14
# 2.6.readOnly
readOnly是挂起的读请求队列,由raft节点持有,pendingReadIndex是一系列还未处理的读请求,通过map的形式存储,建立读请求ID与读请求之间的映射关系。readIndexQueue本质上是一个队列,根据读请求到达的时间顺序先进先出。// ./raft/read_only.go type readOnly struct { option ReadOnlyOption // 读请求可选项 pendingReadIndex map[string]*readIndexStatus // entry的数据为key,保存当前待处理的读请求 readIndexQueue []string // 有序存放读请求的ID } type readIndexStatus struct { req pb.Message // 保存原始的readIndex请求消息 index uint64 // 保存收到读请求时的leader commit索引 acks map[uint64]bool // 保存已应答的节点,用于半数节点通过判断 }1
2
3
4
5
6
7
8
9
10
11
12
# 2.7.Progress
Progress是leader记录其他节点日志同步进度的载体,Match是该节点已同步日志的索引,Next是leader下一次日志同步的索引位置。// ./raft/progress.go type Progress struct { Match, Next uint64 // 已成功复制到跟随者的最高日志索引,发送到跟随者的下一条日志索引 State StateType // 定义领导者与跟随者交互的状态枚举,包含探测、复制、快照 PendingSnapshot uint64 // 快照交互下,待处理的快照索引号,此时复制过程将暂停 RecentActive bool // 最近活跃状态,选举超时后重置为false ProbeSent bool // 探测状态下,为真时raft暂停向该节点发送复制消息,直到重置为false Inflights *Inflights // 跟踪正在传输的消息,限制在途消息的最大数量、每个Progress可用带宽,用于限制流量 IsLearner bool // 表示该进度是否跟踪的学习者节点(非投票成员) }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
# 2.8.raft
raft类是共识算法的抽象,几乎囊括一个raft节点正常运行时必备的所有信息,包括节点任期、投票给谁、管理预写日志、日志同步进度等。// ./raft/raft.go type raft struct { id uint64 Term uint64 // 任期号 Vote uint64 // 投票给哪个节点 readStates []ReadState // 读请求状态 raftLog *raftLog // 日志管理模块 maxMsgSize uint64 // 最大消息大小(字节) maxUncommittedSize uint64 // 最大未提交的数据大小(字节) prs tracker.ProgressTracker // 进度集合,跟踪其他节点的日志同步进度 state StateType // 当前节点角色状态 ... msgs []pb.Message // 当前节点需要发送的消息,算法层会将msgs内容放入Ready发送到readyc channel lead uint64 // 节点的Leader ID leadTransferee uint64 // 记录正在进行领导权移交的目标节点ID pendingConfIndex uint64 // 记录已被提交但未应用到状态机的配置变更的索引位置 uncommittedSize uint64 // 等待提交的未提交的数据大小估计值 readOnly *readOnly // 挂起的读请求列表,等待leader认证身份后响应给客户端ack electionElapsed int // 选举计时,单位为tick heartbeatElapsed int // 心跳计时,单位为tick checkQuorum bool // 是否开启检查集群内节点数量超过半数功能 preVote bool // 是否在选举前进行类型为PreVote的消息投票 heartbeatTimeout int // Leader广播心跳的时间间隔,单位为tick electionTimeout int // candidate/follower进行选举的时间间隔,单位为tick randomizedElectionTimeout int // 随机扰动后的选举时间间隔,范围在[eleTimeout, 2*eleTimeout-1],切换节点重置 disableProposalForwarding bool // 是否关闭转发客户端发来的写请求给其他Leader节点功能 tick func() // 定时器函数 step stepFunc // 节点的状态机处理函数,不同角色的函数实现不同 ... pendingReadIndexMessages []pb.Message // 等待提交到状态机的MsgReadIndex类型的消息列表 }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
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
# 2.9.raftNode
raftNode是在应用层中额外封装一层的raft节点,除持有算法层入口外,还包含与客户端、数据状态机、预写日志持久化模块、通信模块交互的一些能力。// ./contrib/raftexample/raft.go type raftNode struct { proposeC <-chan string // 用于接收客户端发送的写请求 confChangeC <-chan raftpb.ConfChange // 用于接收客户端发送的配置变更请求 commitC chan<- *string // 用于将已提交的日志应用到数据状态机 ... id int // raft节点ID peers []string // 同一集群内其他raft节点的标识信息 ... appliedIndex uint64 // 本节点已应用到状态机的日志索引 node raft.Node // 算法层入口(应用层持有的算法层引用) raftStorage *raft.MemoryStorage // 持久化预写日志存储模块 wal *wal.WAL // 预写日志 ... transport *rafthttp.Transport // raft集群通信模块 ... }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
# 3.2.kvstore
kvStore是一个乞丐版的键值对存储模块,用于以小见大,模拟还原etcd存储系统与raft节点间的交互模式。// ./contrib/raftexample/kvstore.go type kvstore struct { proposeC chan<- string // raftNode.proposeC共用channel,负责向raftNode发送来自客户端的写请求 mu sync.RWMutex // 保护存储介质访问的数据并发安全 kvStore map[string]string // 键值对存储介质,可以理解为数据状态机 ... }1
2
3
4
5
6
7kvstore模块启动时,会注入一个commitC channel,和raft commicC是同一个channel。func newKVStore(snapshotter *snap.Snapshotter, proposeC chan<- string, commitC <-chan *string, errorC <-chan error) *kvstore { s := &kvstore{proposeC: proposeC, kvStore: make(map[string]string), snapshotter: snapshotter} // 同步raftNode提交的日志应用到数据状态机 s.readCommits(commitC, errorC) // 异步监听raftNode commitC提交的日志应用到数据状态机 go s.readCommits(commitC, errorC) return s }1
2
3
4
5
6
7
8kvstore会持续监听commitC channel,获取到raftNode提交的日志数据,将其应用到数据状态机。func (s *kvstore) readCommits(commitC <-chan *string, errorC <-chan error) { for data := range commitC { ... var dataKv kv dec := gob.NewDecoder(bytes.NewBufferString(*data)) // 日志数据解析 if err := dec.Decode(&dataKv); err != nil { log.Fatalf("raftexample: could not decode message (%v)", err) } s.mu.Lock() // 更新到map映射 s.kvStore[dataKv.Key] = dataKv.Val s.mu.Unlock() } ... }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16此外,
kvstore提供Propose方法供客户端发送写请求,kvstore随之通过proposeC将请求发送给raftNode处理。func (s *kvstore) Propose(k string, v string) { var buf bytes.Buffer if err := gob.NewEncoder(&buf).Encode(kv{k, v}); err != nil { log.Fatal(err) } s.proposeC <- buf.String() }1
2
3
4
5
6
7
# 3.httpapi
httpapi是一个http Server,可以理解为低配版的客户端模块。服务运行时,收到的PUT请求会被视为写请求,收到的POST请求会被视为添加节点配置变更请求,收到的DELETE请求会被视为删除节点的集群配置变更请求,写请求会通过kvstore.Propose api进行转发处理。type httpKVAPI struct { store *kvstore confChangeC chan<- raftpb.ConfChange } // 接收用户发送请求并转发 func (h *httpKVAPI) ServeHTTP(w http.ResponseWriter, r *http.Request) { key := r.RequestURI defer r.Body.Close() switch { case r.Method == "PUT": v, err := ioutil.ReadAll(r.Body) ... // 写请求转发 h.store.Propose(key, string(v)) ... case r.Method == "GET": // 处理查询 if v, ok := h.store.Lookup(key); ok { w.Write([]byte(v)) } ... case r.Method == "POST": url, err := ioutil.ReadAll(r.Body) ... // 节点配置变更请求 cc := raftpb.ConfChange{ Type: raftpb.ConfChangeAddNode, NodeID: nodeId, Context: url, } h.confChangeC <- cc ... case r.Method == "DELETE": ... // 移除节点集群配置变更请求 cc := raftpb.ConfChange{ Type: raftpb.ConfChangeRemoveNode, NodeID: nodeId, } h.confChangeC <- cc ... } }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
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
# 4.宏观流程
# 4.1.应用层运行
应用层封装的
raft节点定义为raftNode,raftNode.startRaft被调用后,一个raft节点真正启用,其中包含了应用层和算法层两个部分的初始化和启动过程。func (rc *raftNode) startRaft() { ... // 获取集群内其他raft节点消息 rpeers := make([]raft.Peer, len(rc.peers)) for i := range rpeers { rpeers[i] = raft.Peer{ID: uint64(i + 1)} } // 创建一份raft节点配置,leader的心跳间隔默认为tick,选举时间默认10个tick c := &raft.Config{ ID: uint64(rc.id), ElectionTick: 10, HeartbeatTick: 1, Storage: rc.raftStorage, MaxSizePerMsg: 1024 * 1024, MaxInflightMsgs: 256, MaxUncommittedEntriesSize: 1 << 30, } // 已经执行过日志预写,重启算法层Node if oldwal { rc.node = raft.RestartNode(c) // 未执行过日志预写(新节点),启动算法层的Node } else { startPeers := rpeers if rc.join { startPeers = nil } rc.node = raft.StartNode(c, startPeers) } ... // 启动通信模块 rc.transport.Start() ... go rc.serveRaft() // 异步开启raftNode主循环,用于与算法层的协程建立持续通信的关系 go rc.serveChannels() }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
29
30
31
32
33
34
35
36
37
# 4.2.算法层运行
应用层启动时会初始化一个
raft共识机制的抽象节点结构,统一设置角色、任期,还会初始化集群管理信息,包括其他节点的日志同步进度、非持久化预写日志状态等。func StartNode(c *Config, peers []Peer) Node { ... // 初始化rawNode(根据持久化配置构造当前节点属性) rn, err := NewRawNode(c) ... // 应用节点配置、重置角色、初始化已提交日志索引 rn.Bootstrap(peers) // 共识机制的节点抽象 n := newNode(rn) go n.run() return &n }1
2
3
4
5
6
7
8
9
10
11
12
13
14run方法用于节点启动,方法内部是一个无限循环,通过select语句监听多个通道事件,以便响应不同类型消息。另外,内部基于channel指针替换,保证select多路复用模式下Node.Ready和Node.Advance方法被成对调用。func (n *node) run() { var propc chan msgWithResult var readyc chan Ready var advancec chan struct{} var rd Ready r := n.rn.raft lead := None for { if advancec != nil { // advance channel不为空,说明还在等应用调用Advance接口通知算法层处理完毕 readyc = nil } else if n.rn.HasReady() { // 算法层处理好结果,重置ready channel同时初始化一个Ready消息(msgs内容补充到Ready结构体) rd = n.rn.readyWithoutAccept() readyc = n.readyc } ... // 当前节点是Leader,初始化proc channel接收提案 propc = n.propc ... select { case pm := <-propc: // 处理本地收到的提交值 m := pm.m m.From = r.id err := r.Step(m) ... case m := <-n.recvc: // 处理其他节点发送的提交值(过滤未知节点) if pr := r.prs.Progress[m.From]; pr != nil || !IsResponseMsg(m.Type) { r.Step(m) } case cc := <-n.confc: _, okBefore := r.prs.Progress[r.id] // 处理集群配置变更 cs := r.applyConfChange(cc) ... // 当前节点被移除,propc置为nil,停止接收提案 ... select { // 配置变更移交给算法层 case n.confstatec <- cs: case <-n.done: } case <-n.tickc: // 触发raft的定时逻辑(选举超时、心跳) n.rn.Tick() case readyc <- rd: ... // 等待应用处理完成通知 advancec = n.advancec case <-advancec: // 应用处理完成,重置Ready消息,advance置空,等待下一轮处理 n.rn.Advance(rd) rd = Ready{} advancec = nil ... } } }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
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63流程
1.通过
channel指针替换,保证select多路复用模式下,Node.Ready和Node.Advance成对调用2.保证算法层没有新结果产生时,不会通过
readyc向应用层提交Ready消息,避免流程空转3.接收来自应用层的消息时步入
raft.step方法4.接收已确认可应用的配置变更消息时,对集群配置发起变更
5.接收到应用层的定时
Tick调用时,根据raft节点角色,使用其对应的tick函数处理6.算法层有处理结果后投递到
readyc供应用层raftNode接收更新状态变量,应用层处理完成后向advancec中发送信号量,算法层更新applied index,开启新一轮循环
# 4.3.请求处理
应用层启动时,会异步调用
serveChannels方法处理应用层与算法层的请求交互。该方法可以拆解为并发运行的两部分,一部分负责将应用层接收请求转至算法层处理,另一部分负责从算法层获取执行结果,func (rc *raftNode) serveChannels() { ... // 定期器 ticker := time.NewTicker(100 * time.Millisecond) defer ticker.Stop() // send proposals over raft ... // event loop on raft state machine updates ... }1
2
3
4
5
6
7
8
9
10
11
12面向外部,
raftNode会持续监听proposeC和confChangeC两个channel,获取来自客户端的写请求和配置变更请求,然后调用Node API将请求发送给算法层。go func() { ... // 这里的消息都会在node.Step消费处理 for rc.proposeC != nil && rc.confChangeC != nil { select { // 转发写请求至算法层 case prop, ok := <-rc.proposeC: ... // 调用node的Propose函数把写请求发送给算法层node处理 rc.node.Propose(context.TODO(), []byte(prop)) // 转发配置变更请求至算法层 case cc, ok := <-rc.confChangeC: confChangeCount++ cc.ID = confChangeCount rc.node.ProposeConfChange(context.TODO(), cc) } } // client closed channel; shutdown raft if not already close(rc.stopc) }()1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20面向内部,
raftNode会启动一个定时器,每个tick默认为100ms,然后定期调用Node.Tick方法驱动算法层执行定时函数。此外,raftNode通过Node.Ready和Node.Advance方法与算法层循环交互for { select { // 定时驱动算法层函数执行 case <-ticker.C: rc.node.Tick() // 算法层反馈处理结果驱动应用层持久化 case rd := <-rc.node.Ready(): // 记录预写日志 rc.wal.Save(rd.HardState, rd.Entries) ... // 非持久化预写日志追加 rc.raftStorage.Append(rd.Entries) // 基于通信模块为算法层执行消息发送动作 rc.transport.Send(rd.Messages) // 与数据状态机交互,应用算法层已提交的预写日志 if ok := rc.publishEntries(rc.entriesToApply(rd.CommittedEntries)); !ok { rc.stop() return } ... // 响应算法层,完成一轮交互 rc.node.Advance() ... } }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
26raftNode.publishEntries用于将算法层已提交的日志应用到数据状态机,已提交日志属于写请求时,通过commitC将数据转给kvstore,通过kvstore的消费协程将其真正写入数据状态机;已提交日志属于配置变更请求时,需要调用Node.ApplyConfChange将其真正作用于算法层。func (rc *raftNode) publishEntries(ents []raftpb.Entry) bool { for i := range ents { switch ents[i].Type { case raftpb.EntryNormal: ... s := string(ents[i].Data) select { case rc.commitC <- &s: ... } case raftpb.EntryConfChange: ... rc.confState = *rc.node.ApplyConfChange(cc) ... } // 更新已应用索引 rc.appliedIndex = ents[i].Index ... } return true }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
# 4.4.定时调用
应用层
raftNode在raftNode.serveChannels方法调用Node.Tick方法,向tickc channel中投递消息。func (rc *raftNode) serveChannels() { ... // event loop on raft state machine updates for { select { case <-ticker.C: rc.node.Tick() ... } } }1
2
3
4
5
6
7
8
9
10
11算法层协程在
node.run方法中消费到消息,然后根据当前节点角色切换到对应的定时函数raft.tick中进行定制处理。func (n *node) run() { ... for { ... select { ... case <-n.tickc: n.rn.Tick() ... } } }1
2
3
4
5
6
7
8
9
10
11
12算法层还会组装消息,将消息封装到
Ready结构体中发送给应用层,组装消息的具体方法为raft.send(),针对拉票、读写请求之外的消息填充任期,消息最终会放在raft.msgs。func (r *raft) send(m pb.Message) { m.From = r.id ... if m.Type != pb.MsgProp && m.Type != pb.MsgReadIndex { m.Term = r.Term } r.msgs = append(r.msgs, m) }1
2
3
4
5
6
7
8算法层在与应用层交互时,会有一个基于
raft类构造Ready结构体的过程,刚才的msgs会填充到Ready结构体。func newReady(r *raft, prevSoftSt *SoftState, prevHardSt pb.HardState) Ready { rd := Ready{ Entries: r.raftLog.unstableEntries(), // 未持久化的数据 CommittedEntries: r.raftLog.nextEnts(), // 已提交但未应用的数据 Messages: r.msgs, // 待发送的消息 } ... return rd }1
2
3
4
5
6
7
8
9每一轮算法层投递完
Ready后,会把raft.msgs置空,保证消息不被重复发送到应用层。func (n *node) run() { ... for { if advancec != nil { readyc = nil } ... select { ... case readyc <- rd: n.rn.acceptReady(rd) advancec = n.advancec ... } } }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
# 5.角色切换流程
# 5.1.简介
- 所谓的节点角色切换,除
raft数据结构进行一些状态信息更新外,还有定时函数raft.tick和状态机函数raft.step进行切换。对于leader而言,定时函数中的任务是要向集群中其他节点广播心跳,对于follower和candidate而言,定时函数的任务是发起竞选。当然,不同角色接收到消息时的处理模式也不尽相同,具体的区别会体现在角色专属的状态机函数。 
# 5.2.Step
算法层运行期间会处理来自应用层的请求时,其中的写请求和其他节点请求处理会先步入
raft.Step方法,进行通用的消息前处理,后续才会根据角色进入定制的状态机函数。func (r *raft) Step(m pb.Message) error { // Handle the message term, which may result in our stepping down to a follower. switch { case m.Term == 0: // 本地消息 case m.Term > r.Term: // 收到更高任期的投票消息 if m.Type == pb.MsgVote || m.Type == pb.MsgPreVote { force := bytes.Equal(m.Context, []byte(campaignTransfer)) inLease := r.checkQuorum && r.lead != None && r.electionElapsed < r.electionTimeout // 检查领导者租约,当前仍在任期内(收到合法Leader心跳不久) // 非Leader权限转移 // 防止网络分区中的节点干扰正常集群 if !force && inLease { return nil } } switch { case m.Type == pb.MsgPreVote: // 预投票篇不更新本地term和角色(试探性请求) case m.Type == pb.MsgPreVoteResp && !m.Reject: // 预投票成功,暂不处理(等待正式投票阶段) default: // 其他高任期消息(心跳/日志等),退位为follower if m.Type == pb.MsgApp || m.Type == pb.MsgHeartbeat || m.Type == pb.MsgSnap { r.becomeFollower(m.Term, m.From) // 明确知道新Leader是谁 } else { r.becomeFollower(m.Term, None) // 仅知道term更新,Leader未知 } } case m.Term < r.Term: // 收到低任期消息 // 1.消息的网络延迟导致的过期消息 // 2.脑裂导致当前节点递增term // 3.配置变更期间产生的混乱 if (r.checkQuorum || r.preVote) && (m.Type == pb.MsgHeartbeat || m.Type == pb.MsgApp) { // 旧Leader的心跳、日志请求,强制回复高term使其退位 // 当前节点term大于消息term时,收到旧的leader消息 // 这种不能忽略,否则旧leader认为自己有效,导致双主问题 // 此时响应不带数据的MsgAppResp消息,旧leader收到更高term的响应后主动退位,避免双主问题 r.send(pb.Message{To: m.From, Type: pb.MsgAppResp}) } else if m.Type == pb.MsgPreVote { // 明确拒绝低任期预投票 // PreVote机制引入前,直接拒绝可能导致死锁,引入后显示拒绝以维持协议的正确性 r.send(pb.Message{To: m.From, Term: r.Term, Type: pb.MsgPreVoteResp, Reject: true}) } else { // 其他类型的低term消息不处理 } return nil } // 能走到这里代表消息的任期大等于本节点 switch m.Type { case pb.MsgHup: // 收到HUP消息,说明准备进行选举,此时如果当前节点不是leader if r.state != StateLeader { // 当前节点不存在于集群配置或属于学习者,不能参与选举,保障协议安全性 // 假设5节点集群,节点A被移除(配置变更提交),节点A由于网络分区尚未收到新配置,节点A超时后尝试发起选举,检查自身后不参与 if !r.promotable() { return nil } // 取出[applied+1,committed+1]之间的消息,即已提交但未应用的配置列表 ents, err := r.raftLog.slice(r.raftLog.applied+1, r.raftLog.committed+1, noLimit) ... // 检查是否存在未应用的配置变更 if n := numOfPendingConf(ents); n != 0 && r.raftLog.committed > r.raftLog.applied { // 避免配置变更期间的选举干扰 return nil } // 根据配置选举凸投票或直接选举 if r.preVote { r.campaign(campaignPreElection) } else { r.campaign(campaignElection) } } else { // 本身就是leader,忽略选举消息 } case pb.MsgVote, pb.MsgPreVote: // 收到投票类消息 canVote := r.Vote == m.From || // 已投票给该候选者 (r.Vote == None && r.lead == None) || // 未投票且当前无leader (m.Type == pb.MsgPreVote && m.Term > r.Term) // 更高任期预投票 // 可以投票,同时候选者的日志更新 if canVote && r.raftLog.isUpToDate(m.Index, m.LogTerm) { // 日志更新的候选者请求投票,接受并投给它 r.send(pb.Message{To: m.From, Term: m.Term, Type: voteRespMsgType(m.Type)}) if m.Type == pb.MsgVote { // 如果是正式的选举投票 // 重置竞选超时器 r.electionElapsed = 0 // 保存给哪个节点投票 r.Vote = m.From } } else { // 拒绝投票(日志过旧或已投其他节点) r.send(pb.Message{To: m.From, Term: r.Term, Type: voteRespMsgType(m.Type), Reject: true}) } default: // 其他情况进入各自状态下的定制的状态机函数 err := r.step(r, m) ... } return nil }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
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
# 5.3.becomeLeader
节点需要切换至
leader身份时,会调用raft.becomeLeader方法,进行节点状态切换、节点状态信息重置、未应用配置变更处理,同时会基于当前任期提交一笔空日志,避免出现提交仍回滚问题。func (r *raft) becomeLeader() { ... // 切换状态机函数(后续的心跳、日志消息处理) r.step = stepLeader // 重置节点状态信息 r.reset(r.Term) // 定时函数切换为心跳 r.tick = r.tickHeartbeat // 切换节点状态,标记leader身份 r.lead = r.id r.state = StateLeader // leader自身的进度设置为Replicate模式(leader不需要自身复制日志) r.prs.Progress[r.id].BecomeReplicate() // 配置变更保护,保守地将未完成地配置变更索引更新为最后未知,阻止新提案直到旧日志都提交 // 防止leader刚上任时因未提交的配置变更导致集群分裂 r.pendingConfIndex = r.raftLog.lastIndex() // 基于当前任期提交一条空日志,避免出现提交仍回滚问题 emptyEnt := pb.Entry{Data: nil} if !r.appendEntry(emptyEnt) {...} // 配额特殊处理,空日志不计入未提交日志配额 r.reduceUncommittedSize([]pb.Entry{emptyEnt}) }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22raft.reset()会批量重置节点状态信息,包括term、lead、心跳计时、选举计时、竞选票箱、读请求队列、其他节点日志进度等内容。其中,每次reset都会对选举时间间隔添加随机扰动,产生一个新的随机值。func (r *raft) reset(term uint64) { if r.Term != term { // 更新当前任期 r.Term = term // 清空投票记录 r.Vote = None } r.lead = None r.electionElapsed = 0 r.heartbeatElapsed = 0 // 重置选举超时(增加随机扰动,避免集群同时发起选举) r.resetRandomizedElectionTimeout() // 节点正在进行leader转移时被重置,终止转移流程 r.abortLeaderTransfer() // 清除所有节点的投票记录(用于选举统计) r.prs.ResetVotes() // 重置进度,确保新leader从最新日志开始复制 r.prs.Visit(func(id uint64, pr *tracker.Progress) { *pr = tracker.Progress{ Match: 0, // 已匹配的日志所有 Next: r.raftLog.lastIndex() + 1, // 下一个待发送索引 Inflights: tracker.NewInflights(r.prs.MaxInflight), IsLearner: pr.IsLearner, // 保持learner状态 } // leader自身日志完全匹配 if id == r.id { pr.Match = r.raftLog.lastIndex() } }) // 重置未完成的配置变更索引,此时无待处理的配置变更,防止新leader未应用旧配置变更时发起新的变更 r.pendingConfIndex = 0 // 日志配额控制,限制未提交的日志总量,允许新leader重新积累配额 r.uncommittedSize = 0 // 新建readOnly实例确保旧任期的只读请求不会影响新任期,实现线性一致读 r.readOnly = newReadOnly(r.readOnly.option) }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
29
30
31
32
33
34
35
36
37
38
39leader上任追加空日志
1.
raft规定新leader必须提交一条自己任期内的日志条目,空日志本质上作为新leader的就职宣言,向集群宣告其领导权有效。只有成功复制该日志后,leader才能开始处理客户端请求。2.其实,空日志发挥防御作用,
leader未提交空日志宣示主权时,崩溃选举的新leader可能丢失部分旧日志,违反leader的完备性原则
# 5.4.becomeFollower
节点角色需要变更为
follower时,会调用raft.becomeFollower方法,重置定时处理函数、状态机函数、角色、任期等信息。func (r *raft) becomeFollower(term uint64, lead uint64) { // 重置状态机函数(处理心跳、日志、快照) r.step = stepFollower // 重置任期、投票记录等 r.reset(term) // 重置定时处理函数(定时发起选举) r.tick = r.tickElection // 标记leader和自身角色 r.lead = lead r.state = StateFollower }1
2
3
4
5
6
7
8
9
10
11
# 5.5.becomePreCandidate
节点需要发起预选举时,需要变更为
preCandidate角色,会调用raft.becomePreCandidate方法切换节点的定时处理函状态机函数。其中,预选举节点节点的term不会更新。func (r *raft) becomePreCandidate() { ... // 更新状态机函数,preVote不会递增term,也不会先进行投票,直至preVote结果出来再进行决定 r.step = stepCandidate // 清空投票记录,为新一轮预投票做准备 r.prs.ResetVotes() // 设置选举定时器 r.tick = r.tickElection // 重置leader及角色 r.lead = None r.state = StatePreCandidate }1
2
3
4
5
6
7
8
9
10
11
12预投票
1.
预投票属于探测性阶段,仅仅是信息收集,不会改变集群任何状态,此时更新term,网络分区中的节点任期会不断自增,造成污染2.
预投票可以筛选不必选举的节点——Leader不存在、节点日志不够新、无法获得多数节点响应,同时防止网络分区节点干扰,预投票探测自己有机会获胜才会触发选举3.
预投票与投票属于两阶段确认,确保有资格的节点才会成为Leader
# 5.6.becomeCandidate
节点角色切换为
candidate时调用raft.becomeCandidate(),此时会将节点的定时处理函数重置为驱动竞选函数tickElection、状态机函数切换为stepCandidate,同时更新节点的Leader和state信息。func (r *raft) becomeCandidate() { // 更新状态机函数 r.step = stepCandidate // 增加任期 r.reset(r.Term + 1) // 切换定时函数为驱动选举 r.tick = r.tickElection // 给自己投票 r.Vote = r.id // 更新状态信息 r.state = StateCandidate }1
2
3
4
5
6
7
8
9
10
11
12注意
与
预选举不同,正式的选举流程中需要对term自增
# 5.7.leader->follower
leader处理请求时,会先走进通用的状态机函数raft.Step中,此时发现消息任期更大,会自动退回follower。func (r *raft) Step(m pb.Message) error { switch { case m.Term == 0: // 本地消息 case m.Term > r.Term: ... switch { case m.Type == pb.MsgPreVote: // 应答预投票消息时不修改任期 case m.Type == pb.MsgPreVoteResp && !m.Reject: // 预投票消息未被拒绝处理 default: // 心跳、日志、快照消息,退位记录新leader if m.Type == pb.MsgApp || m.Type == pb.MsgHeartbeat || m.Type == pb.MsgSnap { r.becomeFollower(m.Term, m.From) // 其他消息 } else { r.becomeFollower(m.Term, None) // 仅直到term更新,leader未知 } } ... } ... return nil }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25leader处理消息时,步入leader专属的状态机函数stepLeader,此时leader倘若发现自己和集群的多数派断开联系,也会回退到follower。func stepLeader(r *raft, m pb.Message) error { switch m.Type { case pb.MsgBeat: ... case pb.MsgCheckQuorum: ... // 检查集群可用性 if !r.prs.QuorumActive() { // 超过半数的服务没有活跃变成follower状态 r.becomeFollower(r.Term, None) } ... } return nil }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
# 5.8.follower->candidate
应用层定时调用
Node.Tick函数驱动算法层调用定时函数,raft节点角色为follower或者candidate,其定时函数是tickElection。func (r *raft) tickElection() { r.electionElapsed++ if r.promotable() && r.pastElectionTimeout() { // 有资格选举,达到选举时间 r.electionElapsed = 0 // 发布HUP消息重新开始选举 r.Step(pb.Message{From: r.id, Type: pb.MsgHup}) } }1
2
3
4
5
6
7
8
9
10每次调用
tickElection会把竞选计时器raft.electionElapsed的tick数累加1,超过选举时间间隔,会给本节点推一条MsgHup类型的消息发起预选举,防止脑裂原因导致follower无限选举增大自身term。func (r *raft) pastElectionTimeout() bool { return r.electionElapsed >= r.randomizedElectionTimeout }1
2
3通用状态机函数中,会对
MsgHup类型的消息进行响应,调用raft.compaign方法发起选举。func (r *raft) Step(m pb.Message) error { ... switch m.Type { case pb.MsgHup: // 非Leader收到HUP消息,说明准备进行选举 if r.state != StateLeader { ... //进行选举 if r.preVote { r.campaign(campaignPreElection) } else { r.campaign(campaignElection) } } ... } return nil }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18最终调用
raft.campaign方法进行选举,内部进一步调用raft.becomePreCandidate和raft.becomeCandidate方法切换角色状态。func (r *raft) campaign(t CampaignType) { ... var term uint64 var voteMsg pb.MessageType // 预投票选举 if t == campaignPreElection { // 进入预候选状态 r.becomePreCandidate() // 使用预投票消息类型 voteMsg = pb.MsgPreVote // 不更新任期,但预投票时使用递增的任期号 term = r.Term + 1 // 正式投票选举 } else { // 变为候选状态 r.becomeCandidate() // 使用正式投票消息类型 voteMsg = pb.MsgVote // 正式投票使用当前任期号 term = r.Term } // 自投票阶段,节点给自己投票就获取多数(单节点集群) if _, _, res := r.poll(r.id, voteRespMsgType(voteMsg), true); res == quorum.VoteWon { // 预选举成功,发起正式选举 if t == campaignPreElection { r.campaign(campaignElection) // 正式选举成功,成为leader } else { r.becomeLeader() } return } ... // 整理有投票权的节点ID,向集群里的其他节点发送竞选消息 for _, id := range ids { if id == r.id { // 跳过自己 continue } ... // 广播竞选拉票请求 r.send(pb.Message{Term: term, To: id, Type: voteMsg, Index: r.raftLog.lastIndex(), LogTerm: r.raftLog.lastTerm(), Context: ctx}) } }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
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
# 5.9.candidate->leader/follower
candidate竞选期间收到任期更大的消息,会回退到follower,这部分逻辑位于通用状态机函数,和leader退位回follower的处理路径相同。func (r *raft) Step(m pb.Message) error { switch { ... case m.Term > r.Term: ... switch { case m.Type == pb.MsgPreVote: // 预投票处理 case m.Type == pb.MsgPreVoteResp && !m.Reject: // 正式投票处理 default: // 日志、心跳、快照消息转为follower if m.Type == pb.MsgApp || m.Type == pb.MsgHeartbeat || m.Type == pb.MsgSnap { r.becomeFollower(m.Term, m.From) // 其他消息 } else { r.becomeFollower(m.Term, None) } } ... } ... return nil }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24candidate收到半数以上拒绝票,会回退到follower,否则收到半数以上的赞成票,会进位为leader,这部分处理在candidate的定制状态机函数stepCandidate中。func stepCandidate(r *raft, m pb.Message) error { ... switch m.Type { ... case pb.MsgApp: // 收到日志追加消息说明已经有合法leader,切换角色为follower r.becomeFollower(m.Term, m.From) // always m.Term == r.Term // 追加日志 r.handleAppendEntries(m) case pb.MsgHeartbeat: // 收到心跳消息,说明已经有合法的leader,切换角色为follower r.becomeFollower(m.Term, m.From) // 处理心跳消息 r.handleHeartbeat(m) case pb.MsgSnap: // 收到快照消息,说明已经有合法leader,切换角色为follower r.becomeFollower(m.Term, m.From) // 处理快照消息 r.handleSnapshot(m) case myVoteRespType: // 加入最新一票后投票统计(同意票数、拒绝票数、选举结果) gr, rj, res := r.poll(m.From, m.Type, !m.Reject) switch res { case quorum.VoteWon: // 预选举成功,发起正式选举 if r.state == StatePreCandidate { r.campaign(campaignElection) // 正式选举成功,切换为leader并广播日志 } else { r.becomeLeader() r.bcastAppend() } case quorum.VoteLost: // 选举失败,退回follower r.becomeFollower(r.Term, None) } case pb.MsgTimeoutNow: // 候选者不处理超时,只有follower才处理领导权转移 } return nil }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
29
30
31
32
33
34
35
36
37
38
39
40
41
# 6.写流程
# 6.1.宏观流程
raft节点处理写请求需要经历两阶段提交过程。当前节点为leader时,应用层调用Node.Propose发送写请求给算法层,算法层会添加日志到unstable,同时通过Node.Ready方法让应用层消费到,应用层记录日志到wal,存储日志到raftStorage,然后清除刚添加的日志unstable,通过通信模块转发消息给集群内的其他节点,最后调用Node.Adavance开启下一轮交互。
- 当前节点为
follower,应用层接收到来自leader的同步日志请求,调用Node.Step将请求转发给算法层,算法层同步日志到unstable后,将同步结果封装到Ready结构体,通过Node.Ready方法让应用层消费到,应用层会存储日志到raftStorage,然后清除刚添加的日志unstable,通过通信模块转发给leader,最后调用Node.Advance开启下一轮交互。 
- 再回到
leader,应用层收到来自follower的同步响应日志,会调用Node.Step将请求转发给算法层,算法层检测同步请求得到多数派同意,更新集群内的其他节点的已提交日志同步进度,同时更新自己的已提交日志同步进度,封装成新的commited index到Ready,应用层通过Node.Ready()消费,将已提交的日志应用到数据状态机,广播已提交的日志请求,最后调用Node.Advance开启下一轮交互。、 
# 6.2.应用层发送写请求
客户端向应用层提交写数据请求,
raftNode调用Node.Propose方法将请求传到算法层的propc channel中,消息类型为MsgProp。func (rc *raftNode) serveChannels() { ... // send proposals over raft go func() { confChangeCount := uint64(0) for rc.proposeC != nil && rc.confChangeC != nil { select { // 接收客户端的写请求 case prop, ok := <-rc.proposeC: ... // 转发至算法层处理 rc.node.Propose(context.TODO(), []byte(prop)) case cc, ok := <-rc.confChangeC: ... confChangeCount++ cc.ID = confChangeCount rc.node.ProposeConfChange(context.TODO(), cc) } } // client closed channel; shutdown raft if not already close(rc.stopc) }() ... } func (n *node) Propose(ctx context.Context, data []byte) error { return n.stepWait(ctx, pb.Message{Type: pb.MsgProp, Entries: []pb.Entry{{Data: data}}}) } func (n *node) stepWait(ctx context.Context, m pb.Message) error { return n.stepWithWaitOption(ctx, m, true) } func (n *node) stepWithWaitOption(ctx context.Context, m pb.Message, wait bool) error { ... // 算法层的propc channel ch := n.propc pm := msgWithResult{m: m} select { //写请求发送到算法层管道处理 case ch <- pm: ... } return nil }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
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
# 6.3.不同身份接收到写请求处理
算法层接收到消息,进入状态机处理函数,根据节点的不同身份做出不同的处理,也就是每个身份都有一个专属的状态函数。
// 启动算法层,持续与应用层进行通信交互 func (n *node) run() { ... for { ... select { case pm := <-propc: m := pm.m m.From = r.id // 处理本地收到的提交值,该方法是通用状态机函数,用于专用处理前的消息前置处理 err := r.Step(m) ... } } }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15follower的定制状态机函数是stepFollower(),由于只有leader可以处理写请求,因此会转发写请求到leader。// follower状态机处理函数 func stepFollower(r *raft, m pb.Message) error { switch m.Type { case pb.MsgProp: ... // 转发到leader m.To = r.lead r.send(m) ... } return nil }1
2
3
4
5
6
7
8
9
10
11
12candidate的定制状态机处理函数是stepCandidate(),candidate的出现代表集群没有leader,这里会直接打印错误且不进行额外处理。func stepCandidate(r *raft, m pb.Message) error { ... switch m.Type { case pb.MsgProp: // candidate不能处理写请求 return ErrProposalDropped ... } return nil }1
2
3
4
5
6
7
8
9
10leader的状态机处理函数是stepLeader(),用于处理写请求并同步给follower,根据同步结果进行持久化存储。func stepLeader(r *raft, m pb.Message) error { // These message types do not require any progress for m.From. switch m.Type { ... case pb.MsgProp: ... // 添加预写日志到raftLog.unstable,更新raft.ptrs中自己的Match和Next if !r.appendEntry(m.Entries...) { return ErrProposalDropped } // 广播添加append消息 r.bcastAppend() return nil ... } ... return nil } // 算法层添加预写日志到内存 func (r *raft) appendEntry(es ...pb.Entry) (accepted bool) { // 获取raftLog的最后一条日志的索引 li := r.raftLog.lastIndex() // 向添加的预写日志设置term和index for i := range es { es[i].Term = r.Term es[i].Index = li + 1 + uint64(i) } ... // 添加日志到raftLog.unstable li = r.raftLog.append(es...) // 更新自己的同步进度 r.prs.Progress[r.id].MaybeUpdate(li) // 处理leader写请求的时候是多余的调用,此时只有leader自己写入内存unstable,其他节点并未写入,不可能更新commited index r.maybeCommit() return true } // 算法层添加日志到内存 func (l *raftLog) append(ents ...pb.Entry) uint64 { ... l.unstable.truncateAndAppend(ents) return l.lastIndex() } // 判断添加的预写日志是否需要回滚 func (u *unstable) truncateAndAppend(ents []pb.Entry) { // 先获取要添加的预写日志的第一条数据的索引 after := ents[0].Index switch { case after == u.offset+uint64(len(u.entries)): // 首条日志是下一条要添加的数据,直接添加 // u.offset+uint64(len(u.entries))计算的是已经存到unstable的最后一条日志索引+1 u.entries = append(u.entries, ents...) case after <= u.offset: // 首条日志的索引小于偏移量,需要回滚全部的unstable日志,替换新的偏移量和entries // 简单来说,添加的日志在offset之前,直接丢弃offset后面的entries脏数据 u.offset = after u.entries = ents default: // u.offset<after<uint64(len(u.entries)),回滚after后面的数据 u.entries = append([]pb.Entry{}, u.slice(u.offset, after)...) // 拼接当前的日志 u.entries = append(u.entries, ents...) } }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
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
bcastAppend()遍历集群中的每个节点,调用sendAppend()将要同步的日志、要同步日志的上一条日志的索引和任期、当前leader已提交的索引全部封装在类型为MagApp的meaasge,后续算法层会检测msgs是否有消息并以Ready结构体的形式发送到应用层消费处理。func (r *raft) bcastAppend() { // 遍历集群中节点 r.prs.Visit(func(id uint64, _ *tracker.Progress) { if id == r.id { return } // 封装消息到raft.msgs r.sendAppend(id) }) } func (r *raft) sendAppend(to uint64) { r.maybeSendAppend(to, true) } func (r *raft) maybeSendAppend(to uint64, sendIfEmpty bool) bool { ... m := pb.Message{} m.To = to // 获取上一条日志的term term, errt := r.raftLog.term(pr.Next - 1) // 获取发往节点的已提交索引之后的预写日志,数量总和不超过maxMsgSize ents, erre := r.raftLog.entries(pr.Next, r.maxMsgSize) ... if errt != nil || erre != nil { // send snapshot if we failed to get term or entries ... // 任期及预写日志获取失败时发送快照 } else { // 消息类型为append m.Type = pb.MsgApp // 接收方的commited index m.Index = pr.Next - 1 // 接收方的commited term m.LogTerm = term // 待发送的预写日志 m.Entries = ents // 当前leader的commited index m.Commit = r.raftLog.committed ... } // 封装message r.send(m) return true } func (r *raft) send(m pb.Message) { m.From = r.id // 对拉票、读、写请求之外的消息填充任期信息 if m.Type == pb.MsgVote || m.Type == pb.MsgVoteResp || m.Type == pb.MsgPreVote || m.Type == pb.MsgPreVoteResp { ... } else { if m.Term != 0 { panic(fmt.Sprintf("term should not be set when sending %s (was %d)", m.Type, m.Term)) } // 填充任期信息 if m.Type != pb.MsgProp && m.Type != pb.MsgReadIndex { m.Term = r.Term } } // 追加到raft.msgs r.msgs = append(r.msgs, m) }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
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62注意
算法层启动的循环中会监听
raft.msgs消息,基于newReady()封装上面处理好的msgs消息,转至readyc channel供应用层消费处理
# 6.4.应用层持久化日志与广播日志同步消息
算法层处理好的消息会被放到
readyc channel供应用层捕获处理,应用层通过node.Ready()接收到算法层发送的消息,将写操作记录到日志WAL,同时将算法层存放到unstable的那一条待应用层持久化的entry到raftStorage,最后广播日志同步请求给集群内的其他节点。这里需要注意的是,raftStorage和unstable存储的日志此时都未提交,可能有脏数据,raftStorage.Append和unstable.truncateAndAppend这两个方法分别处理应用层和算法层的日志回滚问题。func (rc *raftNode) serveChannels() { ... // event loop on raft state machine updates for { select { ... case rd := <-rc.node.Ready(): // 持久化预写日志及状态 rc.wal.Save(rd.HardState, rd.Entries) ... // 预写日志追加到storage rc.raftStorage.Append(rd.Entries) // 向其他节点发送日志同步消息 rc.transport.Send(rd.Messages) // 过滤出需要应用到状态机的数据及提交到状态机 if ok := rc.publishEntries(rc.entriesToApply(rd.CommittedEntries)); !ok { rc.stop() return } // 通知算法层发送完毕,开启下一轮交互 rc.node.Advance() ... } } } // 添加日志到raftStorage,主要确保添加的日志是可提交的,存在脏数据时会回滚 func (ms *MemoryStorage) Append(entries []pb.Entry) error { ... ms.Lock() defer ms.Unlock() // 获取storage中存储的第一条日志索引 first := ms.firstIndex() // 获取要追加日志的最后一条索引 last := entries[0].Index + uint64(len(entries)) - 1 // 没有更新日志,结束 if last < first { return nil } // 待追加日志有部分过期的 if first > entries[0].Index { // 回滚first之前已添加过的日志数据 entries = entries[first-entries[0].Index:] } // 计算相对storage第一条日志的偏移量 offset := entries[0].Index - ms.ents[0].Index switch { case uint64(len(ms.ents)) > offset: // 有重叠部分,截断offset之前的 ms.ents = append([]pb.Entry{}, ms.ents[:offset]...) // 追加本地添加的 ms.ents = append(ms.ents, entries...) case uint64(len(ms.ents)) == offset: // 索引连续,直接追加 ms.ents = append(ms.ents, entries...) ... } return nil }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
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
# 6.5.follower和candidate处理同步日志提议
candidate和follower主要在定制化的状态机处理函数中处理日志同步提议。candidate收到任期大于等于自己竞选任期的同步日志请求后,会回退follower,尝试添加日志到raftLog.unstable。func stepCandidate(r *raft, m pb.Message) error { ... switch m.Type { ... case pb.MsgApp: // 收到append消息,说明集群中出现leader,转为follower r.becomeFollower(m.Term, m.From) // always m.Term == r.Term // 添加日志(消息索引不足提交索引时响应当前已提交索引,追加成功响应最后一条追加索引,失败响应拒绝消息) r.handleAppendEntries(m) ... } return nil }1
2
3
4
5
6
7
8
9
10
11
12
13follower收到同步日志请求后,会将选举时间重置,同时尝试将日志添加到raftLog.unstable。func stepFollower(r *raft, m pb.Message) error { switch m.Type { ... case pb.MsgApp: // 收到leader的日志同步消息,重置选举tick计时器 r.electionElapsed = 0 // 标记leader r.lead = m.From // 日志追加到非持久化日志列表 r.handleAppendEntries(m) ... } return nil }1
2
3
4
5
6
7
8
9
10
11
12
13
14follower和candidate共用handleAppendEntries()方法,同步日志且进行数据匹配。当leader发送的待同步消息的前一条日志和本节点已同步的最后一条日志索引匹配上,就接受新的日志同步。匹配不上时就会拒绝日志同步,因为本届点缺失部分leader已经持久化的日志记录,拒绝时会返回本节点最新已提交的索引给leader,方便leader补发本节点缺失的日志数据。最后会调用send()封装同步日志的响应,后续会在算法层获取响应传递给应用层消费处理。func (r *raft) handleAppendEntries(m pb.Message) { // 发送方的index小于本节点已提交的index,告知对方自己的commited index和term // 代表发送方的消息落后 if m.Index < r.raftLog.committed { r.send(pb.Message{To: m.From, Type: pb.MsgAppResp, Index: r.raftLog.committed}) return } // 尝试添加到日志模块 // mlastIndex是当前follower最新日志的index if mlastIndex, ok := r.raftLog.maybeAppend(m.Index, m.LogTerm, m.Commit, m.Entries...); ok { // 添加成功,返回follower已提交的最新索引 r.send(pb.Message{To: m.From, Type: pb.MsgAppResp, Index: mlastIndex}) } else { // 添加失败,响应拒绝同步消息,返回follower已提交的最新索引 r.send(pb.Message{To: m.From, Type: pb.MsgAppResp, Index: m.Index, Reject: true, RejectHint: r.raftLog.lastIndex()}) } } // 添加日志,返回最新日志的index func (l *raftLog) maybeAppend(index, logTerm, committed uint64, ents ...pb.Entry) (lastnewi uint64, ok bool) { // 验证领导者同步的前一条日志是否与本地匹配 if l.matchTerm(index, logTerm) { lastnewi = index + uint64(len(ents)) // 查找传入的数据term冲突的位置,正常情况下是ents[0].index // 领导者日志:[index:term] // 1:1, 2:1, [3:2, 4:2, 5:3]本次传入的 // 跟随者日志: // 1:1, 2:1, 3:1, 4:1 // 冲突位置就是index=3 ci := l.findConflict(ents) switch { case ci == 0: // 无冲突,当前传入的日志都已经添加 case ci <= l.committed: // raft协议保证一旦日志条目被提交不能更改或删除,这也是集群达成共识的状态 // 所以冲突位置发生在已提交区域,说明存在严重不一致,无法修复 l.logger.Panicf("entry %d conflict with committed entry [committed(%d)]", ci, l.committed) default: // 冲突发生在未提交区域,此时未提交日志尚未被集群确认 // 新领导者的日志可以安全覆盖冲突位置后的日志 // 假设提交index=2,冲突index=3,此时覆盖后就是 1:1, 2:1, 3:2, 4:2, 5:3 offset := index + 1 l.append(ents[ci-offset:]...) } // 提交不超出当前日志长度的、已有的日志 l.commitTo(min(committed, lastnewi)) return lastnewi, true } return 0, false }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
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50commitTo安全性考量
1.不能提交超过当前日志长度的索引
lastnewi2.防止提交不存在的日志条目,确保状态机只应用实际存在的日志
3.领导者只能提交当前拥有的日志条目
4.
leader的commited大于接收者的日志长度,说明接收者落后,此时取小值提交可以避免越界
# 6.6.leader收到follower成功响应提交日志
leader每次收到follower的同步日志响应时,会在定制状态机函数中进一步处理。当follower日志进度滞后而拒绝,可以根据Message.RejectHint对缺失的日志进行补发。follower接收同步日志的提议,leader会在raft.mayBeUpdate方法中更新对应节点的日志进度Progress,同时在raft.mayBeCommit中基于多数派原则更新已提交日志的索引,最后广播给follower最新提交日志的索引。func stepLeader(r *raft, m pb.Message) error { ... // 检查消息发送者是否在当前集群 pr := r.prs.Progress[m.From] if pr == nil { return nil } switch m.Type { case pb.MsgAppResp: // 处理append应答消息 pr.RecentActive = true // follower拒绝同时时,说明term/index不匹配 if m.Reject { // 根据follower建议的回退位置更新期望发送的Next日志位置 if pr.MaybeDecrTo(m.Index, m.RejectHint) { // 乐观复制副本状态切换为探测状态 if pr.State == tracker.StateReplicate { pr.BecomeProbe() } // 调整Next索引后尝试重新发送日志 r.sendAppend(m.From) } } else { // follower通过append消息 oldPaused := pr.IsPaused() // 根据响应位置更新follower的同步信息(match、index) if pr.MaybeUpdate(m.Index) { switch { // 跟随者已同步,保守探测状态切换为复制状态 case pr.State == tracker.StateProbe: pr.BecomeReplicate() ... // 本身处于复制状态,释放已确认的inFlight消息空间(本次日志同步成功,释放≤m.Index的日志条目) case pr.State == tracker.StateReplicate: pr.Inflights.FreeLE(m.Index) } // 尝试commit(半数通过) // 基于各节点 Progesss.Match组成数组进行逆序排列 // 下中位数的日志索引就是获得多数派认可可提交的日志索引 // 1.可提交索引≥commited // 2.可提交索引位于非持久化日志索引范围内 // 3.可提交索引的任期与当前节点任期匹配 // 4.计算任期时,根据mayIndex-offset计算逻辑索引对应的物理索引,从而读出内存的对应日志的任期(最大确认日志的任期) if r.maybeCommit() { // leader更新自己的commited index为上一笔预写日志的index // 广播通知其他节点持久化commited index(封装提交消息至msgs) r.bcastAppend() // 该跟随者由于流控暂停发送 } else if oldPaused { // 单独向该节点发送最新状态,避免其长期停滞 r.sendAppend(m.From) } // 继续发送待同步的日志,直到日志处理完毕、流控窗口已满、网络连接不可用 for r.maybeSendAppend(m.From, false) { } // 响应消息的节点是领导权转移目标且日志完全匹配,触发目标立即选举消息 if m.From == r.leadTransferee && pr.Match == r.raftLog.lastIndex() { r.sendTimeoutNow(m.From) } } } ... } return nil }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
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65maybeDecrTo()主要更新follower的日志同步进度,follower拒绝leader日志同步会告知自己当前最新的已提交日志索引commited index,leader会更新这个follower的Next为commited index+1,然后整理补发日志重新发送给follower同步。// 更新拒绝日志同步的节点的同步进度 func (pr *Progress) MaybeDecrTo(rejected, last uint64) bool { // 高效复制模式(单次发多条) if pr.State == StateReplicate { // 拒绝的位置位于Match前,说明此拒绝消息已过期 if rejected <= pr.Match { return false } // 回退到已确认位置的下一条 pr.Next = pr.Match + 1 return true } // 拒绝的不是最后探测的条目,说明拒绝消息过期 if pr.Next-1 != rejected { return false } // 取已拒绝位置和对方已提交位置的较小值 if pr.Next = min(rejected, last+1); pr.Next < 1 { pr.Next = 1 } // 重置探测标记允许重新探测 pr.ProbeSent = false return true }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提交
maybeCommit()获取所有节点的Match,从小到大排序后取中位数,此中位数大于当前leader的commited index且term匹配,会更新leader的commited index。
# 6.7.leader算法层更新提交索引通知应用层
leader算法层收到多数派follower同意后,提交日志及广播给其他节点,广播内部就是就是将提交日志消息封装到msgs,算法层的node.run()循环检测到就绪消息重新newReady(),最后将Ready送到readyc channel供应用层消费。func (rc *raftNode) serveChannels() { ... // event loop on raft state machine updates for { select { ... 又回到这里 // 1.写操作记录到WAL日志 // 2.待持久化的预写日志追击到storage // 3.调用通信模块执行消息发送 // 4.应用算法层已提交的预写日志至状态机 // 5.通知算法层进行下一轮交互 case rd := <-rc.node.Ready(): rc.wal.Save(rd.HardState, rd.Entries) ... rc.raftStorage.Append(rd.Entries) rc.transport.Send(rd.Messages) if ok := rc.publishEntries(rc.entriesToApply(rd.CommittedEntries)); !ok { rc.stop() return } rc.maybeTriggerSnapshot() rc.node.Advance() ... } } }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
26entriesToApply用于去掉已经应用过的Entry,防止每次刷新状态机时重复应用。// 去掉已经应用过的entry func (rc *raftNode) entriesToApply(ents []raftpb.Entry) (nents []raftpb.Entry) { ... // 已提交日志的第一个索引 firstIdx := ents[0].Index // 日志不连续,记录 if firstIdx > rc.appliedIndex+1 { ... } // 存在未应用的日志,获取未应用的子集 if rc.appliedIndex-firstIdx+1 < uint64(len(ents)) { nents = ents[rc.appliedIndex-firstIdx+1:] } return nents }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15publishEntries()主要就是把要应用的Entry直接丢进rc.commitC channel中,kvstore会监听这个channel将Entry放到数据状态机。func (rc *raftNode) publishEntries(ents []raftpb.Entry) bool { for i := range ents { switch ents[i].Type { case raftpb.EntryNormal: ... // 获取待应用的预写日志 s := string(ents[i].Data) select { // 待应用的预写日志传入raft.commitC channel,由kvstore.readCommits()存入状态机 case rc.commitC <- &s: ... } ... } } return true }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
# 7.读流程
# 7.1.整体流程
--- leader 应用层调用Node.ReadIndex发送读请求给算法层,处于ReadOnlyLeaseBased模式会直接封装响应到readyc channel,应用层直接响应客户端的读请求;处于ReadOnlySafe模式读请求会挂起到读请求队列,leader向follower发送心跳,获得多数派认可后封装响应到readyc channel,最后调用Node.Advance开启下一轮交互 --- follower 应用层接收到来自leader的心跳请求,调用Node.Step将请求转发到算法层,算法层将心跳响应结果封装到Ready结构体,通过Node.Ready方法让应用层消费到,应用层通过通信模块转给leader,之后调用Node.Adavance开启下一轮交互 --- leader follower的响应会被streamReader.decodeLoop捕获放入recvc channel中,这个recvc channel中的消息最终会在startPeer中消费,通过transport层的EtcdServer.Process()间接调用算法层node.Step(),将消息转交给算法层的recvc供底层获取处理,leader接收到半数以上节点的Ack则判断自身身份合法,封装读请求响应消息通过Node.Ready让应用层消费到,应用层收到消息后读取数据状态机,对客户端的读请求进行响应1
2
3
4
5
6
7
8
# 7.2.应用层发送读请求
应用层
raftNode调用Node.ReadIndex方法,通过recvc向算法层发送一条类型为ReadIndex的消息。func (n *node) ReadIndex(ctx context.Context, rctx []byte) error { return n.step(ctx, pb.Message{Type: pb.MsgReadIndex, Entries: []pb.Entry{{Data: rctx}}}) } func (n *node) step(ctx context.Context, m pb.Message) error { return n.stepWithWaitOption(ctx, m, false) } func (n *node) stepWithWaitOption(ctx context.Context, m pb.Message, wait bool) error { if m.Type != pb.MsgProp { select { // 不是写请求,放入recvc待算法层处理 case n.recvc <- m: return nil ... } } ... return nil }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
# 7.3.不同身份接收读请求处理
follower收到读请求后,会直接转发给leader,让leader处理好后给自己响应处理结果,最后再响应给客户端。func stepFollower(r *raft, m pb.Message) error { switch m.Type { ... case pb.MsgReadIndex: ... m.To = r.lead // 读请求转发至leader r.send(m) case pb.MsgReadIndexResp: // 处理响应消息 r.readStates = append(r.readStates, ReadState{Index: m.Index, RequestCtx: m.Entries[0].Data}) } return nil }1
2
3
4
5
6
7
8
9
10
11
12
13
14leader接收到读请求后,会根据不同读模式进行处理。ReadOnlyLeaseBased模式下,会认为Leader还在有效期内,无需其他节点认证,直接通知应用层响应客户端的读请求;ReadOnlySafe模式下,读请求会挂起到读请求队列,先向所有节点广播context带上读请求的ID的心跳,标识这是一笔特殊的心跳请求,用于leader身份认证。func stepLeader(r *raft, m pb.Message) error { // These message types do not require any progress for m.From. switch m.Type { ... case pb.MsgReadIndex: // 单节点直接处理 if r.prs.IsSingleton() { if resp := r.responseToReadIndexReq(m, r.raftLog.committed); resp.To != None { r.send(resp) } return nil } ... sendMsgReadIndexResponse(r, m) return nil } ... return nil } // 根据不同读模式处理 func sendMsgReadIndexResponse(r *raft, m pb.Message) { switch r.readOnly.option { // 安全模式 case ReadOnlySafe: // 请求添加到readOnly只读队列,记录当前commited Index作为基准 r.readOnly.addRequest(r.raftLog.committed, m) // 本地节点的响应信息加入到只读请求的ACK队列 r.readOnly.recvAck(r.id, m.Entries[0].Data) // 广播心跳消息等待其他节点确认(msgs) r.bcastHeartbeatWithCtx(m.Entries[0].Data) // 租约模式 case ReadOnlyLeaseBased: // 不需要多数派确认,基于租约机制认为领导者在任期内,直接返回commited Index if resp := r.responseToReadIndexReq(m, r.raftLog.committed); resp.To != None { r.send(resp) } } }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
29
30
31
32
33
34
35
36
37
38
39
40
# 7.4.follower收到leader的广播心跳
follower收到leader的广播心跳后,会重置自己的选举时间,更新自己的commited index,最后响应leader的心跳请求。func stepFollower(r *raft, m pb.Message) error { switch m.Type { ... case pb.MsgHeartbeat: r.electionElapsed = 0 r.lead = m.From // 处理心跳消息 r.handleHeartbeat(m) ... } return nil } func (r *raft) handleHeartbeat(m pb.Message) { // 提交消息到leader commited index r.raftLog.commitTo(m.Commit) // 封装心跳响应消息至msgs r.send(pb.Message{To: m.From, Type: pb.MsgHeartbeatResp, Context: m.Context}) }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
# 7.5.leader收到follower响应
leader收到多数派的follower响应自己的心跳,说明leader至少在收到这笔读请求的时刻具备合法身份,拥有最新数据,因此会对这笔请求及之前挂起的读请求批量做出响应。func stepLeader(r *raft, m pb.Message) error { ... // 检查消息发送者是否在集群 pr := r.prs.Progress[m.From] if pr == nil { return nil } switch m.Type { ... case pb.MsgHeartbeatResp: ... // 针对当前消息已经应答的节点数量少于半数不处理 if r.prs.Voters.VoteResult(r.readOnly.recvAck(m.From, m.Context)) != quorum.VoteWon { return nil } // 调用advance函数获取该读请求前的一系列读请求,同时将这些读请求从readOnly.pendingReadIndex和readIndexQueue删除 rss := r.readOnly.advance(m) // 遍历准备丢弃的读请求 for _, rs := range rss { req := rs.req // 本地请求 if req.From == None || req.From == r.id { r.readStates = append(r.readStates, ReadState{Index: rs.index, RequestCtx: req.Entries[0].Data}) // 外部请求应答 } else { r.send(pb.Message{To: req.From, Type: pb.MsgReadIndexResp, Index: rs.index, Entries: req.Entries}) } } ... } return nil }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
29
30
31
32