公司动态

基于Raft分布式Kv存储:RPC

📅 2026/7/31 19:48:17
基于Raft分布式Kv存储:RPC
在这个项目里RPC 可以理解为整个分布式 KV 系统的“远程通信层”。它负责让不同进程、不同端口上的对象能够互相调用方法。RPC 本身不实现 Raft 算法也不存储 KV 数据它负责把请求送过去、调用正确的方法再把结果送回来。一、项目中有两类 RPC客户端 Clerk │ Get / Put / Append ▼ KVServer │ RequestVote / AppendEntries / InstallSnapshot ▼ 其他 Raft 节点对应两份.proto通信双方RPC 服务方法Clerk → KVServerkvServerRpcGet、PutAppendRaft → RaftraftRpcRequestVote、AppendEntries、InstallSnapshot1. Clerk 与 KVServer 之间的 RPC接口定义在 [kvServerRPC.proto (line 43)](C:/Users/LENOVO/Desktop/KVstorageBaseRaft-cpp-main/src/raftRpcPro/kvServerRPC.proto:43)service kvServerRpc { rpc PutAppend(PutAppendArgs) returns(PutAppendReply); rpc Get(GetArgs) returns(GetReply); }它负责把用户的 KV 操作发送给远程服务器。GetRPC客户端请求message GetArgs { bytes Key 1; bytes ClientId 2; int32 RequestId 3; }作用分别是Key要查询的键。ClientId客户端唯一标识。RequestId该客户端的请求序号用于重试和去重。bool ok m_servers[server]-Get(args, reply);这个调用看起来像普通 C 函数但m_servers[server]代表远程服务器实际会经过网络发送。PutAppendRPC客户端可以发起两种修改Put(key, value); Append(key, value);最终都被转换成PutAppend(key, value, Put); PutAppend(key, value, Append);请求中包含message PutAppendArgs { bytes Key 1; bytes Value 2; bytes Op 3; bytes ClientId 4; int32 RequestId 5; }ClientId RequestId很重要。因为 RPC 可能超时客户端可能重发同一个请求。如果没有请求编号一个Append请求重试两次就可能真的追加两次。RPC 负责传递编号KVServer 负责判断是否重复执行。2. RPC 帮助客户端寻找 LeaderClerk 一开始不一定知道哪个服务器是 Leader。在 [clerk.cpp (line 21)](C:/Users/LENOVO/Desktop/KVstorageBaseRaft-cpp-main/src/raftClerk/clerk.cpp:21) 中它会循环请求服务器bool ok m_servers[server]-Get(args, reply); if (!ok || reply.err() ErrWrongLeader) { server (server 1) % m_servers.size(); continue; }这里有两种失败!okRPC 通信失败例如连接失败、服务器宕机。ErrWrongLeader网络通信成功但请求发给了 Follower。这两种失败不能混为一谈。ok true只表示请求成功发出并成功收到了可解析的响应不表示业务操作一定成功。Clerk 会继续尝试其他服务器找到 Leader 后缓存其编号到m_recentLeaderId下次优先联系它。3. KV 请求最终会进入 RaftKVServer 收到PutAppendRPC 后不会直接修改本地数据库。Op op; op.Operation args-op(); op.Key args-key(); op.Value args-value(); op.ClientId args-clientid(); op.RequestId args-requestid(); m_raftNode-Start(op, raftIndex, term, isLeader);执行过程是Clerk 发起 Put RPC ↓ Leader KVServer 收到请求 ↓ 把请求转换成 Op ↓ Raft::Start() 写入 Leader 日志 ↓ 通过 AppendEntries RPC 复制到其他节点 ↓ 多数节点复制成功 ↓ 日志提交并发送到 applyCh ↓ KVServer 执行 Put ↓ 通过 RPC 向 Clerk 返回成功所以客户端 RPC 和 Raft RPC 是串联起来的。读取也经过共识顺序有助于提供线性一致性读取。二、Raft 节点之间的 RPCservice raftRpc { rpc AppendEntries(AppendEntriesArgs) returns(AppendEntriesReply); rpc InstallSnapshot(InstallSnapshotRequest) returns(InstallSnapshotResponse); rpc RequestVote(RequestVoteArgs) returns(RequestVoteReply); }这三个 RPC 基本覆盖了 Raft 节点间的全部通信。1.RequestVote选举 Leader当 Follower 长时间没有收到 Leader 心跳就会Follower ↓ 选举超时 Candidate ↓ currentTerm 向所有其他节点发送 RequestVote RPC请求携带Term CandidateId LastLogIndex LastLogTerm接收节点检查Candidate 的任期是否足够新。当前任期是否已经投票。Candidate 的日志是否至少和自己一样新。响应包含Term VoteGranted候选人获得超过半数的票后成为 Leader。没有RequestVoteRPC各节点就无法跨进程完成投票和 Leader 选举。2.AppendEntries心跳和日志复制这是项目里使用最频繁的 Raft RPC作用有两个。作为心跳Leader 周期性发送没有日志的AppendEntriesEntries 空它告诉 FollowerLeader 还活着不要发起新选举。Follower 收到合法心跳后会重置选举定时器。作为日志复制客户端提交命令后Leader 把日志通过AppendEntries发给 FollowerTerm LeaderId PrevLogIndex PrevLogTerm Entries LeaderCommitFollower 检查上一条日志是否匹配本地 PrevLogIndex 位置的 term 请求中的 PrevLogTerm匹配则追加日志不匹配则返回失败。Leader 根据返回结果更新m_nextIndex[server] m_matchIndex[server]如果复制失败Leader 会调整nextIndex再次发送更早的日志。当一条日志被多数节点确认后Leader 推进commitIndex这条命令才真正提交。3.InstallSnapshot同步快照如果某个 Follower 落后太多Leader 需要的旧日志可能已经被压缩成快照。这时不能再用普通AppendEntries补日志而要发送InstallSnapshot RPC请求携带LeaderId Term LastSnapShotIncludeIndex LastSnapShotIncludeTerm Data其中Data是 KV 数据、客户端去重信息等序列化后的快照。Follower 收到后会安装快照数据。删除快照覆盖的旧日志。更新快照索引和任期。更新commitIndex、lastApplied。把快照交给上层 KVServer 恢复状态。它解决的是“落后节点快速追赶”的问题。三、底层 RPC 框架具体做了什么业务代码只需要写stub-AppendEntries(controller, args, response, nullptr);底层实际完成了很多工作。1. 定义统一消息格式message RpcHeader { bytes service_name 1; bytes method_name 2; uint32 args_size 3; }一次请求大致是Header长度 RpcHeader protobuf序列化后的请求参数例如service_name raftRpc method_name AppendEntries args_size 120 args AppendEntriesArgs的二进制数据2. 客户端序列化并发送从方法描述符获取服务名和方法名。序列化请求参数。构造 RPC 请求头。连接目标 IP 和端口。通过 TCP 发送数据。阻塞等待响应。将响应反序列化到reply。通过RpcController记录通信错误。3. 服务端注册和查找服务KVServer 启动时同时注册两个服务provider.NotifyService(this); provider.NotifyService(this-m_raftNode.get());也就是KvServer对象 → Get、PutAppend Raft对象 → RequestVote、AppendEntries、InstallSnapshotRpcProvider保存服务名、方法名、服务对象之间的映射。4. 服务端反序列化和方法分发解析 RpcHeader ↓ 得到 service_name ↓ 得到 method_name ↓ 找到注册的 Service 对象 ↓ 创建对应的 request 和 response ↓ 反序列化 request ↓ service-CallMethod(...)protobuf 生成的CallMethod再把请求分发到正确的重写方法。5. 把响应发回来业务函数填写response后执行done-Run();它会调用RpcProvider::SendRpcResponse(...)将响应序列化并通过 TCP 发送给调用者。四、RPC 不负责什么RPC 只是通信机制它不负责决定谁是 Leader。判断是否应该投票。检查日志是否匹配。决定日志何时提交。修改 KV 数据。处理重复客户端请求。制作或应用快照。这些分别由 Raft、KVServer 和 Clerk 的业务逻辑负责。最准确的划分是RPC把某个对象的方法调用跨网络送到另一个进程 Raft保证多个节点对命令顺序达成一致 KVServer执行命令并维护键值数据 Clerk发起请求、寻找 Leader、失败重试