基于Raft分布式Kv存储:Clerk

📅 2026/7/31 16:41:29
基于Raft分布式Kv存储:Clerk
一、Clerk保存了哪些信息Clerk主要有四个成员std::vectorstd::shared_ptrraftServerRpcUtil m_servers; std::string m_clientId; int m_requestId; int m_recentLeaderId;1.m_serversstd::vectorstd::shared_ptrraftServerRpcUtil m_servers;保存所有 KVServer 的 RPC 客户端对象。例如集群有三个节点node0: 127.0.0.1:8000 node1: 127.0.0.1:8001 node2: 127.0.0.1:8002那么m_servers[0] - 访问 node0 的 RPC 对象 m_servers[1] - 访问 node1 的 RPC 对象 m_servers[2] - 访问 node2 的 RPC 对象注意这些对象不代表对应节点一定是 Leader它们只是“访问服务器的客户端代理”。2.m_clientId每个Clerk创建时都会生成一个客户端 IDClerk::Clerk() : m_clientId(Uuid()), m_requestId(0), m_recentLeaderId(0) {}作用是区分不同客户端。例如客户端 AclientId abc123 客户端 BclientId xyz789服务端可以根据ClientId RequestId判断一个请求是否已经执行过。3.m_requestId每个客户端请求都有一个递增的请求编号m_requestId; auto requestId m_requestId;例如Put(a, 1) - RequestId 1 Append(a, 2) - RequestId 2 Get(a) - RequestId 3这个编号非常重要因为 RPC 可能出现这种情况客户端发送 Append ↓ 服务器已经执行成功 ↓ 响应在网络中丢失 ↓ 客户端以为失败再次发送 Append如果第二次请求使用新的编号服务器可能会再次执行Append造成数据重复。所以一次逻辑请求在重试过程中必须始终使用同一个RequestId。4.m_recentLeaderId记录最近一次成功处理请求的服务器编号。int m_recentLeaderId;例如上一次 node2 成功m_recentLeaderId 2下一次请求就优先访问 node2。它只是一个性能优化并不保证 node2 现在仍然是 Leader。如果 Leader 发生变化node2 会返回ErrWrongLeader然后 Clerk 再尝试其他节点。二、测试客户端main截图上方的main主要是一个测试程序大致流程是int main() { Clerk client; client.Init(test.conf); client.Put(...); client.Append(...); std::string value client.Get(...); }它做了三件事第一步创建 ClerkClerk client;此时会调用构造函数m_clientId Uuid(); m_requestId 0; m_recentLeaderId 0;客户端拥有了自己的身份但还不知道服务器地址。第二步初始化服务器连接client.Init(test.conf);Init()会读取配置文件中的所有节点node0ip node0port node1ip node1port node2ip node2port然后为每个节点创建一个 RPC 客户端代理。第三步调用 KV 操作client.Put(...) client.Append(...) client.Get(...)这些函数看起来像本地函数调用但实际上内部都会经过 RPC 网络通信。例如client.Put(key, value);实际上会走Clerk::Put() ↓ Clerk::PutAppend() ↓ RPC 调用 KvServer::PutAppend()三、Clerk::Init()核心代码是MprpcConfig config; config.LoadConfigFile(configFileName.c_str());这一步加载配置文件。然后循环读取节点地址for (int i 0; i INT_MAX - 1; i) { std::string node node std::to_string(i); std::string nodeIp config.Load(node ip); std::string nodePortStr config.Load(node port); if (nodeIp.empty()) { break; } ipPortVt.emplace_back( nodeIp, atoi(nodePortStr.c_str()) ); }假设配置文件是node0ip127.0.0.1 node0port8000 node1ip127.0.0.1 node1port8001 node2ip127.0.0.1 node2port8002读取之后得到ipPortVt { {127.0.0.1, 8000}, {127.0.0.1, 8001}, {127.0.0.1, 8002} }然后为每个节点创建auto* rpc new raftServerRpcUtil(ip, port); m_servers.push_back( std::shared_ptrraftServerRpcUtil(rpc) );raftServerRpcUtil的构造函数中又创建 protobuf Stubstub new raftKVRpcProctoc::kvServerRpc_Stub( new MprpcChannel(ip, port, false) );这里有三层对象Clerk ↓ raftServerRpcUtil ↓ kvServerRpc_Stub ↓ MprpcChannel其中Clerk负责业务层重试raftServerRpcUtil封装 RPC 调用kvServerRpc_Stubprotobuf 自动生成的客户端代理MprpcChannel真正负责序列化和 TCP 通信。false表示允许延迟连接。创建MprpcChannel时可以先不连接真正调用 RPC 时再连接服务器。四、PutAppend()1.Put()和Append()只是包装函数void Clerk::Put(std::string key, std::string value) { PutAppend(key, value, Put); } void Clerk::Append(std::string key, std::string value) { PutAppend(key, value, Append); }它们最后都会进入同一个函数PutAppend(key, value, op);区别只是Put - 覆盖原值 Append - 在原值后追加这样可以减少重复代码。2. 给一次逻辑请求分配 RequestIdm_requestId; auto requestId m_requestId;这里的requestId是本次逻辑操作的编号。注意它在while循环外面只增加一次。例如第一次发送RequestId 10 第二次重试RequestId 10 第三次重试RequestId 10不能写成while (true) { m_requestId; }否则每次重试都会变成新请求重复检测就失效了。3. 选择第一次访问的服务器auto server m_recentLeaderId;如果之前 node1 成功过那么server 1这次优先访问 node1。如果是第一次运行m_recentLeaderId 0;所以第一次默认访问 node0。4. 构造 protobuf 请求raftKVRpcProctoc::PutAppendArgs args; args.set_key(key); args.set_value(value); args.set_op(op); args.set_clientid(m_clientId); args.set_requestid(requestId);最终请求里面包含key : 要操作的键 value : 要写入或追加的值 op : Put 或 Append clientId : 当前客户端 ID requestId : 当前请求编号对应的 protobuf 定义在message PutAppendArgs { bytes Key 1; bytes Value 2; bytes Op 3; bytes ClientId 4; int32 RequestId 5; }5. 发起 RPC 调用raftKVRpcProctoc::PutAppendReply reply; bool ok m_servers[server]-PutAppend( args, reply );这句代码表面上只是调用一个普通 C 函数但内部调用链是mermaid flowchart TD A[Clerk::Put] -- B[Clerk::PutAppend] B -- C[构造 PutAppendArgs] C -- D[raftServerRpcUtil::PutAppend] D -- E[kvServerRpc_Stub::PutAppend] E -- F[MprpcChannel::CallMethod] F -- G[序列化请求] G -- H[TCP 发送到 KvServer] H -- I[RpcProvider 分发服务和方法] I -- J[KvServer::PutAppend RPC入口] J -- K[KvServer::PutAppend 业务函数] K -- L[Raft::Start] L -- M[Raft复制并提交日志] M -- N[KV状态机执行] N -- O[生成 PutAppendReply] O -- P[TCP返回响应] P -- Q[Clerk反序列化并处理结果] 在客户端封装中实际代码是MprpcController controller; stub-PutAppend(controller, args, reply, nullptr); return !controller.Failed();这里的ok只表示 RPC 通信是否成功不代表业务一定成功。五、服务器端收到请求后做什么服务器端的入口是 protobuf 规定的 RPC 函数void KvServer::PutAppend( google::protobuf::RpcController* controller, const PutAppendArgs* request, PutAppendReply* response, google::protobuf::Closure* done )它会调用真正的业务函数KvServer::PutAppend(request, response); done-Run();1. 先把 RPC 参数转换成 Raft 命令Op op; op.Operation args-op(); op.Key args-key(); op.Value args-value(); op.ClientId args-clientid(); op.RequestId args-requestid();这里的Op是项目内部使用的 Raft 日志命令。也就是说客户端的 RPC 请求不会直接修改本地 KV 数据而是先变成Raft 日志中的一条命令2. 调用Raft::Start()m_raftNode-Start( op, raftIndex, _, isleader );Start()的作用是把这条命令交给 Raft。如果当前节点不是 Leaderif (!isleader) { reply-set_err(ErrWrongLeader); return; }于是响应会返回给 Clerk当前节点不是 Leader然后 Clerk 换下一个服务器继续尝试。3. 如果当前节点是 LeaderLeader 会把命令加入自己的日志并通过 Raft 的AppendEntriesRPC 复制给其他节点Leader 本地追加日志 ↓ 发送 AppendEntries ↓ Follower 接收日志 ↓ 多数节点确认 ↓ 日志提交 ↓ 各节点 ApplyMsg ↓ KV 状态机执行命令所以PutAppend()的最终一致性不是由 Clerk 完成的而是由 Raft 完成的。Clerk 只负责把请求送到某个服务器并在失败时换节点重试。六、waitApplyCh的作用Leader 接收到客户端请求后不能在Raft::Start()返回时立即告诉客户端成功。因为Start() 返回只说明命令已经提交到 Leader 的日志中不一定已经复制到多数节点也不一定已经真正执行。所以服务端会根据raftIndex建立等待通道waitApplyCh[raftIndex]然后等待 Raft 提交并应用这条日志。chForRaftIndex-timeOutPop( CONSENSUS_TIMEOUT, raftCommitOp );Raft 应用线程收到命令后会执行GetCommandFromRaft(message)然后把执行结果通知给对应的waitApplyCh。于是请求处理过程是RPC线程 Start() 创建 waitApplyCh 等待 Raft 应用结果 Raft应用线程 收到 ApplyMsg 执行 Put/Append 通知 waitApplyCh RPC线程 被唤醒 检查是不是自己的请求 设置 reply 返回客户端七、为什么要判断ClientId RequestIdKVServer 中维护了std::unordered_mapstd::string, int m_lastRequestId;含义是每个客户端最近一次已经执行的 RequestId判断函数是return RequestId m_lastRequestId[ClientId];如果发现请求已经执行过就不再重复执行。例如客户端发送ClientId A RequestId 5 Operation Append Key x Value abc假设服务器已经执行成功但是响应丢失Clerk 再次发送相同请求ClientId A RequestId 5 Operation Append Key x Value abc服务端发现m_lastRequestId[A] 5于是判断这是重复请求不再执行第二次。否则结果就可能从x helloabc错误地变成x helloabcabc这就是RequestId对Append操作尤其重要的原因。八、Clerk 如何处理返回结果核心判断是if (!ok || reply.err() ErrWrongLeader) { server (server 1) % m_servers.size(); continue; }这里包含两种失败。情况一ok false表示 RPC 通信失败例如TCP 连接失败发送失败接收失败请求序列化失败响应反序列化失败。此时客户端不知道服务器有没有执行成功因此仍然使用相同的RequestId重试。情况二reply.err() ErrWrongLeader表示网络通信成功服务器也返回了响应但是业务结果是当前节点不是 Leader此时 Clerk 继续访问下一个节点。例如第一次访问 node0ErrWrongLeader 第二次访问 node1ErrWrongLeader 第三次访问 node2OK成功后m_recentLeaderId server; return;下一次请求就优先访问 node2。情况三reply.err() OK表示RPC 通信成功 当前节点正确处理了请求 Raft 已经提交并应用了命令于是PutAppend()返回用户程序认为写入成功。要特别区分ok true只表示 RPC 层通信成功。reply.err() OK才表示 KV 业务层请求成功。九、Get()和PutAppend()的逻辑基本相同它的流程是生成 RequestId ↓ 构造 GetArgs ↓ 优先访问最近 Leader ↓ 调用 RPC ↓ RPC失败或 ErrWrongLeader ↓ 换下一个节点重试 ↓ ErrNoKey 返回空字符串 ↓ OK 返回 value客户端侧std::string value client.Get(name);服务器返回OK - 返回实际 value ErrNoKey - 返回 ErrWrongLeader - Clerk 换节点重试这个项目中的Get也会经过 Raftm_raftNode-Start(op, raftIndex, _, isLeader);这样可以保证读操作也具有线性一致性而不是直接读取某个可能落后的 Follower。