一、Clerk保存了哪些信息
Clerk主要有四个成员:
std::vector<std::shared_ptr<raftServerRpcUtil>> m_servers; std::string m_clientId; int m_requestId; int m_recentLeaderId;1.m_servers
std::vector<std::shared_ptr<raftServerRpcUtil>> 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创建时都会生成一个客户端 ID:
Clerk::Clerk() : m_clientId(Uuid()), m_requestId(0), m_recentLeaderId(0) {}作用是区分不同客户端。
例如:
客户端 A:clientId = abc123 客户端 B:clientId = 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(...); }它做了三件事:
第一步:创建 Clerk
Clerk 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()) ); }假设配置文件是:
node0ip=127.0.0.1 node0port=8000 node1ip=127.0.0.1 node1port=8001 node2ip=127.0.0.1 node2port=8002读取之后,得到:
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_ptr<raftServerRpcUtil>(rpc) );raftServerRpcUtil的构造函数中又创建 protobuf Stub:
stub = new raftKVRpcProctoc::kvServerRpc_Stub( new MprpcChannel(ip, port, false) );这里有三层对象:
Clerk ↓ raftServerRpcUtil ↓ kvServerRpc_Stub ↓ MprpcChannel其中:
Clerk:负责业务层重试;
raftServerRpcUtil:封装 RPC 调用;
kvServerRpc_Stub:protobuf 自动生成的客户端代理;
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. 给一次逻辑请求分配 RequestId
m_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。
如果当前节点不是 Leader:
if (!isleader) { reply->set_err(ErrWrongLeader); return; }于是响应会返回给 Clerk:
当前节点不是 Leader然后 Clerk 换下一个服务器继续尝试。
3. 如果当前节点是 Leader
Leader 会把命令加入自己的日志,并通过 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 + RequestId
KVServer 中维护了:
std::unordered_map<std::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 继续访问下一个节点。
例如:
第一次访问 node0:ErrWrongLeader 第二次访问 node1:ErrWrongLeader 第三次访问 node2:OK成功后:
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也会经过 Raft:
m_raftNode->Start(op, &raftIndex, &_, &isLeader);这样可以保证读操作也具有线性一致性,而不是直接读取某个可能落后的 Follower。