Raft 共识协议工程实现:从选举到日志复制的分布式一致性构建

一、分布式系统的数据同步难题
在分布式系统里,数据一致性是最让人头疼的问题。客户端写入一条数据,得保证多个副本上都有这份数据,这样才能既持久又可用。但现实很骨感:网络会分区、节点会宕机、消息还会乱序。这些故障是常态,一旦处理不好,不同节点上的数据就会对不上号。
如果没有一致性协议兜底,集群面对故障时只有两条路:要么停服保一致,要么接着跑但数据乱了。这两条在生产环境里都是死路。
Raft 协议把共识问题拆成了三块:选主、日志复制、安全保证。这么拆之后,理解难度和实现复杂度都比 Paxos 低了不少,但一致性语义没打折。工程上,Raft 的状态机模型很清晰,代码写出来容易测、容易验。这也是 etcd、Consul、TiKV 这些生产系统都选 Raft 的原因。
不过,Raft 论文看着简单,真落地上坑不少。选举超时的抖动怎么设、日志怎么压缩和快照、线性一致性读怎么搞、网络分区后怎么恢复,这些论文里轻描淡写的地方,恰恰是生产环境里最容易踩雷的。
二、Raft 状态机与日志复制机制
Raft 把每个节点看作一个有限状态机,角色就三种:Follower、Candidate、Leader。Leader 负责处理客户端请求,把日志复制到 Follower;Follower 被动收日志;Candidate 是选举期间的临时角色,超时了就发起选举。
stateDiagram-v2
[*] –> Follower : 节点启动
Follower –> Candidate : 选举超时<br/>未收到 Leader 心跳
Candidate –> Leader : 获得多数派投票
Candidate –> Follower : 收到更高 Term 的心跳
Candidate –> Follower : 选举超时<br/>未获多数票
Leader –> Follower : 发现更高 Term 的请求
state Leader {
[*] –> 接收客户端请求
接收客户端请求 –> 追加本地日志
追加本地日志 –> 并行复制到Follower
并行复制到Follower –> 等待多数派确认
等待多数派确认 –> 提交日志
提交日志 –> 应用到状态机
应用到状态机 –> 响应客户端
}
state Follower {
[*] –> 接收AppendEntries
接收AppendEntries –> 日志一致性检查
日志一致性检查 –> 追加日志或拒绝
追加日志或拒绝 –> 确认提交索引
}
日志复制的流程其实不复杂:Leader 把客户端请求打包成日志条目,存到本地,然后并行发给所有 Follower。等这条日志在多数节点上都存上了,Leader 就把它标记为已提交(Committed),再通知 Follower 更新提交索引。一旦提交,日志就能安全地应用到状态机,结果对客户端可见。
日志一致性检查是安全性的核心。AppendEntries RPC 里带着 prev_log_index 和 prev_log_term,Follower 在追加日志前,得先看看本地对应位置的日志跟 RPC 里的信息对得上不。对不上就拒绝,Leader 得回退重试。这种逐条回退的机制保证了:只要两个日志在某个索引上一样,那这个索引之前的所有日志也都一样。
三、Raft 核心模块的 Rust 实现
下面这段代码实现了 Raft 的核心数据结构跟关键逻辑,包括状态转换、日志追加和选举。
use std::collections::HashMap;
/// Raft 节点角色
#[derive(Debug, Clone, Copy, PartialEq)]
pub enum Role {
Follower,
Candidate,
Leader,
}
/// 日志条目
#[derive(Debug, Clone)]
pub struct LogEntry {
pub term: u64,
pub index: u64,
pub command: Vec<u8>,
}
/// Raft 持久化状态(变更前必须持久化到稳定存储)
#[derive(Debug, Clone)]
pub struct PersistentState {
pub current_term: u64,
pub voted_for: Option<u64>,
pub log: Vec<LogEntry>,
}
/// Raft 易失状态
#[derive(Debug)]
pub struct VolatileState {
pub commit_index: u64,
pub last_applied: u64,
}
/// Leader 专有的易失状态
#[derive(Debug)]
pub struct LeaderState {
/// 每个 Follower 的下一个复制索引
pub next_index: HashMap<u64, u64>,
/// 每个 Follower 的已匹配索引
pub match_index: HashMap<u64, u64>,
}
/// Raft 节点核心结构
pub struct RaftNode {
pub id: u64,
pub peers: Vec<u64>,
pub role: Role,
pub persistent: PersistentState,
pub volatile: VolatileState,
pub leader_state: Option<LeaderState>,
}
impl RaftNode {
pub fn new(id: u64, peers: Vec<u64>) -> Self {
Self {
id,
peers,
role: Role::Follower,
persistent: PersistentState {
current_term: 0,
voted_for: None,
log: Vec::new(),
},
volatile: VolatileState {
commit_index: 0,
last_applied: 0,
},
leader_state: None,
}
}
/// 处理 RequestVote RPC
/// 核心规则:每个 Term 只投一票,且候选人的日志至少与自己一样新
pub fn handle_request_vote(
&mut self,
candidate_id: u64,
candidate_term: u64,
candidate_last_log_index: u64,
candidate_last_log_term: u64,
) -> (bool, u64) {
// 如果候选人的 Term 更高,先回退到 Follower
if candidate_term > self.persistent.current_term {
self.step_down(candidate_term);
}
if candidate_term < self.persistent.current_term {
return (false, self.persistent.current_term);
}
// 检查是否已投票给其他候选人
let grant_vote = match self.persistent.voted_for {
Some(voted_id) if voted_id != candidate_id => false,
_ => {
// 日志完整性检查:候选人的日志至少与自己一样新
let my_last = self.last_log();
let candidate_is_up_to_date = candidate_last_log_term > my_last.term
|| (candidate_last_log_term == my_last.term
&& candidate_last_log_index >= my_last.index);
candidate_is_up_to_date
}
};
if grant_vote {
self.persistent.voted_for = Some(candidate_id);
// 持久化 voted_for(生产环境中写入稳定存储)
}
(grant_vote, self.persistent.current_term)
}
/// 处理 AppendEntries RPC
/// 返回 (success, current_term)
pub fn handle_append_entries(
&mut self,
leader_term: u64,
prev_log_index: u64,
prev_log_term: u64,
entries: Vec<LogEntry>,
leader_commit: u64,
) -> (bool, u64) {
if leader_term < self.persistent.current_term {
return (false, self.persistent.current_term);
}
// 收到合法 Leader 的心跳,回退到 Follower
if self.role != Role::Follower {
self.step_down(leader_term);
}
// 日志一致性检查
if prev_log_index > 0 {
match self.get_log_entry(prev_log_index) {
Some(entry) if entry.term == prev_log_term => {}
_ => return (false, self.persistent.current_term),
}
}
// 追加新日志条目(处理冲突:截断不一致的部分)
for entry in entries {
match self.get_log_entry(entry.index) {
Some(existing) if existing.term == entry.term => {
// 日志一致,跳过
}
Some(_) => {
// 冲突:截断从此索引开始的所有日志
self.persistent.log.truncate(
(entry.index as usize).saturating_sub(1),
);
self.persistent.log.push(entry);
}
None => {
self.persistent.log.push(entry);
}
}
}
// 更新提交索引
if leader_commit > self.volatile.commit_index {
let last_new_index = self.persistent.log.last()
.map(|e| e.index)
.unwrap_or(0);
self.volatile.commit_index =
leader_commit.min(last_new_index);
}
(true, self.persistent.current_term)
}
/// Leader 推进提交索引
/// 规则:找到被多数节点复制的最高日志索引
pub fn advance_commit_index(&mut self) {
if self.role != Role::Leader {
return;
}
let leader_state = self.leader_state.as_ref().unwrap();
let mut match_indices: Vec<u64> = leader_state
.match_index
.values()
.copied()
.collect();
match_indices.push(self.last_log().index);
match_indices.sort_by(|a, b| b.cmp(a));
// 多数派中位数即为可提交索引
let majority_idx = self.peers.len() / 2;
let new_commit = match_indices[majority_idx];
// 只能提交当前 Term 的日志(Raft 安全性约束)
if new_commit > self.volatile.commit_index {
if let Some(entry) = self.get_log_entry(new_commit) {
if entry.term == self.persistent.current_term {
self.volatile.commit_index = new_commit;
}
}
}
}
fn step_down(&mut self, new_term: u64) {
self.persistent.current_term = new_term;
self.persistent.voted_for = None;
self.role = Role::Follower;
self.leader_state = None;
}
fn last_log(&self) -> &LogEntry {
self.persistent.log.last().unwrap_or(&LogEntry {
term: 0,
index: 0,
command: Vec::new(),
})
}
fn get_log_entry(&self, index: u64) -> Option<&LogEntry> {
if index == 0 {
return None;
}
self.persistent.log.get((index – 1) as usize)
}
}
这段代码基本照着 Raft 论文的状态定义和 RPC 规则写的。handle_request_vote 里的日志完整性检查,保证了只有日志更完整的候选人才能拿到票;handle_append_entries 里的日志截断逻辑,处理了 Leader 切换后可能出现的日志不一致;advance_commit_index 里“只提交当前 Term 日志”的约束,是 Raft 安全性证明的关键。
四、Raft 的工程代价跟 Paxos 的对比
Raft 的强 Leader 模型确实好理解,但局限性也摆在那儿。
写入瓶颈是最直接的。所有客户端请求都得经过 Leader,Leader 的吞吐量就是整个集群的写入上限。跨数据中心部署的时候,客户端跟 Leader 之间的网络延迟,直接影响写入性能。Paxos 的 Multi-Paxos 变体允许多个节点同时当不同日志槽位的 Leader,理论上写入并行度能更高,但实现复杂度上去了一大截。
Leader 切换的可用性间隙也是个问题。Leader 一挂,集群得经历一轮选举才能恢复写入。选举耗时看超时配置,一般 150ms-300ms。这段时间集群对客户端来说就是不可用的。对可用性要求高的场景,可以用预选举(Pre-Vote)机制,减少不必要的 Term 递增,缩短恢复时间。
日志膨胀是另一个坑。Raft 的日志只增不减,生产环境必须搞快照:把状态机当前状态序列化存起来,然后把快照里的日志条目截掉。快照传输数据量大,慢速网络上 Follower 容易长时间滞后。etcd 的做法是增量快照传输,只发差异部分。
适用边界:Raft 适合中小规模集群(3-9 节点)、强一致性要求、读多写少的场景。大规模集群、跨地域部署、超高写入吞吐的场景,得考虑 Paxos 变体或者 EPaxos 这种去中心化协议。
五、总结
Raft 把共识问题拆成选主、日志复制、安全保证三块,给分布式一致性提供了一个工程上能落地的方案。强 Leader 模型虽然简化了实现,但也带来了写入瓶颈和可用性间隙,得靠预选举、快照传输、线性一致性读这些工程手段来补。
落地建议:先从 etcd/raft-rs 这些成熟实现入手,先验证业务场景跟 Raft 模型匹不匹配,再根据性能瓶颈决定要不要定制优化。核心原则就一个:共识协议的选择不是找理论最优解,而是在工程约束下做务实的权衡。
网硕互联帮助中心






评论前必须登录!
注册