news 2026/7/31 16:41:28

基于Raft分布式Kv存储:Clerk

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
基于Raft分布式Kv存储:Clerk

一、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"

这就是RequestIdAppend操作尤其重要的原因。

八、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。

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/7/31 16:41:01

3分钟解锁《鸣潮》120FPS高帧率:WaveTools终极工具箱完整指南

3分钟解锁《鸣潮》120FPS高帧率&#xff1a;WaveTools终极工具箱完整指南 【免费下载链接】WaveTools &#x1f9f0;鸣潮工具箱 项目地址: https://gitcode.com/gh_mirrors/wa/WaveTools 还在为《鸣潮》默认的60FPS限制而烦恼吗&#xff1f;想要让高端显卡充分发挥性能&…

作者头像 李华
网站建设 2026/7/31 16:39:26

C语言关键字深度解析:从语法基础到高级应用

1. 从“Hello, World!”到理解编译器&#xff1a;为什么我们需要关键字&#xff1f; 如果你写过C语言的第一个程序&#xff0c;那你一定见过 int main() 和 return 0; 。你可能知道 int 是“整数”&#xff0c; return 是“返回”&#xff0c;但有没有想过&#xff0c;…

作者头像 李华
网站建设 2026/7/31 16:38:30

Akagi雀魂助手:你的智能麻将AI教练真的能提升游戏水平吗?

Akagi雀魂助手&#xff1a;你的智能麻将AI教练真的能提升游戏水平吗&#xff1f; 【免费下载链接】Akagi 支持雀魂、天鳳、麻雀一番街、天月麻將&#xff0c;能夠使用自定義的AI模型實時分析對局並給出建議&#xff0c;內建Mortal AI作為示例。 Supports Majsoul, Tenhou, Riic…

作者头像 李华
网站建设 2026/7/31 16:37:12

M4Markets:注重效率的使用者更在意的技术架构,这里做个视角盘点

对新手与注重稳健体验的外汇内容读者而言&#xff0c;“能看懂”往往比“堆概念”更重要。以M4Markets为例&#xff0c;以下重点写清解释是否通俗、规则是否易查、提示是否前置&#xff0c;以及服务是否具备连续性。外汇相关平台的价值&#xff0c;体现在长期一致性与信息呈现的…

作者头像 李华
网站建设 2026/7/31 16:36:16

【单片机毕业设计推荐】基于 STM32 或 51 单片机的水质多参数监测与自动换水控制系统设计 基于 STM32 或 51 单片机的水体 pH、温度、浑浊度智能监测装置设计(021504)

文章目录20 个相关毕业设计备选题目项目研究背景摘要总体方案核心功能基础功能核心功能辅助功能技术路线项目演示关于我们项目案例源码获取温馨提示&#xff1a;本人主页置顶文章(点我)有 CSDN 平台官方提供的学长联系方式的名片&#xff01; 温馨提示&#xff1a;本人主页置顶…

作者头像 李华