在保证线性一致性的情况下如何写Kv
要在 Raft 系统中保证 KV 写操作的线性一致性,最核心的原则是:
RPC 收到
Put/Append后不能直接修改 KV,也不能在写入 Leader 本地日志后立即返回成功;必须等命令被 Raft 提交并应用到 KV 状态机后,才能返回OK。
一、完整写入流程
假设客户端发送:
Put("x", "100") ClientId = C7 RequestId = 12完整过程应该是:
客户端发送写请求 ↓ Leader调用Raft::Start() ↓ 写入Leader本地日志 ↓ 复制给Follower ↓ 多数节点复制成功 ↓ 日志变成Committed ↓ 通过applyChan交给KVServer ↓ KVServer检查请求是否重复 ↓ 真正修改KV数据库 ↓ 更新lastRequestId ↓ 通知等待中的RPC线程 ↓ RPC返回OK对这个项目来说,最直观的线性化点就是:
已提交命令在
GetCommandFromRaft()中真正应用到 KV 状态机的时刻。
二、RPC线程不能直接写KV
PutAppend()收到请求后,只负责构造命令:
Op op; op.Operation = args->op(); op.Key = args->key(); op.Value = args->value(); op.ClientId = args->clientid(); op.RequestId = args->requestid();然后提交给 Raft:
m_raftNode->Start(op, &raftIndex, &term, &isleader);这里的Start()通常只表示:
Leader接受了命令,并尝试把它追加到Raft日志它不代表:
命令已经得到多数节点确认 命令已经提交 KV数据库已经修改因此不能这样写:
m_raftNode->Start(...); reply->set_err(OK); // 错误:此时还没有提交否则 Leader 可能刚写入本地日志就宕机,命令后来被新 Leader 覆盖,但客户端已经收到成功,线性一致性就被破坏了。
三、真正修改KV的位置
项目是在收到 Raft 的ApplyMsg后修改 KV:
void KvServer::GetCommandFromRaft(ApplyMsg message) { Op op; op.parseFromString(message.Command); if (!ifRequestDuplicate(op.ClientId, op.RequestId)) { if (op.Operation == "Put") { ExecutePutOpOnKVDB(op); } if (op.Operation == "Append") { ExecuteAppendOpOnKVDB(op); } } SendMessageToWaitChan(op, message.CommandIndex); }对应代码:[kvServer.cpp (line 166)](C:/Users/LENOVO/Desktop/KVstorageBaseRaft-cpp-main/src/raftCore/kvServer.cpp:166)
这里先去重,再执行:
Raft已经提交 ↓ 检查是否重复 ↓ 修改KV ↓ 通知RPC线程这才是正确的写入路径。
四、图片中的timeOutPop()是什么
RPC线程把命令交给 Raft 后,会通过日志下标找到对应的等待队列:
chForRaftIndex->timeOutPop( CONSENSUS_TIMEOUT, &raftCommitOp );它在等待:
Raft 应用线程通知我,这个日志位置上的命令已经提交并应用了。
这里存在两种结果。
情况一:等待到了Apply消息
也就是进入图片的else:
if (raftCommitOp.ClientId == op.ClientId && raftCommitOp.RequestId == op.RequestId) { reply->set_err(OK); } else { reply->set_err(ErrWrongLeader); }为什么不能只看到相同的raftIndex就返回成功?
假设旧 Leader 在日志位置 10 写入:
index=10:客户端C7的请求12但还没提交,旧 Leader 就失去领导权。新 Leader可能用另一条命令覆盖这个位置:
index=10:客户端C9的请求20此时 RPC线程等到了index=10的应用通知,但应用的不是自己的命令。
所以必须检查:
raftCommitOp.ClientId == op.ClientId raftCommitOp.RequestId == op.RequestId只有两者都相同,才能证明:
被应用的确实是当前客户端的当前请求然后才能返回OK。
五、情况二:等待超时
图片中的代码是:
if (!chForRaftIndex->timeOutPop(...)) { if (ifRequestDuplicate(op.ClientId, op.RequestId)) { reply->set_err(OK); } else { reply->set_err(ErrWrongLeader); } }超时只表示:
RPC线程在规定时间内没有收到Apply通知不表示命令一定失败。
命令可能已经成功应用,只是:
Apply通知到达得比较晚 等待队列通知丢失 RPC线程恰好先超时 网络或线程调度发生延迟因此超时后再次检查:
ifRequestDuplicate(op.ClientId, op.RequestId)它的判断逻辑是:
return RequestId <= m_lastRequestId[ClientId];已经执行过
如果返回true:
lastRequestId[C7] >= 12说明状态机已经执行过这个请求,所以可以返回:
reply->set_err(OK);注意:这不是“把重复请求重新执行一次”,而是:
请求已经执行过,这次只补发成功响应尚未确认执行
如果返回false,只能说明:
目前没有证据证明请求已经应用不能确定它最终会不会提交。因此不能返回OK,而是返回一个可重试错误:
reply->set_err(ErrWrongLeader);这里的ErrWrongLeader不一定真的表示“节点不是 Leader”,更多是在告诉客户端:
当前执行结果不确定,请使用相同的
ClientId + RequestId重试。
六、客户端怎么重试
客户端创建一个新逻辑请求时,只增加一次RequestId:
m_requestId++; auto requestId = m_requestId; while (true) { args.set_clientid(m_clientId); args.set_requestid(requestId); // 不断尝试不同节点 }比如:
第一次发送:(C7, 12) 超时后重试:(C7, 12) 换Leader重试:(C7, 12)不能变成:
第一次发送:(C7, 12) 第一次重试:(C7, 13) 第二次重试:(C7, 14)否则服务端会把它们当成三个不同操作,导致Append重复执行。
七、KV和去重表必须一起更新
执行Put时:
void KvServer::ExecutePutOpOnKVDB(Op op) { m_mtx.lock(); m_skipList.insert_set_element(op.Key, op.Value); m_lastRequestId[op.ClientId] = op.RequestId; m_mtx.unlock(); }这里同时更新:
KV数据 lastRequestId去重信息这是必要的。不能出现:
KV已经修改 但lastRequestId没有更新否则同一个请求重试时会被再次执行。
逻辑上,它们应该是状态机的一次原子状态转换:
(KV状态, 去重状态) ↓ 同时应用一条已提交命令 ↓ (新KV状态, 新去重状态)八、这段代码保证成功写入的依据
只有下面两种情况能够返回OK:
1. 等到了Apply消息,并且ClientId、RequestId都匹配 2. 等待超时,但去重表证明这个请求已经应用过以下情况不能返回成功:
刚调用Start() 只写入了Leader本地日志 只知道自己目前是Leader 等待到相同日志下标,但不是相同请求 超时且去重表里找不到请求