拓冰建站拓冰建站
首页 / 资讯中心 / 正文

第3讲:Raft共识算法(上)——Leader选举与日志复制

前两讲我们实现了节点通信和数据分片但有一个根本问题没解决当节点故障或网络分区时谁来保证数据一致性这一讲我们来实现分布式系统的核心——Raft共识算法。Raft是Paxos的工程化替代以可理解性著称。我们将实现它的两个核心模块Leader选举和日志复制。一、Raft核心概念1.1 节点状态Raft中的节点有三种角色┌─────────────┐ │ Follower │ ← 被动响应接收Leader的日志 └──────┬──────┘ │ 选举超时发起选举 ▼ ┌─────────────┐ │ Candidate │ ← 发起选举拉票 └──────┬──────┘ │ 获得多数票 ▼ ┌─────────────┐ │ Leader │ ← 处理客户端请求复制日志 └─────────────┘1.2 Term任期Term 1 Term 2 Term 3 ├─────┤├──────┤├────────┤ ↑ ↑ Leader Leader 宕机 当选每个Term最多只有一个LeaderTerm是单调递增的逻辑时钟节点通过Term判断谁的信息更新1.3 日志结构Log: ┌────┬────┬────┬────┬────┬────┐ │ 1 │ 2 │ 3 │ 4 │ 5 │ 6 │ ├────┼────┼────┼────┼────┼────┤ │SET │SET │DEL │SET │SET │DEL │ │x3 │y5 │ x │z7 │x8 │ y │ └────┴────┴────┴────┴────┴────┘ ↑ commitIndex: 已提交到多数节点的位置二、Raft核心数据结构# minikv/raft/types.py from enum import Enum, auto from dataclasses import dataclass, field from typing import Any, Dict, List, Optional import time import uuid class NodeRole(Enum): FOLLOWER auto() CANDIDATE auto() LEADER auto() dataclass class LogEntry: 日志条目 index: int 0 # 日志索引 term: int 0 # 创建时的任期 command: str # 命令如 SET, DEL key: str # 操作的 key value: Any None # 操作的 value timestamp: float 0.0 # 时间戳 def __post_init__(self): if not self.timestamp: self.timestamp time.time() dataclass class RaftState: Raft 持久化状态 current_term: int 0 # 当前任期 voted_for: Optional[str] None # 当前任期投票给了谁 log: List[LogEntry] field(default_factorylist) # 快照后续实现 snapshot: Optional[bytes] None snapshot_index: int 0 snapshot_term: int 0 dataclass class RaftVolatileState: Raft 易失状态 commit_index: int 0 # 已提交的最大日志索引 last_applied: int 0 # 已应用到状态机的最大日志索引 # Leader 特有 next_index: Dict[str, int] field(default_factorydict) # 发给每个Follower的下一条日志索引 match_index: Dict[str, int] field(default_factorydict) # 每个Follower已匹配的最高日志索引 dataclass class RequestVoteArgs: 请求投票参数 term: int candidate_id: str last_log_index: int last_log_term: int dataclass class RequestVoteResult: 请求投票结果 term: int vote_granted: bool dataclass class AppendEntriesArgs: 附加日志参数 term: int leader_id: str prev_log_index: int prev_log_term: int entries: List[LogEntry] leader_commit: int dataclass class AppendEntriesResult: 附加日志结果 term: int success: bool conflict_index: int 0 conflict_term: int 0三、Raft核心实现3.1 Raft节点# minikv/raft/raft_node.py import random import threading import time import logging from typing import Dict, List, Optional, Callable from .types import * logger logging.getLogger(__name__) class RaftNode: Raft 共识节点 实现 - Leader 选举 - 日志复制 - 安全性保证 def __init__(self, node_id: str, peers: List[str], election_timeout: tuple (150, 300), heartbeat_interval: float 0.05): Args: node_id: 本节点ID peers: 集群中其他节点ID列表 election_timeout: 选举超时范围(ms) heartbeat_interval: 心跳间隔(秒) self.node_id node_id self.peers peers self.election_timeout_range election_timeout self.heartbeat_interval heartbeat_interval # Raft 状态 self.state RaftState() self.volatile RaftVolatileState() # 角色 self.role NodeRole.FOLLOWER # 选举计时器 self.election_timer None self.election_deadline time.time() self._random_election_timeout() # 状态机由上层实现 self.state_machine: Optional[Callable] None # 消息发送接口由传输层注入 self.send_message: Optional[Callable] None # 控制 self.running False self.lock threading.Lock() logger.info(fRaftNode[{node_id}] initialized with peers: {peers}) def start(self): 启动节点 self.running True # 启动选举超时检测 self._start_election_timer() # 启动应用循环 self._start_apply_loop() logger.info(fRaftNode[{self.node_id}] started as FOLLOWER) def stop(self): 停止节点 self.running False if self.election_timer: self.election_timer.cancel() # Leader 选举 def _start_election_timer(self): 启动选举定时器 def check_election(): while self.running: time.sleep(0.01) # 10ms 精度 self._check_election_timeout() thread threading.Thread(targetcheck_election, daemonTrue) thread.start() def _check_election_timeout(self): 检查选举超时 if self.role NodeRole.LEADER: return if time.time() self.election_deadline: self._start_election() def _start_election(self): 开始选举 with self.lock: self.role NodeRole.CANDIDATE self.state.current_term 1 self.state.voted_for self.node_id logger.info(f开始选举term{self.state.current_term}) # 重置选举超时 self.election_deadline time.time() self._random_election_timeout() # 收集投票 votes_received 1 # 自己投自己 last_log_index len(self.state.log) last_log_term self.state.log[-1].term if self.state.log else 0 args RequestVoteArgs( termself.state.current_term, candidate_idself.node_id, last_log_indexlast_log_index, last_log_termlast_log_term ) # 向所有peer发送投票请求 for peer in self.peers: if self.send_message: self.send_message(peer, RequestVote, args) # 等待投票结果异步处理回复 def handle_request_vote(self, args: RequestVoteArgs, from_node: str) - RequestVoteResult: 处理投票请求 投票规则 1. 如果 args.term currentTerm拒绝 2. 如果 votedFor 为空或等于 candidateId且候选人的日志至少和自己一样新同意 with self.lock: result RequestVoteResult(termself.state.current_term, vote_grantedFalse) # 规则1任期检查 if args.term self.state.current_term: return result # 如果发现更高的任期转为Follower if args.term self.state.current_term: self.state.current_term args.term self.role NodeRole.FOLLOWER self.state.voted_for None # 规则2日志新旧检查 my_last_index len(self.state.log) my_last_term self.state.log[-1].term if self.state.log else 0 log_up_to_date ( args.last_log_term my_last_term or (args.last_log_term my_last_term and args.last_log_index my_last_index) ) can_vote ( (self.state.voted_for is None or self.state.voted_for args.candidate_id) and log_up_to_date ) if can_vote: self.state.voted_for args.candidate_id result.vote_granted True self.election_deadline time.time() self._random_election_timeout() logger.info(f投票给 {args.candidate_id} (term{args.term})) return result def handle_vote_response(self, result: RequestVoteResult, from_node: str): 处理投票回复 with self.lock: if self.role ! NodeRole.CANDIDATE: return # 如果发现更高任期转为Follower if result.term self.state.current_term: self.state.current_term result.term self.role NodeRole.FOLLOWER self.state.voted_for None return if result.vote_granted: # 统计票数 # (简化在实际实现中需要维护一个计票器) # 如果获得多数票成为Leader self._become_leader() def _become_leader(self): 成为Leader self.role NodeRole.LEADER # 初始化Leader状态 last_log_index len(self.state.log) for peer in self.peers: self.volatile.next_index[peer] last_log_index 1 self.volatile.match_index[peer] 0 logger.info(f成为 Leader (term{self.state.current_term})) # 立即发送心跳 self._broadcast_heartbeat() # 启动心跳定时器 self._start_heartbeat() def _random_election_timeout(self) - float: 生成随机的选举超时时间秒 min_ms, max_ms self.election_timeout_range return random.randint(min_ms, max_ms) / 1000.0 # 日志复制 def _start_heartbeat(self): 启动心跳 def heartbeat_loop(): while self.running and self.role NodeRole.LEADER: self._broadcast_heartbeat() time.sleep(self.heartbeat_interval) thread threading.Thread(targetheartbeat_loop, daemonTrue) thread.start() def _broadcast_heartbeat(self): 广播心跳空AppendEntries with self.lock: for peer in self.peers: self._send_append_entries(peer) def _send_append_entries(self, peer: str): 向指定peer发送AppendEntries prev_log_index self.volatile.next_index[peer] - 1 prev_log_term 0 if prev_log_index 0 and prev_log_index len(self.state.log): prev_log_term self.state.log[prev_log_index - 1].term # 获取要发送的日志条目 entries [] next_idx self.volatile.next_index[peer] if next_idx len(self.state.log): entries self.state.log[next_idx - 1:] args AppendEntriesArgs( termself.state.current_term, leader_idself.node_id, prev_log_indexprev_log_index, prev_log_termprev_log_term, entriesentries, leader_commitself.volatile.commit_index ) if self.send_message: self.send_message(peer, AppendEntries, args) def handle_append_entries(self, args: AppendEntriesArgs, from_node: str) - AppendEntriesResult: 处理AppendEntries请求 包含两种场景 1. 心跳entries为空确认Leader存活 2. 日志复制entries不为空追加新日志 with self.lock: result AppendEntriesResult( termself.state.current_term, successFalse ) # 规则1任期检查 if args.term self.state.current_term: return result # 重置选举超时收到合法Leader的消息 self.election_deadline time.time() self._random_election_timeout() # 如果发现更高任期转为Follower if args.term self.state.current_term: self.state.current_term args.term self.role NodeRole.FOLLOWER self.state.voted_for None # 规则2日志一致性检查 if args.prev_log_index 0: if args.prev_log_index len(self.state.log): # 日志不够长 result.conflict_index len(self.state.log) 1 result.conflict_term None return result prev_entry self.state.log[args.prev_log_index - 1] if prev_entry.term ! args.prev_log_term: # 任期冲突 result.conflict_term prev_entry.term # 找到该任期的第一个日志 idx args.prev_log_index - 1 while idx 0 and self.state.log[idx - 1].term result.conflict_term: idx - 1 result.conflict_index idx 1 return result # 规则3追加新日志 if args.entries: # 删除冲突的日志 insert_index args.prev_log_index for i, entry in enumerate(args.entries): log_idx insert_index i if log_idx len(self.state.log): if self.state.log[log_idx - 1].term ! entry.term: # 截断冲突日志 self.state.log self.state.log[:log_idx - 1] self.state.log.append(entry) else: self.state.log.append(entry) # 更新commitIndex if args.leader_commit self.volatile.commit_index: self.volatile.commit_index min( args.leader_commit, len(self.state.log) ) result.success True return result def handle_append_entries_response(self, result: AppendEntriesResult, peer: str): 处理AppendEntries回复 with self.lock: if self.role ! NodeRole.LEADER: return if result.term self.state.current_term: # 发现更高任期降级 self.state.current_term result.term self.role NodeRole.FOLLOWER self.state.voted_for None return if result.success: # 更新匹配状态 self.volatile.match_index[peer] len(self.state.log) self.volatile.next_index[peer] len(self.state.log) 1 # 尝试推进commitIndex self._update_commit_index() else: # 失败回退nextIndex if result.conflict_term: # 找到该任期的最后一个日志 for i in range(len(self.state.log) - 1, -1, -1): if self.state.log[i].term result.conflict_term: self.volatile.next_index[peer] i 1 break else: self.volatile.next_index[peer] result.conflict_index def _update_commit_index(self): 更新commitIndex # 检查是否存在一个索引N使得多数节点的matchIndex N for n in range(self.volatile.commit_index 1, len(self.state.log) 1): if self.state.log[n - 1].term ! self.state.current_term: continue # 只能提交当前任期的日志 count 1 # 自己 for peer in self.peers: if self.volatile.match_index.get(peer, 0) n: count 1 if count len(self.peers) // 2: # 多数派 self.volatile.commit_index n logger.debug(fcommitIndex 推进到 {n}) # 客户端请求 def propose(self, command: str, key: str, value: Any None) - bool: 客户端提议一个操作 只有Leader能接受客户端请求 with self.lock: if self.role ! NodeRole.LEADER: return False # 创建日志条目 entry LogEntry( indexlen(self.state.log) 1, termself.state.current_term, commandcommand, keykey, valuevalue ) self.state.log.append(entry) # 广播给所有Follower for peer in self.peers: self._send_append_entries(peer) return True def _start_apply_loop(self): 启动应用循环将已提交的日志应用到状态机 def apply_loop(): while self.running: time.sleep(0.001) with self.lock: while self.volatile.last_applied self.volatile.commit_index: self.volatile.last_applied 1 entry self.state.log[self.volatile.last_applied - 1] if self.state_machine: self.state_machine(entry.command, entry.key, entry.value) thread threading.Thread(targetapply_loop, daemonTrue) thread.start() # 状态查询 def get_state(self) - dict: 获取Raft状态 with self.lock: return { node_id: self.node_id, role: self.role.name, current_term: self.state.current_term, log_size: len(self.state.log), commit_index: self.volatile.commit_index, last_applied: self.volatile.last_applied, } def is_leader(self) - bool: 判断是否为Leader return self.role NodeRole.LEADER四、Raft集成到传输层4.1 Raft消息处理# minikv/raft/raft_service.py import threading import logging from typing import Dict, Optional from .raft_node import RaftNode from .types import * from ..transport.message import Message, MessageType logger logging.getLogger(__name__) class RaftService: Raft 服务层 将 Raft 节点与传输层集成 def __init__(self, node_id: str, peers: List[str], transport): self.node_id node_id self.transport transport # 创建 Raft 节点 self.raft RaftNode(node_id, peers) # 设置消息发送接口 self.raft.send_message self._send_raft_message # 注册消息处理器 self.transport.message_handler self._handle_message def start(self): 启动服务 self.raft.start() logger.info(Raft service started) def stop(self): 停止服务 self.raft.stop() def _send_raft_message(self, target: str, msg_type: str, args): 发送 Raft 消息 message Message( msg_typeself._map_raft_msg_type(msg_type), sender_idself.node_id, receiver_idtarget, bodyself._serialize_args(msg_type, args) ) self.transport.send_to(target, message) def _handle_message(self, msg: Message, conn): 处理收到的消息 if msg.msg_type MessageType.REQUEST_VOTE: args RequestVoteArgs(**msg.body) result self.raft.handle_request_vote(args, msg.sender_id) reply msg.reply(self._serialize_result(RequestVote, result)) conn.send(reply) elif msg.msg_type MessageType.VOTE_RESPONSE: result RequestVoteResult(**msg.body) self.raft.handle_vote_response(result, msg.sender_id) elif msg.msg_type MessageType.APPEND_ENTRIES: args AppendEntriesArgs(**msg.body) result self.raft.handle_append_entries(args, msg.sender_id) reply msg.reply(self._serialize_result(AppendEntries, result)) conn.send(reply) elif msg.msg_type MessageType.APPEND_RESPONSE: result AppendEntriesResult(**msg.body) self.raft.handle_append_entries_response(result, msg.sender_id) def _map_raft_msg_type(self, raft_type: str) - MessageType: 映射 Raft 消息类型到传输层消息类型 mapping { RequestVote: MessageType.REQUEST_VOTE, VoteResponse: MessageType.VOTE_RESPONSE, AppendEntries: MessageType.APPEND_ENTRIES, AppendResponse: MessageType.APPEND_RESPONSE, } return mapping.get(raft_type, MessageType.PING) def _serialize_args(self, msg_type: str, args) - dict: 序列化参数 if isinstance(args, RequestVoteArgs): return { term: args.term, candidate_id: args.candidate_id, last_log_index: args.last_log_index, last_log_term: args.last_log_term } elif isinstance(args, AppendEntriesArgs): return { term: args.term, leader_id: args.leader_id, prev_log_index: args.prev_log_index, prev_log_term: args.prev_log_term, entries: [vars(e) for e in args.entries], leader_commit: args.leader_commit } return {} def _serialize_result(self, msg_type: str, result) - dict: 序列化结果 if isinstance(result, RequestVoteResult): return {term: result.term, vote_granted: result.vote_granted} elif isinstance(result, AppendEntriesResult): return { term: result.term, success: result.success, conflict_index: result.conflict_index, conflict_term: result.conflict_term } return {}五、完整演示# examples/raft_demo.py import time import threading import logging import sys logging.basicConfig( levellogging.INFO, format%(asctime)s [%(levelname)s] %(name)s: %(message)s ) sys.path.insert(0, ..) from minikv.raft.raft_node import RaftNode, NodeRole class SimulatedTransport: 模拟传输层用于演示 def __init__(self): self.nodes {} self.message_log [] def register(self, node_id: str, raft: RaftNode): self.nodes[node_id] raft raft.send_message lambda target, msg_type, args: self._deliver( node_id, target, msg_type, args ) def _deliver(self, sender: str, target: str, msg_type: str, args): self.message_log.append(f{sender} → {target}: {msg_type}) if target in self.nodes: target_raft self.nodes[target] if msg_type RequestVote: result target_raft.handle_request_vote(args, sender) # 模拟异步回复 threading.Timer(0.01, self._deliver_vote_response, args[sender, target, result]).start() elif msg_type AppendEntries: result target_raft.handle_append_entries(args, sender) threading.Timer(0.01, self._deliver_append_response, args[sender, target, result]).start() def _deliver_vote_response(self, sender, target, result): if sender in self.nodes: self.nodes[sender].handle_vote_response(result, target) def _deliver_append_response(self, sender, target, result): if sender in self.nodes: self.nodes[sender].handle_append_entries_response(result, target) def demo_raft_election(): 演示Leader选举 print( * 60) print(️ Raft Leader 选举演示) print( * 60) transport SimulatedTransport() # 创建3节点集群 nodes {} for i in range(1, 4): node_id fnode-{i} peers [fnode-{j} for j in range(1, 4) if j ! i] raft RaftNode( node_idnode_id, peerspeers, election_timeout(100, 200), # 快速选举 heartbeat_interval0.03 ) transport.register(node_id, raft) raft.start() nodes[node_id] raft # 等待选举完成 print(\n⏳ 等待Leader选举...) time.sleep(1) # 查看谁成为了Leader leader None for node_id, raft in nodes.items(): state raft.get_state() print(f {node_id}: role{state[role]}, term{state[current_term]}) if state[role] LEADER: leader node_id print(f\n 选举结果: Leader {leader}) # 测试日志复制 if leader: print(\n 测试日志复制:) leader_raft nodes[leader] # 提交一些操作 leader_raft.propose(SET, name, Alice) leader_raft.propose(SET, age, 30) leader_raft.propose(SET, city, Beijing) time.sleep(0.5) print(\n 各节点日志状态:) for node_id, raft in nodes.items(): state raft.get_state() print(f {node_id}: log_size{state[log_size]}, fcommit{state[commit_index]}, applied{state[last_applied]}) # 测试Leader故障 if leader: print(f\n 模拟Leader故障: {leader}) nodes[leader].stop() time.sleep(1.5) print(\n 重新选举:) for node_id, raft in nodes.items(): if raft.running: state raft.get_state() print(f {node_id}: role{state[role]}, term{state[current_term]}) if state[role] LEADER: print(f 新Leader: {node_id}) # 清理 for raft in nodes.values(): raft.stop() def demo_log_replication(): 演示日志复制细节 print(\n * 60) print( 日志复制细节演示) print( * 60) transport SimulatedTransport() # 创建集群 nodes {} for i in range(1, 4): node_id fnode-{i} peers [fnode-{j} for j in range(1, 4) if j ! i] raft RaftNode(node_id, peers, election_timeout(200, 400)) transport.register(node_id, raft) raft.start() nodes[node_id] raft time.sleep(1) # 找到Leader leader_id None for node_id, raft in nodes.items(): if raft.is_leader(): leader_id node_id break if leader_id: leader nodes[leader_id] print(f\nLeader: {leader_id}) print(\n逐步提交日志:) for i in range(5): cmd fSET key{i}value{i} leader.propose(SET, fkey{i}, fvalue{i}) time.sleep(0.1) print(f 提交 [{i1}/5]: {cmd}) time.sleep(0.3) print(\n最终日志状态:) for node_id, raft in nodes.items(): state raft.get_state() print(f {node_id}: {state[log_size]} 条日志, f已提交到 {state[commit_index]}) for raft in nodes.values(): raft.stop() if __name__ __main__: demo_raft_election() demo_log_replication()六、测试# tests/test_raft.py import unittest import threading import time from minikv.raft.raft_node import RaftNode, NodeRole from minikv.raft.types import * class TestRaftElection(unittest.TestCase): Raft选举测试 def setUp(self): self.nodes {} self.transport self._create_transport() for i in range(1, 4): node_id fnode-{i} peers [fnode-{j} for j in range(1, 4) if j ! i] raft RaftNode(node_id, peers, election_timeout(50, 100), heartbeat_interval0.02) self.transport.register(node_id, raft) raft.start() self.nodes[node_id] raft time.sleep(0.5) def tearDown(self): for raft in self.nodes.values(): raft.stop() def _create_transport(self): class TestTransport: def __init__(self): self.nodes {} def register(self, node_id, raft): self.nodes[node_id] raft raft.send_message lambda t, mt, a: self._deliver(node_id, t, mt, a) def _deliver(self, sender, target, msg_type, args): if target in self.nodes: target_raft self.nodes[target] if msg_type RequestVote: result target_raft.handle_request_vote(args, sender) threading.Timer(0.005, self._deliver_vote_response, args[sender, target, result]).start() elif msg_type AppendEntries: result target_raft.handle_append_entries(args, sender) threading.Timer(0.005, self._deliver_append_response, args[sender, target, result]).start() def _deliver_vote_response(self, sender, target, result): if sender in self.nodes: self.nodes[sender].handle_vote_response(result, target) def _deliver_append_response(self, sender, target, result): if sender in self.nodes: self.nodes[sender].handle_append_entries_response(result, target) return TestTransport() def test_elect_one_leader(self): 测试集群恰好选举出一个Leader leaders sum(1 for raft in self.nodes.values() if raft.is_leader()) self.assertEqual(leaders, 1) def test_leader_heartbeat(self): 测试Leader发送心跳 leader [r for r in self.nodes.values() if r.is_leader()][0] followers [r for r in self.nodes.values() if not r.is_leader()] # Leader应该定期发送心跳保持Followers的选举超时重置 time.sleep(0.2) # Followers不应该超时发起选举 for f in followers: self.assertFalse(f.is_leader()) def test_log_replication(self): 测试日志复制 leader [r for r in self.nodes.values() if r.is_leader()][0] # 提交日志 leader.propose(SET, x, 1) leader.propose(SET, y, 2) time.sleep(0.3) # 所有节点应该有相同的日志 log_sizes [len(r.state.log) for r in self.nodes.values()] self.assertEqual(len(set(log_sizes)), 1) class TestRaftSafety(unittest.TestCase): Raft安全性测试 def test_election_safety(self): 测试选举安全性每个Term最多一个Leader # 模拟网络分区 pass # 后续实现 def test_log_matching(self): 测试日志匹配特性 pass # 后续实现 if __name__ __main__: unittest.main()七、总结这一讲我们实现了Raft共识算法的核心组件功能Leader选举​随机超时、投票机制、Term管理日志复制​AppendEntries、一致性检查、冲突解决安全性​日志匹配、Leader完整性、多数派提交状态机​已提交日志自动应用到状态机关键成果3节点集群能自动选出LeaderLeader故障后能重新选举日志能可靠复制到所有节点通过多数派保证一致性下一讲我们将完成Raft的剩余部分——安全性证明、成员变更、日志压缩与快照。开发之余的小工具推荐​处理 Base64、JWT 解析、JSON 格式化、Crontab 计算、PDF 合并压缩这些碎片需求我常用一个纯前端本地工具箱zz365.top子页 PDF 大师PDF 大师 - zz365工具箱。所有计算在浏览器完成文件不上传服务器关页即清。免费、无登录、无广告适合开发者当常驻标签页。
分享:

看完干货,该让你的企业上线了

免费需求沟通 · 48 小时内出具建站方案 · 河南本地可上门