基于Raft分布式Kv存储:读屏障

📅 2026/7/31 9:44:05
基于Raft分布式Kv存储:读屏障
读屏障Read Barrier可以理解为在真正读取 KV 之前先在 Raft 的全局操作顺序中确定一个安全位置B并等待本节点的状态机至少执行到B然后才能读取。它不是 C 中的一把锁也不是 CPU 的内存屏障而是分布式系统中的一个逻辑检查点。一、为什么需要读屏障假设客户端已经完成写操作Put(x, 100) - 返回 OK随后客户端发起Get(x)根据线性一致性Get必须读到x 100或者读到在它之后发生的更新但绝对不能读到Put之前的旧值。然而节点的 Raft 状态与 KV 状态之间可能存在延迟Raft日志已经提交到 index10 KV状态机只应用到 index8即commitIndex 10 lastApplied 8如果直接读取本地 KV就会漏掉日志 9、10 中的写操作。所以必须等待lastApplied 安全读取位置这个等待过程就是读屏障的重要组成部分。二、读屏障到底“挡”住了什么假设日志是index7Put(x, 100) index8Append(x, A) index9Get(x) ← 读屏障状态机严格按顺序处理应用 index7x 100 应用 index8x 100A 到达 index9允许执行Get因此当Get对应的日志 9 到达 Apply 层时可以确定排在日志 9 前面的所有已提交操作都已经应用到 KV。读屏障像一条分界线屏障之前必须完成 -------------------- Read Barrier 屏障之后可以尚未执行此时读取x至少能够看到x 100A三、读屏障需要解决两个问题一个完整的读屏障通常需要同时解决两件事。1. 确认当前节点有资格处理读取节点本地可能仍然认为自己是 Leader但实际上已经因为网络分区失去多数派。例如S1旧Leader被隔离 S2新Leader S3Follower新 LeaderS2已经完成Put(x, 200)旧 LeaderS1还保存着x 100如果S1直接处理Get就会返回旧值。因此读屏障必须确认当前Leader仍然可以和多数节点通信把Get写入 Raft 日志可以完成这种确认因为该日志只有被多数节点接受后才能提交。2. 确认本地KV已经追上安全位置即使节点确实是 LeaderRaft线程和状态机线程也可能不同步commitIndex 20 lastApplied 17日志已经提交到 20但 KV只执行到 17。直接读取仍然可能得到旧值。所以还必须等待lastApplied readIndex总结起来确认Leader仍然有效 等待状态机执行到安全位置 可以进行线性一致读四、本项目如何实现读屏障本项目采用最直观的方式把Get自身作为一条 Raft 日志提交。第一步构造Get命令Op op; op.Operation Get; op.Key args-key(); op.ClientId args-clientid(); op.RequestId args-requestid();第二步写入Raftm_raftNode-Start( op, raftIndex, term, isLeader );假设 Raft 分配raftIndex 20这里日志 20 就是这个Get的读屏障。但Start()返回不代表日志已经提交只表示当前节点尝试将它添加到本地日志。第三步等待Apply通知RPC线程调用chForRaftIndex-timeOutPop( CONSENSUS_TIMEOUT, raftCommitOp );等待过程是Get进入日志20 ↓ 复制到多数节点 ↓ 日志20被提交 ↓ Apply线程依次处理日志 ↓ 状态机执行完日志119 ↓ Apply线程到达日志20 ↓ 通知Get RPC线程只有 Apply线程到达日志 20才能证明屏障前面的写操作已经完成。五、这个项目的Apply过程Raft提交命令后ReadRaftApplyCommandLoop()会收到ApplyMsgauto message applyChan-Pop(); if (message.CommandValid) { GetCommandFromRaft(message); }然后进入void KvServer::GetCommandFromRaft(ApplyMsg message)对于Put/Append它会修改 KVif (!ifRequestDuplicate(op.ClientId, op.RequestId)) { if (op.Operation Put) { ExecutePutOpOnKVDB(op); } if (op.Operation Append) { ExecuteAppendOpOnKVDB(op); } }对于Get这里不修改业务 KV而是通知等待这个日志下标的 RPC线程SendMessageToWaitChan(op, message.CommandIndex);因为 Apply循环按照日志顺序运行当它到达Get时前面的Put/Append已经处理完成。六、为什么还要核对请求身份RPC线程收到通知后检查if (raftCommitOp.ClientId op.ClientId raftCommitOp.RequestId op.RequestId) { // 执行读取 }不能只判断日志下标相同。例如旧 Leader将当前 Get 放到了日志 20index20(C7, RequestId12, Get)但还没有提交旧 Leader就失去了领导权。新 Leader可能覆盖日志 20index20(C9, RequestId6, Put)RPC线程最终等到了日志 20 的 Apply通知但被应用的并不是自己的 Get。因此需要检查日志下标相同 ClientId相同 RequestId相同确认提交的确实是自己的请求后读屏障才算建立成功。七、什么时候真正读取KV屏障建立后RPC线程调用ExecuteGetOpOnKVDB(op, value, exist);内部在锁的保护下读取m_mtx.lock(); if (m_skipList.search_element(op.Key, *value)) { *exist true; } m_mtx.unlock();整体顺序是Get日志提交 ↓ Apply线程到达Get日志 ↓ 说明屏障之前的写已应用 ↓ RPC线程收到匹配通知 ↓ 加锁读取KV ↓ 返回结果八、线性化点在哪里可以区分两个重要时刻Get日志提交/应用这个时刻建立了安全读屏障屏障之前的写一定已经提交并应用锁内读取KV这个时刻确定了实际返回值可以将它看作当前实现中更直接的读操作线性化点。例如Get的读屏障建立 ↓ 另一个Put并发执行 ↓ Get真正读取KV如果 Get 读到了新 Put 的值也不违反线性一致性因为 Get 和 Put 在时间上重叠可以把 Get 的线性化点放在 Put 之后。九、ReadIndex也是一种读屏障每次把Get写入日志虽然简单但成本较高创建日志 复制给Follower 等待多数派确认 日志持续增长生产系统常用ReadIndex1. Leader向多数节点发送心跳 2. 确认自己仍然是当前任期的Leader 3. 记录当前安全commitIndex为readIndex 4. 等待本地lastApplied readIndex 5. 读取本地KV假设readIndex 100 lastApplied 96暂时不能读需要等待lastApplied 100然后才能读取。ReadIndex没有新增Get日志但同样建立了日志100之前的命令已经提交 本地状态机已经执行到日志100所以它也是读屏障。十、读屏障不等于加互斥锁二者解决的问题完全不同。互斥锁解决本机多个线程不能同时不安全地访问KV读屏障解决这个节点读取的数据在整个Raft集群中是否足够新 当前节点是否仍然有资格提供线性一致读只有锁std::lock_guardstd::mutex lock(m_mtx); return kv[key];只能保证没有本机数据竞争不能防止读到Follower旧数据 读到旧Leader旧数据 读到尚未应用完的状态因此通常需要Raft读屏障 本地互斥锁