尧图网络 高端网站定制 · 原创设计
免费咨询热线
400-888-6620
免费获取方案
从零实现基于Raft的KV存储:日志存储、选举与复制实战
简介一套基于Raft算法的分布式键值存储系统完整实现面向分布式系统初学者、计算机专业学生及需要设计与实现一致性服务的开发者。项目以Java为主要语言通过Raft协议保障节点间日志一致与高可用覆盖领导者选举、日志复制、成员变更、节点状态管理等核心模块并借助Netty构建节点间RPC通信帮助读者理解分布式一致性的实际落地方法与底层细节。压缩包共161个文件其中149个Java源码文件构成系统主体另有properties与xml配置文件、单元测试及说明文档整体仅209KB文件数量适中且结构清晰便于逐模块通读。目前已有265人学习浏览资源包含日志存储、节点管理、RPC通信等完整工程代码可作为课程设计、毕业设计或系统学习Raft算法的参考资料既能对照算法原理梳理实现思路也能在此基础上扩展功能或开展二次开发。1. 为什么把 Raft 共识塞进 KV 存储比想象中难面试聊 Raft 时大家都能把论文里的角色转换图背出来真正动手把 Raft 塞进一个可用的键值存储时卡住的第一步往往是日志落盘。这套项目源码我拆过一遍核心代码量不大但在NodeImpl、FileLogStore、KVDatabaseImpl几个类之间来回切换时你会发现分布式系统的难点根本不在算法本身而在「日志怎么存」「消息怎么传」「状态机怎么应用」这三件事的耦合关系。项目本身是一个完整可运行的基于 Raft 的 KV 存储支持节点角色切换、日志复制、Netty 通信和简单的 GET/PUT 操作。如果你正想自己写一个简化版 Raft 或者需要一份能跑通的参考实现这套代码比读十篇论文都实用。2. Raft 选举与日志复制从角色状态机到节点间协商2.1 三个角色在 NodeImpl 里如何落地NodeImpl.java是核心状态机所在。实际代码里并没有用复杂的状态模式而是用一个简单的NodeRole枚举加几个状态变量来维护public class NodeImpl { // 当前角色FOLLOWER, CANDIDATE, LEADER private volatile NodeRole role NodeRole.FOLLOWER; // 当前任期每次收到更大的任期号立即更新 private volatile long currentTerm 0; // 当前任期投票给了谁null 表示未投 private volatile String votedFor null; // 日志条目LogEntry 包含 index、term、operation、value private final ListLogEntry log new CopyOnWriteArrayList(); // 选举定时器 private ScheduledFuture? electionTimer; // 心跳定时器 private ScheduledFuture? heartbeatTimer; }这段代码的关键不是字段本身而是volatile的使用。currentTerm和votedFor可能被 RPC 工作线程和定时器线程同时读写volatile保证了可见性。CopyOnWriteArrayList是为了让日志读取在遍历期间不被并发修改干扰虽然写成本高但在节点日志量不大时很合适。角色转换的触发点有三个选举超时后从 Follower 变成 Candidate收到多数派投票后从 Candidate 变成 Leader发现更高任期时任何角色都退回 Follower。这些判断逻辑全部收敛在NodeImpl.changeRole()方法里而不是散落在各处排错时只需要盯着一个方法看。2.2 日志复制流程与提交判定Leader 接收到客户端写请求后会把操作封装成LogEntry追加到自己的日志里然后通过AppendEntriesRPC 广播给所有 Follower。这个项目的处理逻辑简化成下面这样// 处理 AppendEntries 请求返回是否成功 public boolean appendEntries(long leaderTerm, int prevLogIndex, int prevLogTerm, ListLogEntry entries, long leaderCommit) { // 1. 任期校验如果 leader 任期小于当前任期拒绝 if (leaderTerm currentTerm) { return false; } // 2. 一致性检查前一日志必须存在且任期匹配 if (prevLogIndex 0 !matchLog(prevLogIndex, prevLogTerm)) { return false; } // 3. 处理冲突本地日志与 leader 日志冲突时截断本地日志 for (LogEntry entry : entries) { LogEntry local getLogEntry(entry.getIndex()); if (local ! null local.getTerm() ! entry.getTerm()) { // 删除当前条目及其后所有日志 truncateFrom(entry.getIndex()); } // 追加新条目 log.add(entry); } // 4. 更新提交索引 if (leaderCommit commitIndex) { commitIndex Math.min(leaderCommit, lastIndex()); } return true; }注意这里有两个小细节。matchLog不仅要看 index还要看 term否则会出现“日志长度一样但内容不同”的问题。截断日志后必须连带日索引和文件日志一起清理只清内存列表会导致重启后旧日志重新加载时出现游标错位。提交判定是由 Leader 发起的Leader 在每次成功复制到多数派后推进commitIndex并在下一条AppendEntries里携带leaderCommitFollower 收到后才会应用日志到状态机。参数含义整理如下参数类型说明leaderTermlong发送方任期用于拒绝旧 Leader 消息prevLogIndexint前一条日志的索引必须存在才能继续prevLogTermlong前一条日志的任期防止不同任期的相同索引冲突entriesList待复制的新日志条目可为空表示心跳leaderCommitlong提交索引Follower 据此提交本地日志2.3 选举超时与心跳参数的实测调法在本地模拟三节点时我把心跳间隔设成 100ms选举超时随机在 150~300ms 之间// 心跳固定间隔 int heartbeatIntervalMs 100; // 选举超时随机范围避免多个节点同时超时导致选票分裂 int electionTimeoutMin 150; int electionTimeoutMax 300; private long electionTimeoutMs() { return ThreadLocalRandom.current().nextLong(electionTimeoutMin, electionTimeoutMax 1); }这个范围不是拍脑袋定的。如果选举超时随机器范围过窄三节点可能同时变成 Candidate谁都拿不到多数票日志会一直刷 Election timeout如果范围过宽节点故障恢复后要等很久才能触发重新选举。之前测试时把最小值设成 150ms三个节点在同一机器上运行网络延迟本身就会造成时间差分裂概率很低。提示在真实网络环境里节点间 RTT 可能超过 50ms建议把选举超时下限设在 300ms 以上否则频繁的 pre-vote 会把 Raft 变成“心跳选举”。3. 日志存储层FileLogStore 和 EntryIndexFile 的配合3.1 内存日志与文件日志的边界AbstractLogStore定义了日志存储的抽象接口LogImpl和FileLogStore分别对应内存和文件实现。很多人一开始会把整个日志全放内存节点重启后直接从快照恢复这是个陷阱Raft 需要保证已提交日志在崩溃后不丢失内存日志一旦进程退出就全没了。但这个项目有意思的是LogImpl并不是没用的它在运行期负责缓存热数据FileLogStore负责把每条日志刷到磁盘两者之间用EntryIndexFile做索引。我建议你把它理解成两级存储内存里维护完整日志列表写入时同步追加到文件读取时优先从内存返回。真正的持久化文件只需要记录日志条目不承载读取路径这样磁盘 IO 可以用顺序写优化读请求仍然走内存延迟极低。3.2 日志条目的序列化设计文件日志存储的关键在于序列化格式。项目里使用的是自描述的「长度前缀」协议二进制结构如下public byte[] serialize(LogEntry entry) throws IOException { ByteArrayOutputStream byteOut new ByteArrayOutputStream(); DataOutputStream dataOut new DataOutputStream(byteOut); dataOut.writeLong(entry.getIndex()); // 8 字节日志索引 dataOut.writeLong(entry.getTerm()); // 8 字节任期号 byte[] opBytes entry.getOperation().name().getBytes(StandardCharsets.UTF_8); dataOut.writeInt(opBytes.length); // 4 字节操作名长度 dataOut.write(opBytes); // 操作名如 PUT / DELETE byte[] valueBytes entry.getValue().getBytes(StandardCharsets.UTF_8); dataOut.writeInt(valueBytes.length); // 4 字节值长度 dataOut.write(valueBytes); // 值内容 dataOut.flush(); return byteOut.toByteArray(); }为什么不用 Java 自带的ObjectOutputStream两个原因。第一ObjectOutputStream会把类名、类版本号等元数据写进去日志文件无法跨语言解析第二它的写入头部没有长度信息追加日志后如果想随机读取某条日志必须逐个读才知道边界在哪里。上面这种格式每条日志长度可以直接算出来8 8 4 操作名长度 4 值长度配合索引文件可以做到 O(1) 定位。反序列化时注意一点DataInputStream读writeUTF字符串有 65535 字节上限所以这里手写长度前缀避免 value 一长就抛异常。这属于那种不跑大数据量永远发现不了的坑。3.3 索引文件为什么不能省如果只有日志文件定位第 1000 条日志需要从头读取 999 条节点重启后状态机恢复会非常慢。EntryIndexFile相当于一张映射表日志索引 - (文件偏移量, 日志长度)。每次追加日志时只需要记录偏移和长度public class EntryIndexFile { private final MapLong, IndexEntry indexMap new ConcurrentHashMap(); public void append(long index, long offset, int length) { indexMap.put(index, new IndexEntry(offset, length)); } // 读取第 index 条日志时直接从文件偏移量开始读 length 个字节 public byte[] read(FileChannel channel, long index) throws IOException { IndexEntry entry indexMap.get(index); if (entry null) { return null; } ByteBuffer buffer ByteBuffer.allocate(entry.length); channel.read(buffer, entry.offset); return buffer.array(); } }索引文件在 Raft 场景下还有一个用途日志截断。当 Follower 发现日志冲突时需要删除从 conflictIndex 开始的所有日志此时只需要在indexMap里删除对应 key不必真的去磁盘文件上抹掉数据。旧数据会在后续日志压缩或快照生成时被整体重写覆盖。方案追加耗时随机读取耗时重启恢复时间复杂度纯内存日志O(1)O(1)丢失全部日志低顺序文件 全量扫描O(1)O(N)慢需从头读中顺序文件 索引O(1)O(1)快按索引跳读中高这个项目选的是第三种。文件日志只有一条追加路径没有随机写所以磁盘性能可以达到很高的吞吐而索引文件本身很小几百 MB 日志的索引才几 MB可以常驻内存。4. 基于 Netty 的 RPC 通信与节点组管理4.1 ChannelGroup 维护节点连接节点之间的消息传递没有用 HTTP而是直接跑 Netty TCP。ChannelGroup.java在这里的作用不只是收集连接的 Channel还承担了「节点 ID 到连接」的映射职责。看代码能发现它更像是 Channel 注册表public class ChannelGroup { private final ConcurrentMapString, Channel channels new ConcurrentHashMap(); public void addChannel(String nodeId, Channel channel) { Channel old channels.put(nodeId, channel); if (old ! null old ! channel) { old.close(); // 同一个节点重复建连时关闭旧连接 } } public void removeChannel(String nodeId) { Channel removed channels.remove(nodeId); if (removed ! null) { removed.close(); } } public void send(String nodeId, Object message) { Channel channel channels.get(nodeId); if (channel ! null channel.isActive()) { channel.writeAndFlush(message); } } }这里值得注意的地方是old.close()。在分布式环境里节点可能因为网络抖动重连如果不关闭旧连接Leader 会给同一个 Follow 发两套心跳Follower 处理消息时会因为消息乱序导致日志重复或任期混乱。消息体我建议用ByteBuf包装的自定义协议而不是 Java 序列化。这个项目实际走的也是二进制协议第一个字节表示消息类型0x01 选举投票0x02 AppendEntries0x03 客户端请求后面跟具体的 JSON 或自定义二进制载荷。异网传输时不要依赖ObjectOutputStream一旦两端版本不一致连接直接挂断。4.2 消息处理链与 EntryGenerationHandlerEntryGenerationHandler.java看起来像是处理日志条目生成的 Handler实际上它在 RPC 链路里充当了「协议解析和分发」的角色。Netty 中每个入站消息都会经过 Pipeline 里的各个 Handler在这个项目里它大致是这样的public class EntryGenerationHandler extends ChannelInboundHandlerAdapter { private final NodeImpl raftNode; Override public void channelRead(ChannelHandlerContext ctx, Object msg) { if (msg instanceof RaftMessage) { RaftMessage message (RaftMessage) msg; switch (message.getType()) { case REQUEST_VOTE: raftNode.handleRequestVote(message); break; case APPEND_ENTRIES: raftNode.handleAppendEntries(message); break; case CLIENT_PUT: // 客户端请求直接进入日志生成流程 raftNode.propose(logEntryFromMessage(message)); break; default: ctx.fireChannelRead(msg); } } else { ctx.fireChannelRead(msg); } } }channelRead的第一个参数是ChannelHandlerContext不是Channel。很多初学者会在这里直接写ctx.writeAndFlush回包容易把响应发到错误的连接上。正确做法是在消息里携带fromNodeId和requestId回包时根据fromNodeId从ChannelGroup找到对应连接再发送。4.3 KV 数据库读写如何串起 Raft 日志到了这一层才算把整个系统串起来。KVDatabaseImpl对上层暴露put/get/delete接口但内部核心逻辑是「先写 Raft 日志再应用状态机」public CompletableFutureString put(String key, String value) { // 构造一条日志表示对 key 执行 PUT 操作 LogEntry entry new LogEntry(Operation.PUT, key, value); // 将日志提交给 Raft 节点进入复制流程 CompletableFutureString future new CompletableFuture(); raftNode.propose(entry, future); return future; }这里propose是一个异步方法。只有当选票多数派复制成功且日志进入 committed 状态future.complete才会被调用客户端才能收到成功响应。如果节点不是 Leaderpropose需要把请求转发给当前 Leader或者直接返回错误让客户端重试。这个项目选择的是在响应里带上leaderId由客户端自己重连。实际测试中一个简单的PUT请求从客户端发出到多数派落盘延迟大约在几毫秒到十几毫秒之间瓶颈主要在文件日志的fsync。如果你想要更高吞吐可以把fsync改成批量刷盘但这会降低故障恢复的可靠性需要自己权衡。5. 本地模拟分布式集群验证选举、故障转移与数据一致性5.1 用 IDEA 多开进程模拟三节点这个项目是 IDEA 模块kraft.iml可以直接导入 IntelliJ IDEA。模拟分布式集群不需要三台服务器一台机器多开三个进程即可。在 IDEA 的 Run Configuration 里勾选Allow multiple instances分别给三个实例传不同的启动参数# 节点 1 mainClass: com.kraft.KVStoreServer programArgs: --node.id1 --server.port8081 --data.dir./data/node1 # 节点 2 mainClass: com.kraft.KVStoreServer programArgs: --node.id2 --server.port8082 --data.dir./data/node2 # 节点 3 mainClass: com.kraft.KVStoreServer programArgs: --node.id3 --server.port8083 --data.dir./data/node3三个节点初始通过peers配置文件互相发现。启动顺序无所谓因为 Raft 会等你等不到响应后自动重试。启动完成后再开一个终端直接通过 KV 客户端向8081端口写入数据。注意IDEA 多开进程时每个进程的工作目录默认是同一个必须用--data.dir区分日志文件否则三个节点会互相覆盖索引文件。5.2 用 KVDatabaseImplTest 验证强一致性项目里自带的KVDatabaseImplTest可以改造成一致性验证工具。核心思路是向任意节点写入一个值然后从所有节点读取确认值相同接着杀掉 Leader 节点继续写入确认另一个节点能接管。Test public void testRaftConsistency() throws Exception { // 连接三个节点 KVStoreClient client1 new KVStoreClient(127.0.0.1:8081); KVStoreClient client2 new KVStoreClient(127.0.0.1:8082); KVStoreClient client3 new KVStoreClient(127.0.0.1:8083); // 向 8081 写入数据 client1.put(order:1001, paid); // 等待日志复制完成 Thread.sleep(300); // 三个节点应当返回一致结果 assertEquals(paid, client2.get(order:1001)); assertEquals(paid, client3.get(order:1001)); }这个测试跑通之后再试故障转移手动杀掉当前 Leader 进程观察剩下两个节点中是否会有一个在选举超时后变成新 Leader。注意选举超时最短是 150ms所以在杀掉进程后大概需要 150~300ms 才能选出新 Leader客户端在此时段内会收到NOT_LEADER错误。重试逻辑需要自己实现否则用户请求直接失败。5.3 本地模拟时最容易踩的五个坑第一日志目录不隔离。三节点共用data目录是最常见的问题启动后索引文件互相覆盖表现为日志复制偶尔成功偶尔失败。第二时钟漂移。Raft 假设机器时钟大致同步本地模拟时一台机器上的三个进程时钟肯定一致但如果你在多个物理机部署ntp没配好会引发选举超时误判。第三选举超时过短。同一台机器上三个进程竞争 CPUGIL 或 GC 暂停会导致心跳延迟如果超时下限设成 100ms很容易发生无主分裂。第四fsync频率。FileLogStore如果每条日志都刷盘三节点本机模拟时吞吐还能接受跨网络跑时吞吐会断崖式下跌。建议做成可配置的flushInterval数据安全性要求不高时每 10ms 刷一次。第五快照没实现。项目目前没有日志压缩和快照生成长时间运行后日志文件会无限增长。我自己的做法是加一个compact()方法定期把已提交日志里的最新 KV 状态写入快照文件然后把旧日志截断。这一步不属于 Raft 必选功能但生产环境一定需要。本文还有配套的精品资源点击获取
RELATED

相关推荐

混凝土企业ERP选型避坑指南:从业务全貌到落地细节

混凝土企业ERP选型避坑指南:从业务全貌到落地细节

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

📅 2026/9/14 13:31:43
Quick Reference 中的 Adobe Premiere Pro 键盘快捷键速查:按菜单与面板组织的完整清单

Quick Reference 中的 Adobe Premiere Pro 键盘快捷键速查:按菜单与面板组织的完整清单

Quick Reference 中的 Adobe Premiere Pro 键盘快捷键速查:按菜单与面板组织的完整清单 【免费下载链接】reference 面向开发者的技术速查清单(Cheat Sheets)集合,整理常见技术、工具与开发流程,帮助快速查阅关键信息&…

📅 2026/9/14 13:31:43
Java EE项目源码解析:Maven工程与Servlet三层架构实践

Java EE项目源码解析:Maven工程与Servlet三层架构实践

简介:面向Java EE开发者的Qimo项目设计源码,涵盖完整的企业级Web应用工程,适合正在学习Java EE、准备课程设计或希望了解项目整体架构的开发者。压缩包共185个文件,大小11.55MB,主要包含53个Java源文件、73个JPG图片、…

📅 2026/9/14 13:26:42
MORE NEWS

更多资讯

📰

Quarkdown 粗体(Strong)解析全解:从 strong.md 测试夹具到 Strong AST 节点的完整管线

Quarkdown 粗体(Strong)解析全解:从 strong.md 测试夹具到 Strong AST 节点的完整管线 【免费下载链接】quarkdown 🪐 Markdown with superpowers: from ideas to papers, presentations, websites, books, and knowledge bases. …

📰

JDK 怎么生成 compile-commands 编译数据库供 clangd 等索引器使用

JDK 怎么生成 compile-commands 编译数据库供 clangd 等索引器使用 【免费下载链接】jdk JDK main-line development https://openjdk.org/projects/jdk 项目地址: https://gitcode.com/GitHub_Trending/jd/jdk 如果你要阅读或修改 JDK 仓库里的 C/C 原生代码&#xff0…

📰

PDF书签编辑、PDF合并与解除限制:免费工具箱PDF补丁丁新手指南

PDF书签编辑、PDF合并与解除限制:免费工具箱PDF补丁丁新手指南 【免费下载链接】PDFPatcher PDF补丁丁——PDF工具箱,可以编辑书签、剪裁旋转页面、解除限制、提取或合并文档,探查文档结构,提取图片、转成图片等等 项目地址: ht…

📰

AI驱动游戏出海:专属语言引擎与买量策略实战指南

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

📰

Flask+CodeMirror+subprocess网页版Python编辑器

简介:这是一份基于Flask与CodeMirror构建的网页版Python编辑器项目源码,源自程序设计课程大作业,适合需要完成在线代码编辑、远程实验或课程设计展示的开发者参考。后端由Python Flask提供路由、登录认证与文件管理,前端通过HTML、…

📰

AI时代开发者转型指南:从代码工人到AI架构师

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

TODAY

今日更新

THIS WEEK

本周精选

THIS MONTH

本月热门

读完文章,想聊聊您的网站?

告诉我们您的行业与需求,资深顾问一对一梳理方案与报价,全程免费。

📞 💬