尧图建网站 尧图建网站 YAOTU WEB BUILD 免费咨询
ARTICLE DETAIL

资讯详情

深耕网站建设与建站编程的一线实战洞察。

Day50-Raft协议:Leader选举、日志复制与Java实现方式

Day50-Raft协议:Leader选举、日志复制与Java实现方式 当库存服务出现了两个 Leader同一批订单被扣了两次库存时。 你是不是看着ZooKeeper 的会话超时日志发呆为什么会这样——不是代码 bug是分布式共识没选对。直到遇见 Raft我才明白一致性协议也可以像讲故事一样被理解。今天这篇文章用一个能跑起来的 Java 代码把 Leader 选举、日志复制、安全性这三个核心环节一次讲透。一、Raft 到底解决了什么问题分布式系统里有条铁律网络是不可靠的机器是会挂的。当集群被网络切成两半脑裂两边都可能继续接收写请求数据就分叉了。Raft 的做法很朴素任何时刻集群里最多只有一个 Leader所有写操作都必须经过 LeaderLeader 提交日志前必须让多数派节点确认。这里隐藏着 Raft 的灵魂用“多数派”这个数学保证把不确定的网络变成确定的承诺。只要有超过半数节点存活系统就能对外服务。三个核心状态角色职责触发条件Follower被动接收 Leader 心跳复制日志默认状态Candidate发起选举拉票超时没收到心跳Leader处理客户端写请求分发日志赢得多数票二、手写一个教学版 Raft 节点下面这段代码是整个教学实现的地基节点状态、日志条目、任期管理。它依赖纯 JDK复制下来就能跑。import java.util.*; import java.util.concurrent.*; import java.util.concurrent.atomic.AtomicInteger; /** * 教学版 Raft 节点仅保留核心共识逻辑不含网络 RPC。 * JDK 版本17 */ public class RaftNode { enum State { FOLLOWER, CANDIDATE, LEADER } // 日志条目任期 操作指令生产环境会包含 index、term、command、clientId 等 record LogEntry(int term, String command) {} private final int id; // 节点 ID private final ListRaftNode peers; // 集群其他节点 private volatile State state State.FOLLOWER; private volatile int currentTerm 0; // 当前任期单调递增 private volatile Integer votedFor null; // 本轮投给谁了 private final ListLogEntry log new ArrayList(); // 提交索引 已应用索引 private volatile int commitIndex 0; private volatile int lastApplied 0; // Leader 额外维护每个节点的下一个日志索引和匹配索引 private final MapInteger, Integer nextIndex new ConcurrentHashMap(); private final MapInteger, Integer matchIndex new ConcurrentHashMap(); private final ScheduledExecutorService scheduler Executors.newScheduledThreadPool(2); private final Random random new Random(); public RaftNode(int id, ListRaftNode peers) { this.id id; this.peers peers; } public int id() { return id; } public State state() { return state; } public int currentTerm() { return currentTerm; } public ListLogEntry log() { return log; } public int commitIndex() { return commitIndex; } public void start() { // 随机选举超时150~300ms避免活锁 resetElectionTimer(); } private void resetElectionTimer() { int timeout 150 random.nextInt(150); scheduler.schedule(this::onElectionTimeout, timeout, TimeUnit.MILLISECONDS); } // 选举超时触发 private synchronized void onElectionTimeout() { if (state State.LEADER) return; becomeCandidate(); resetElectionTimer(); } private void becomeCandidate() { currentTerm; votedFor id; state State.CANDIDATE; System.out.printf(Node %d 成为 Candidate任期 %d%n, id, currentTerm); int votesNeeded (peers.size() 1) / 2 1; AtomicInteger votes new AtomicInteger(1); // 先投自己 int lastLogIndex log.size(); int lastLogTerm lastLogIndex 0 ? log.get(lastLogIndex - 1).term() : 0; for (RaftNode peer : peers) { peer.requestVote(currentTerm, id, lastLogIndex, lastLogTerm) .ifPresent(granted - { if (granted) { if (votes.incrementAndGet() votesNeeded state State.CANDIDATE) { becomeLeader(); } } }); } } private void becomeLeader() { state State.LEADER; System.out.printf(Node %d 成为 Leader任期 %d%n, id, currentTerm); for (RaftNode peer : peers) { nextIndex.put(peer.id(), log.size() 1); matchIndex.put(peer.id(), 0); } // 立刻并周期性发送心跳空 AppendEntries scheduler.scheduleAtFixedRate(this::sendHeartbeats, 0, 50, TimeUnit.MILLISECONDS); } // ---- 下一段代码会实现这里的方法 ---- public OptionalBoolean requestVote(int term, int candidateId, int lastLogIndex, int lastLogTerm) { return null; } public void appendEntries(int term, int leaderId, int prevLogIndex, int prevLogTerm, ListLogEntry entries, int leaderCommit) {} public void sendHeartbeats() {} }这段代码刻意省略了真实网络用直接方法调用模拟 RPC。教学场景中这样更容易看清逻辑真正上线请用 gRPC 或 Netty并加上持久化 WAL。三、Leader 选举一次“少数服从多数”的投票选举的关键规则只有两条任期优先Candidate 的任期必须大于等于本地任期日志不旧于自己Candidate 的最后一条日志任期不能比本地小任期相同则索引不能比本地小。/** * 处理 RequestVote RPC。 * 返回 true 表示投给 candidate。 */ public synchronized OptionalBoolean requestVote(int term, int candidateId, int lastLogIndex, int lastLogTerm) { // 发现更高任期立刻退位 if (term currentTerm) { currentTerm term; state State.FOLLOWER; votedFor null; } if (term currentTerm) { return Optional.of(false); // 过期候选人 } int myLastIndex log.size(); int myLastTerm myLastIndex 0 ? log.get(myLastIndex - 1).term() : 0; boolean logOk (lastLogTerm myLastTerm) || (lastLogTerm myLastTerm lastLogIndex myLastIndex); boolean canVote votedFor null || votedFor candidateId; if (canVote logOk) { votedFor candidateId; state State.FOLLOWER; System.out.printf(Node %d 在任期 %d 投给 Node %d%n, id, currentTerm, candidateId); return Optional.of(true); } return Optional.of(false); }注意synchronized真实实现不会这么简单但教学代码中用它保证任期、投票、日志的原子性观察避免竞态让人头晕。四、日志复制把“承诺”同步到多数派Leader 收到客户端写请求后先把日志追加到本地然后发送AppendEntries给所有 Follower。只有当多数派包括 Leader 自己都确认后这条日志才被视为“已提交”然后应用到状态机。/** * 处理 AppendEntries RPC心跳 日志复制。 */ public synchronized void appendEntries(int term, int leaderId, int prevLogIndex, int prevLogTerm, ListLogEntry entries, int leaderCommit) { if (term currentTerm) { return; // 拒绝过期 Leader } if (term currentTerm) { currentTerm term; votedFor null; } state State.FOLLOWER; // 收到合法 Leader 心跳退位或保持 Follower resetElectionTimer(); // 1. 校验 prevLog 是否匹配 if (prevLogIndex 0) { if (log.size() prevLogIndex) { return; // 日志太短让 Leader 回退 nextIndex } int actualTerm log.get(prevLogIndex - 1).term(); if (actualTerm ! prevLogTerm) { // 任期冲突删除不一致的日志 while (log.size() prevLogIndex) { log.remove(log.size() - 1); } return; } } // 2. 追加新日志 int insertPos prevLogIndex; for (LogEntry entry : entries) { if (insertPos log.size()) { if (log.get(insertPos).term() ! entry.term()) { // 冲突条目截断 while (log.size() insertPos) { log.remove(log.size() - 1); } log.add(entry); } } else { log.add(entry); } insertPos; } // 3. 提交索引推进 if (leaderCommit commitIndex) { commitIndex Math.min(leaderCommit, log.size()); applyCommitted(); } } /** 应用到状态机简化打印 */ private void applyCommitted() { while (lastApplied commitIndex) { lastApplied; System.out.printf(Node %d 应用日志 #%d: %s%n, id, lastApplied, log.get(lastApplied - 1)); } } /** Leader 发送心跳与复制 */ public synchronized void sendHeartbeats() { if (state ! State.LEADER) return; for (RaftNode peer : peers) { int nextIdx nextIndex.getOrDefault(peer.id(), log.size() 1); int prevLogIndex nextIdx - 1; int prevLogTerm prevLogIndex 0 ? log.get(prevLogIndex - 1).term() : 0; // 取 prevLogIndex 之后的日志作为 entries ListLogEntry entries new ArrayList(); for (int i nextIdx - 1; i log.size(); i) { entries.add(log.get(i)); } peer.appendEntries(currentTerm, id, prevLogIndex, prevLogTerm, entries, commitIndex); } // Leader 推进 commitIndex找到 matchIndex 中多数派确认的索引 updateLeaderCommit(); } private void updateLeaderCommit() { ListInteger match new ArrayList(); match.add(log.size()); // Leader 自己 match.addAll(matchIndex.values()); match.sort(Comparator.reverseOrder()); int majority (peers.size() 1) / 2 1; int n match.get(majority - 1); // 第 N 大的匹配索引 // Raft 安全性只能提交当前任期的日志 if (n commitIndex log.get(n - 1).term() currentTerm) { commitIndex n; applyCommitted(); } }五、安全性Raft 比 Paxos 好在哪里Paxos 是共识协议的祖师爷但出了名的难懂。Raft 的设计者 Diego Ongaro 说了一句话我记了很久“Raft 和 Paxos 解决的问题相同但 Raft 被设计成更容易被工程师理解和实现。”两者本质区别维度PaxosRaft学习曲线陡峭Multi-Paxos 细节多平缓状态机清晰Leader 概念隐式需额外推导显式三种角色直接映射日志顺序允许间隙和乱序强有序AppendEntries 顺序保证工程实现容易出错更容易做对Raft 的两个安全铁律选举限制只有日志足够新的节点才能当选 Leader。这保证了新 Leader 一定包含所有已提交的日志。提交限制Leader 只能提交自己任期内的日志。旧任期的日志必须跟随新任期日志一起被多数派确认后才会间接提交。第二个限制非常反直觉但它是 Raft 的灵魂补丁。少了它一个短暂旧任期的 Leader 可能已经“复制到多数派”的日志会被新任期的 Leader 覆盖从而破坏一致性。六、跑起来一个三节点最小集群下面给你一个极简的启动脚本把上面的类串起来。你可以把requestVote、appendEntries、sendHeartbeats补全后运行public class RaftDemo { public static void main(String[] args) throws InterruptedException { ListRaftNode nodes new ArrayList(); for (int i 1; i 3; i) { nodes.add(new RaftNode(i, Collections.emptyList())); } // 互相注入 peer 引用 for (int i 0; i 3; i) { ListRaftNode peers new ArrayList(nodes); peers.remove(i); nodes.set(i, new RaftNode(i 1, peers)); } nodes.forEach(RaftNode::start); TimeUnit.SECONDS.sleep(2); RaftNode leader nodes.stream() .filter(n - n.state() RaftNode.State.LEADER) .findFirst() .orElseThrow(); // 模拟客户端写请求 leader.log().add(new RaftNode.LogEntry(leader.currentTerm(), set foobar)); TimeUnit.SECONDS.sleep(1); nodes.forEach(n - System.out.printf(Node %d commitIndex%d log%s%n, n.id(), n.commitIndex(), n.log())); } }这段启动类为了展示流程做了简化真实集群需要把直接方法调用换成 RPC并把日志持久化到磁盘。建议不要自己手写 Raft 上生产。etcd、Consul、Nacos 的 Raft 实现都经历过大规模打磨。手写是学习选型才是正经事。先深刻理解协议再判断框架源码里的设计取舍。多数派不是“过半数节点”而是“过半数投票权”。生产环境可以配置加权节点也可以配置 Learner只复制不投票做只读副本。搞错这个脑裂会来得比你想的快。日志一定要 WAL 持久化。网络分区恢复后节点重启如果没有持久化日志它会以为自己是一张白纸然后被Leader同步时可能把旧状态带回来。内存版 Raft 仅供学习。共识不是让所有人一样而是让多数人承诺一样Raft 用任期、多数派、日志复制这三板斧把复杂的共识问题拆成了能讲清楚的故事。记住一句话共识协议保证的不是绝对正确而是在混乱中给出一个确定性的答案。下篇预告Day 51《分布式ID生成方案终极对比雪花算法/美团Leaf/滴滴Tinyid》系列文章回顾90篇JAVA高级工程师深度进阶文章从实践案例、项目代码出发用事实说话版权声明本文为「老梁」原创出品90天Java后端AI系列第50篇转载请注明出处。
返回列表