当前位置: 首页 > news >正文

基于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。

http://www.cnnetsun.cn/news/3769843.html

相关文章:

  • AI驱动项目交付提速40%的关键配置,飞书管理员绝不会告诉你的6个参数
  • 3分钟解锁《鸣潮》120FPS高帧率:WaveTools终极工具箱完整指南
  • C语言关键字深度解析:从语法基础到高级应用
  • Akagi雀魂助手:你的智能麻将AI教练真的能提升游戏水平吗?
  • M4Markets:注重效率的使用者更在意的技术架构,这里做个视角盘点
  • 【单片机毕业设计推荐】基于 STM32 或 51 单片机的水质多参数监测与自动换水控制系统设计 基于 STM32 或 51 单片机的水体 pH、温度、浑浊度智能监测装置设计(021504)
  • 如何用UI-TARS Desktop实现智能GUI自动化:5个步骤让AI成为你的数字助手
  • UE5.1蓝图优化MetaHuman性能:强制LOD与头发卡片化实战
  • Notion 接入 RAG 知识库,为什么不应该在每次问答时实时读取?
  • 10分钟完成视频制作:Auto-Video-Generator让创意不再等待
  • 别再缝合 Tushare、yfinance 和 AkShare 了!用 QuantDash 统一 A股/港股/美股量化数据接入
  • AI不是替代混音师,而是淘汰不会用AI的混音师:2024行业薪酬调研显示,掌握AI多轨协同者时薪暴涨2.8倍
  • 如何高效管理你的数字记忆:WeChatMsg微信聊天记录导出与数据分析终极指南
  • Materials Studio文件操作指南:CIF、PDB、MOL格式转换与结构处理
  • jQuery语法基础与实战技巧全解析
  • TexTools-Blender UV智能选择工具:如何快速解决翻转、重叠和相同UV岛问题
  • Input Overlay终极指南:5分钟学会在直播中展示键盘鼠标操作
  • 2026年最火的10款AI思维导图软件推荐
  • 微信小程序开店找哪家公司?普通商家更该看这几个标准
  • 把Word、PDF直接变成PPT,这个功能省的不是一点点时间
  • 你们的增长是拿钱砸出来的,还是靠自己长出来的?
  • 印度云服务器市场:全球增长最快的数据中心枢纽之一
  • 鸿蒙掌上驾考宝典应用开发11:Swiper 轮播组件——从驾考题库切换看滑动容器
  • 传统成套柜厂,如何靠智能化摆脱价格战?
  • 象棋AI连线工具终极指南:3分钟实现智能对弈分析
  • 终极指南:在Windows 10上部署Android子系统的完整教程
  • RedFoxHub:一站式新媒体数据平台实战指南,告别多平台爬虫维护之苦
  • uni-app小程序大文件分片上传与断点续传实战
  • RAG 混合检索:关键词 + 语义
  • 智测云联ZhiCloud智慧云平台如何实现7x24小时无人值守监测