【adviskv】#04:PUT 链路 在 Storage 内部的 Replica 的设计与实现(上)

日期: 2026.8.8

写在开头的开头

AdvisKV 是我用 C++17 从零写的一个分布式 KV 存储系统原型,包含 CatalogTopoStorageSDK 四个主要模块,覆盖建库建表、路由查询、KV 读写、副本数扩缩容、坏副本替换、RaftWAL/snapshot 恢复,跑了 gtest + E2E 测试和 benchmark

欢迎来给个 star~

项目地址:advisedy/adviskv: C++17 分布式 KV 存储原型,包含 Raft、Meta/SDM、分片路由与恢复链路 | A runnable C++17 distributed KV prototype

写在开头

已经是秋招了,最近基本上无心准备秋招,心浮躁的很,特来决定继续更新 AdvisKV

之前的内容讲过 PUT 的一生,考虑以后打算要讲解清楚更详细的内容,所以应该要接下来我们不得不继续搞定这个 Storage 内部的这个硬骨头,同时也是我重构了好多次,思考过很久之后的模块,Replica

说实在的,这个肯定很难讲解清楚,因为这个不是一条清晰的链路就可以说明白的。里面涉及的模块相对来讲会稍微复杂一点。

我们之前在第二篇的时候讲过关于 PUT 的链路,但是当时是走的整体链路,我们并没有深入讲解 Replica 的内部细节,关于 PUT 是怎么样处理的。

今天就继续按照之前的逻辑去继续讲解。

我的想法是还是接着之前的 PUT 链路去走,由于里面的模块太多了,所以每走到一个模块之后,就尽量的去详细讲解一下模块的内容,这样应该思路更通顺一点。

先看一眼这个整体的走的流程,大致如下:

StorageServiceImpl-> ReplicaManager-> Replica::put-> ReplicaLoop-> RaftCore-> PersistEngine-> ReplicaMessageDispatcher-> ReplicaLoop 处理 response-> commit-> ReplicaApplyTask-> ReplicaApplier-> KvStateMachine-> MapEngine

然后由于 Replica 的设计本身就有点繁琐,Raft 的设计同样也是这样,所以分成两期,这一期主要讲 Replica 的内容,下一期去主要讲关于 Raft 的内容,本期对于 Raft 的内容应该就是大致略过的那种。

另外,GETsnapshot、恢复、成员变更虽然不是一次普通 PUT 的主路径,但后面我也会进行讲解。

一、PUT 的整体链路

whiteboard_exported_image

继续跟着我们之前的思路,此时 SDK 已经根据 Route 找到了某个 Leader Replica,并把 PUT RPC 发到了它所在的 Storage

完整链路可以先记成下面这样:

StorageServiceImpl::Put│▼
ReplicaManager::get_replica_by_id│▼
Replica::put│├── 参数、本地状态、recovering 检查│▼
ProposeCall│▼
ReplicaLoop::sync_submit│▼
submit_queue_ 批量收集 proposal│▼
RaftCore::propose_batch│├── 追加本地 Raft log├── 产生 RaftEffects├── 写 WAL└── 生成 AppendEntries 消息│▼
ReplicaMessageDispatcher│▼
其他 Replica 处理 AppendEntries│▼
AppendEntriesResponse 回到 Leader│▼
RaftCore 更新复制进度并推进 commit_index│▼
ReplicaLoop::resolve_pending_proposals│▼
Replica::put 返回 OK│▼
ReplicaApplyTask│▼
ReplicaApplier│▼
KvStateMachine│▼
MapEngine

另外,这里有一个前置知识:PUT 成功返回,意味着 commit 成功。

二、Storage 如何找到目标 Replica

2.1 Replica 逻辑定义

adviskv 的数据层次是:

Table└── Shard└── 一个 Raft Group├── Replica A├── Replica B└── Replica C

一个 Shard 是一个 Raft Group,而每个 Storage 节点上承载的,是这个 group 的一个 Replica

同一个 Storage 进程可能同时放着很多不同 ShardReplica,例如:

table=1, shard=0, replica=0
table=1, shard=2, replica=1
table=3, shard=1, replica=0

所以 Storage 收到 PUT 后,不能只根据 table_idshard_index 找对象,而是要根据完整的 ReplicaID

ReplicaID = table_id + shard_index + replica_seq

ReplicaManager 里维护了个索引:

std::unordered_map<ReplicaID, ReplicaPtr, ReplicaIDHash> replica_map_;

ReplicaManager 可以使用 get_replica_by_id 查询到对应的 Replica

2.2 ReplicaManager 的作用

StorageServiceImpl::Put 的主要工作是:

收到 PutRequest-> 解码 ReplicaID-> ReplicaManager::get_replica_by_id-> 找不到则返回 REPLICA_NOT_FOUND-> 找到则调用 Replica::put

但是不仅是一个查询作用,由于 ReplicaManager 的内容比较简单,所以我就直接现在这里大致讲解 ReplicaManager 的功能了:


class ReplicaManager {
public://......Status add_replica(const ReplicaInitParam& param);Status delete_replica(const ReplicaID& replica_id);Status add_member(const ReplicaID& leader_replica_id, const PeerMember& member);Status remove_member(const ReplicaID& leader_replica_id, const ReplicaID& replica_id);void start_tick();void recover();private:ReplicaInitParam fill_param_runtime(ReplicaInitParam param) const;//......ReplicaRuntimeOptions runtime_options_;ReplicaMetaPersistEngine meta_persist_;
};

其实作为 ReplicaManager 的作用,在 PUT 链路里面是查找到对应的 Replica;如果是在扩缩容的场景下,会提供给 NodeAgent 一个函数方便回调。

这里又简单提一嘴 NodeAgent,目前就简单的理解成一个定时会给 TopoRPC 的一个东西就好了,只不过在 Topo 希望扩缩容的时候,不是 Topo 打给 Storage,而是 Storage 打向 Topo,接收到期望态之后自己做处理,这个自己做处理,就是 NodeAgent 接受到了之后会去调用例如 add_memberremove_member 这种操作,ReplicaManager 就是负责提供这个回调函数的。

然后除了刚才说的两个函数,还有 start_tick。也就是说,Replica 里面会有一个 tick 操作,定时触发,而这个定时触发,我们不是在每一个 Replica 里面塞了一个定时器,而是直接放在 ReplicaManager 搞了个定时任务,去定时触发所有的 Replica,不然每一个 Replica 都会搞一个线程,消耗太大了。

还有一个就是 recover,负责拉起来所有的 Replica


void ReplicaManager::recover() {std::unique_lock locker(mutex_);{//......std::vector<std::filesystem::path> meta_paths = meta_persist_.scan_replica_meta_files();for (const auto& meta_path : meta_paths) {ReplicaMetaPayload payload;Status status = meta_persist_.load_replica_meta(meta_path, payload);//...ReplicaInitParam& param = payload.init_param;param = fill_param_runtime(param);//...auto replica = std::make_shared<Replica>();status = replica->init(param);//...}}for (const auto& [_, replica] : replica_map_) {Status status = replica->recover();//...}LOG_INFO("all replicas recovered");
}

以上内容,就是 ReplicaManager 负责的内容了。

所以现在需要再扯回来,继续进行 PUT 链路的细节部分。

三、Replica::put 进入 Replica

3.1 Replica 各个模块

一个 Replica 里面持有下面这些对象:

Replica├── StateMachine(KvStateMachine)│     └── MapEngine├── RaftCore├── PersistEngine├── ReplicaLoop├── ReplicaMessageDispatcher├── ReplicaApplier├── ReplicaApplyTask├── ReplicaReadIndexChecker└── ReplicaSnapshotCoordinator

因此,Replica 不是 MapEngine 的一层薄包装。它把一个副本运行所需要的本地状态、Raft、持久化、网络、apply、读检查和 snapshot 都组合在了一起。

3.2 Replica::put 的前置检查

当前 Replica::put 进入正式 proposal 之前,会做几层检查:

PutParam 参数检查│
获取 OperGate 操作 guard│
检查 ReplicaLocalState 是否为 RUNNING│
检查 RaftCore 是否处于 RECOVERING│
创建 ProposeCall

源码主线可以简化成:

Status Replica::put(const PutParam& param) {RETURN_IF_INVALID_PARAM(param)RETURN_IF_OPER_GUARD_ACQUIRE_FAILED(oper_gate_)RETURN_IF_INVALID_STATUS(ensure_local_state_running()){std::lock_guard lock(raft_core_mutex_);if (raft_core_->is_recovering()) {return Status::IS_RECOVERING("replica is recovering");}}ProposeCall call{ProposeParam::write(WriteOpType::PUT, param.key, param.value),K_REPLICA_PROPOSAL_COMMIT_TIMEOUT};loop_->sync_submit(&call);RETURN_IF_INVALID_STATUS(call.status)notify_apply_task();return Status::OK();
}

先说一下这个 RETURN_IF_OPER_GUARD_ACQUIRE_FAILED(oper_gate_),这个是因为我们需要删除 Replica 的时候呢,得确保他所有的操作都是执行完了的,或者说起码得让他不再继续执行操作了。

否则可能会出现,ReplicaManager 已经把 Replica 在磁盘级别删了,但是 Replica 内部可能还有一些操作,导致数据残留。例如这种,所以我们需要确保 Replica 的操作都执行了,就用了一个 OperGuard

这里关于 ReplicaRaft 的状态内容可以等到下一篇的时候细讲。

所以 Replica 主要是走了一个检查功能 + 投递到 loop 里面去,在 loop 里面再串行处理这个 Raft 相关内容。

四、ProposeCall 如何进入 ReplicaLoop

4.1 ReplicaLoop 的目的

这里得扯一下了,这个设计的缘由是因为我最开始的是时候实现 Raft 没有想这么多,直接一股脑塞到这个 Raft 里面,(因为感觉各个模块都比较耦合,想着先直接实现了再说吧),但是结果搞了超级无敌屎山,1000+ 行的源文件,过了两三天再看我直接原地失忆,所以不得不对这坨东西进行了重构,目前才相对来讲比较好受一点。

RaftCore 里面有很多相互关联的状态:

  • current term
  • Leader/Follower/Candidate
  • Raft log
  • commit_index
  • last_applied
  • membership
  • Leader 视角下每个 follower 的复制进度;
  • electionheartbeat 的计时状态。

如果 PUT 请求线程、Raft RPC 回调线程和 tick 线程都直接修改这些状态,就会产生复杂的并发交错。(当时的屎山简直了,锁也是满天飞,并发问题更是多的不敢想)

ReplicaLoop 的核心思路是:

PUT proposal
Raft RPC response
Raft tick
成员变更
snapshot 状态变化
以上内容,放入│▼ReplicaLoop│▼
一次串行执行一个 Raft 相关 step

有点类似线程池,具体来说是把 Raft 状态变化串行起来。

当前实现仍然有 raft_core_mutex_state_machine_mutex_persist_snapshot_mutex_ReplicaLoop 不是为了让整个 Replica 变成无锁结构,而是让 Raft 状态的主要修改有一个比较明确的串行入口。

4.2 ReplicaLoop 的接口

ReplicaLoop 对外提供两个入口:

void async_submit(Event event);
void sync_submit(Call call);

Event 表示异步事件,Call 表示调用者等待需要拿回结果。

PUT 使用的是 ProposeCall。它的路径大致是:

Replica::put-> loop_->sync_submit(ProposeCall)-> ReplicaLoop::handle_call(ProposeCall)-> enqueue_proposal_and_wait-> submit_queue_-> proposal 被真正处理-> 等待 TimeoutWaiter

这里的同步等待不是“把整个 Raft 流程在当前线程里执行完”,而是当前线程提交 proposal 后挂起,等后续 commit 或超时唤醒。

4.3 为什么普通 proposal 要进入 submit_queue_

ReplicaLoop 内部有两套队列:

  • SerialTaskRunner:串行执行已经准备好的任务;
  • BatchDispatchQueue:暂存外部提交的 proposal 和配置变更。

多个 PUT/DEL 在很短时间内到达时,BatchDispatchQueue 会把它们收集起来,搞了两个边界:

  • 最多 64 个 proposal
  • 或等待大约 200 微秒。

然后由 ReplicaLoop 一次调用 RaftCore::propose_batch

PUT 1 ─┐
PUT 2 ─┼── submit_queue_
PUT 3 ─┘       │▼RaftCore::propose_batch│一次追加多条日志

这就是当前代码中 proposal batch 的来源。它不是把多个客户端请求合并成一条 KV 操作,而是把多个独立的 ProposeParam 一次交给 RaftCore 处理。

4.4 EventCall

ReplicaLoopEvent 定义包括:

Event 产生原因
TickEvent RaftTickTask 定期触发
AppendResponseEvent AppendEntries RPC 成功返回
VoteResponseEvent RequestVote RPC 成功返回
SnapshotResponseEvent InstallSnapshot RPC 成功返回
AppendSendFailedEvent AppendEntries 发送失败
SnapshotSendFailedEvent InstallSnapshot 发送失败
PublishSnapshotEvent snapshot 发布状态发生变化

Event 不要求提交者等待处理结果。典型路径是:

RPC worker 收到 AppendEntriesResponse│▼
构造 AppendResponseEvent│▼
async_submit│▼
ReplicaLoop::handle_event

Call 则包含输入、输出和 Status。大致包括:

  • RequestVoteCall
  • AppendEntriesCall
  • BuildReadIndexCall
  • PrepareInstallSnapshotCall
  • CommitInstallSnapshotCall
  • AppendResponseCall
  • SnapshotResponseCall
  • ProposeCall
  • AddMemberCall
  • RemoveMemberCall

使用 Call,一般都是因为 RPC handler 需要把 result 回填给 StorageServiceImplLeader 发出的 RPC response 通常使用 Event,是因为网络线程只需要把结果通知给 ReplicaLoop

五、RaftCore 如何把 PUT 变成 RaftEffects

5.1 接口边界

PUT 进入 ReplicaLoop 后,真正执行 Raft 判断的是 RaftCore

输入:proposal / tick / RPC request / RPC responseRaftCore 计算:输出:RaftEffects

RaftCore 不直接写 WAL,也不直接调用 gRPC。它把副作用放到 RaftEffects 中:

struct RaftEffects {std::optional<RaftMeta> hard_state;std::vector<LogEntry> entries_to_append;std::optional<std::vector<LogEntry>> entries_to_rewrite;std::vector<RaftMessage> messages;
};

这个 RaftEffects,主要是交给 Replica 层面,在设计的时候并不打算交给 Raft 层面去做,否则会有持久化,RPC 层面的事情耦合在一起,太复杂了。

对一次正常 PUT 来说,可以把结果理解成:

RaftEffects├── entries_to_append│       └── 一条 PUT LogEntry└── messages└── 发给其他副本的 AppendEntries

5.2 ReplicaLoop::run_step 的职责

ReplicaLoop 通过 run_step 包住一次 RaftCore 操作:

拿住 raft_core_mutex_│
调用 RaftCore::propose / propose_batch / handle_xxx│
得到 RaftEffects│
PersistEngine 持久化 effects│
ReplicaMessageDispatcher 发送 effects.messages

源码主线可以简化成:

Status ReplicaLoop::run_step(RaftStepFunc&& step) {RaftEffects effects;{std::lock_guard lock(context_.raft_core_mutex);Status status = step(effects);RETURN_IF_INVALID_STATUS(persist_raft_effects(effects))RETURN_IF_INVALID_STATUS(status)}return context_.message_dispatcher.async_send(std::move(effects.messages));
}

这里其实就是一套固定的流程搭配:

先把 RaftCore 产生的持久化副作用落盘,再发送网络消息。

这样其他副本收到 AppendEntries 时,Leader 本地已经先把对应日志写进 WAL,重启后仍然有机会恢复这条日志。

由于这个流程固定,唯一的变动就是 step 具体的内容,所以把他当做参数传进来,做一层封装处理。

六、PersistEngine 负责把本地结果落盘

PersistEngine 主要负责三类持久化内容:

内容 作用
WAL 保存 Raft 日志
raft_meta 保存 current termvoted_forRaft hard state
snapshot 保存已经 apply 的状态机和当时的 membership

正常来说,一次 PUT 产生 entries_to_append 后,ReplicaLoop 会调用 PersistEngineappend_wal_batchPersistEngine 写完 WAL 后还会 fsync

因此 PUT 的本地持久化路径是:

RaftCore::propose_batch-> RaftEffects.entries_to_append-> ReplicaLoop::persist_raft_effects-> PersistEngine::append_wal_batch-> fsync

PersistEngine 保存的是 LogEntry,不是直接保存最终的 map 内容。恢复时,系统可以先恢复 snapshot,再读取 snapshot 之后的 WAL,重新把日志 apply 到状态机。

如果 WALhard statesnapshot 的关键持久化失败,Replica 会通过 context 中的 fault_if_fail 回调,把本地状态变成 FAULTED

// 负责WAL,snapshot,raft_meta的落盘
class PersistEngine {
public:PersistEngine(const std::string& data_dir, const ReplicaID& replica_id);~PersistEngine();Status init();Status close();Status append_wal(const LogEntry& entry);Status append_wal_batch(const std::vector<LogEntry>& entries);Status read_wal_batch(std::vector<LogEntry>& entries);Status rewrite_wal(const std::vector<LogEntry>& entries);// 和下面的truncate_wal_to_offset区分一下,这个是代码业务层这边调用的,用来截取到内存里的wal// 而truncate_wal_to_offset是用来截取磁盘里的wal的Status truncate_wal(const LogIndex& snapshot_index);Status save_raft_meta(const RaftMeta& meta);Status load_raft_meta(RaftMeta& meta) const;Status load_snapshot_meta(SnapshotPtr& snap) const;Status for_each_snapshot_kv(const KvVisitor& fn) const;Status read_snapshot_chunk(uint64 offset, size_t max_bytes, std::string& data, bool& eof) const;Status append_snapshot_chunk(const InstallSnapshotParam& param);Status finish_snapshot_receive();Status write_snapshot(const StateMachine& state_machine, const std::vector<RaftMember>& members = {});Status clear_wal();struct RecoverResult {SnapshotPtr snapshot;RaftMeta raft_meta;std::vector<LogEntry> wal_entries;bool need_recover{false};};Status recover(RecoverResult& result);private://...
};

这里先大概展示一下 Replica 里的 PersistEngine 提供的各个接口,主要就是 WAL + raft_meta + snapshot 这几个。之后的篇章我们再详谈这个。

七、ReplicaMessageDispatcher 如何发送 AppendEntries

7.1 DispatcherRPC

ReplicaLoop 里面发 RPC 的是 ReplicaMessageDispatcher。内部搞了个 Context,可以使得 ReplicaLoop 获取到这个 MessageDispatcher

主线是:

RaftEffects.messages│▼
ReplicaMessageDispatcher::async_send│▼
ThreadPool│▼
RaftSender│▼
其他 Storage 的 Raft RPC

当前 RPC worker 数量是 8 个。

代码内容还是挺简单的:

Status ReplicaMessageDispatcher::async_send(std::vector<RaftMessage> messages) {for (const RaftMessage& msg : messages) {RETURN_IF_INVALID_STATUS(async_send_one(msg))}return Status::OK();
}Status ReplicaMessageDispatcher::async_send_one(const RaftMessage& msg) {if (!rpc_pool_.started()) {return Status::ERROR("raft rpc workers are not started");}rpc_pool_.submit([this, msg]() { send_task(msg); });return Status::OK();
}
//这里简化描述一下,大概展示一下代码内容其实就是这样的
void ReplicaMessageDispatcher::send_task(RaftMessage msg) {switch (msg.type) {case RaftMessageType::REQUEST_VOTE: {RequestVoteResult result;Status status = sender_.send_request_vote(msg.target, msg.vote_param, result);callback(VoteResponseEvent{msg.target.replica_id, result});break;}case RaftMessageType::APPEND_ENTRIES: {AppendEntriesResult result;Status status = sender_.send_append_entries(msg.target, msg.append_param, result);if (status.ok()) {callback(AppendResponseEvent{msg.target.replica_id, msg.append_param, result});} else {callback(AppendSendFailedEvent{msg.target.replica_id, msg.append_param, status});}break;}}
}void ReplicaMessageDispatcher::callback(Event event) {if (!event_callback_) {LOG_WARN("drop raft response event because event sink is not set");return;}event_callback_(std::move(event));
}

然后这里面的 callback 就是一个回调函数,会去投递 ReplicaLoop 里面的投递 Event 函数。

Dispatcher 支持:

  • RequestVote
  • AppendEntries
  • InstallSnapshot

当然这个里面还是有一些同步的操作的,只不过没有展示出来,在 GET 链路里面会有处理。

由于 PUT 链路基本上就是 AppendEntries 操作,在之后 RPC worker 调用 RaftSender 发送 AppendEntries

Leader Dispatcher-> RaftSender::send_append_entries-> Follower StorageServiceImpl::AppendEntries-> Follower Replica::handle_append_entries-> Follower ReplicaLoop-> Follower RaftCore-> AppendEntriesResponse

会发现还是这样同样的逻辑:关于收到了 AppendEntriesRPC 请求,Replica 也还是会放到 ReplicaLoop 里面,只不过我们不需要同步等待,这是异步操作,封装成 Event 里面去处理就好了。

AppendEntriesResponse│▼
AppendResponseEvent│▼
event_callback_│▼
Leader ReplicaLoop::async_submit

另外提一句,这个 callback 是在 Replica::init 创建 Dispatcher 时传进去的。这个 callback 的工作就是把 Event 重新送回本 ReplicaReplicaLoop

另外,如果 AppendEntries 发送失败时,Dispatcher 会生成 AppendSendFailedEvent

网络/RPC 失败-> AppendSendFailedEvent-> ReplicaLoop-> RaftCore::handle_append_send_failed

RequestVote 发送失败当前主要记录日志;AppendEntriesInstallSnapshot 的发送失败会转换成相应的失败 Event,把结果交回 ReplicaLoop 层进行处理。

八、Follower 如何接住 AppendEntries

到目前为止,我们只跟了 Leader 侧的 PUT。现在把视线转到 follower

Follower 收到 AppendEntries RPC 后,StorageServiceImpl 也会根据 request 里的 ReplicaID 找到本机 Replica,然后调用:

replica->handle_append_entries(param, result);

Replica::handle_append_entries 的主要流程是:

获取 OperGate
检查本地状态│
构造 AppendEntriesCall│
loop_->sync_submit(&call)│
ReplicaLoop::handle_call│
RaftCore::handle_append_entries│
回填 AppendEntriesResult│
notify_apply_task

这里的 Call 是同步的,因为当前 Storage RPC handler 需要拿到:

  • response term
  • response success
  • follower 当前 last_log_index
Status Replica::handle_append_entries(const AppendEntriesParam& param, AppendEntriesResult& result) {RETURN_IF_OPER_GUARD_ACQUIRE_FAILED(oper_gate_)RETURN_IF_INVALID_STATUS(ensure_local_state_running())// 这里是作为follower那边的handle,会更新commit_idx{ADVISKV_METRICS_TIMER("storage_replica_handle_append_entries_raft_step");AppendEntriesCall call{param};loop_->sync_submit(&call);RETURN_IF_INVALID_STATUS(call.status)result = call.result;}notify_apply_task();return Status::OK();
}

这里可以简单看一下这个函数的内容,就是创建了个 AppendEntriesCall,然后投递到 loop 里面。另外,AppendEntriesCall 里面会有 AppendEntriesResult,装着我们期望获得的结果。

FollowerRaftCore 处理内容可以先简单理解为:

  1. 检查 Leader term
  2. 检查 prev_log_indexprev_log_term 是否匹配;
  3. 追加或重写本地日志;
  4. 根据 leader_commit 更新本地 commit_index
  5. 生成 response

RaftCore 处理完后,Replica::handle_append_entries 调用 notify_apply_task,让 follower 后台把已经 commit 的日志 apply 到自己的状态机。

九、response 如何推进 commit

8.1 response 回到 LeaderReplicaLoop

LeaderAppendResponseEvent 进入 ReplicaLoop 后,会执行:

ReplicaLoop::handle_event(AppendResponseEvent)-> ReplicaLoop::run_step-> RaftCore::handle_append_response-> 更新该 follower 的复制进度-> 尝试推进 commit_index-> resolve_pending_proposals

这里前面的逻辑就比较重复了,经典的 Loop 投递 Event,然后最终会走到一个 resolve_pending_proposals 的地方。这个是负责处理我们的写操作的,之前我们的逻辑里面是有关于 PUT 操作的时限的,逻辑上来讲,等到 commit_index 推进到了目标的 index 位置,就可以返回成功了。

ReplicaLoop 内部有:

std::multimap<LogIndex, std::shared_ptr<TimeoutWaiter>>pending_proposals_;

当一批 proposalRaftCore 接受,但它们对应的日志还没有 commit 时,ReplicaLoop 会把 waiterlog index 放进这个 map

proposal 对应 index = 10
commit_index = 9│▼
pending_proposals_[10] = waiter

后续只要 commit_index 推进到 10,resolve_pending_proposals 就会完成这个 waiter

log_index <= commit_index-> complete_proposal_commit(waiter, OK)-> condition_variable 唤醒-> Replica::put 继续执行

代码如下:


// 只有在runner内部访问,目前并没有并发方面的危险
void ReplicaLoop::resolve_pending_proposals() {LogIndex commit_index;{std::lock_guard lock(context_.raft_core_mutex);commit_index = context_.raft_core.commit_index();}for (auto it = pending_proposals_.begin(); it != pending_proposals_.end();) {LogIndex log_index = it->first;std::shared_ptr<TimeoutWaiter>& waiter = it->second;bool cancelled = false;{std::lock_guard lock(waiter->mutex);cancelled = waiter->cancelled;}if (cancelled) {it = pending_proposals_.erase(it);continue;}if (log_index > commit_index) {break;}complete_proposal_commit(waiter, Status::OK());it = pending_proposals_.erase(it);}
}

大致的逻辑就是这样,我们会把每一次的 waiter 存起来,然后去逐个遍历判断,如果达到了,就可以把 waiter 标记成功然后返回了。

十、commit 之后的 ApplyTask

9.1 ReplicaApplyTask 如何被触发

Replica::putproposal waiter 返回成功后调用:

notify_apply_task();

Replica 里定义的 ReplicaApplyTask 是一个 BackgroundTask,当前周期约为 5 ms。run 代码大致如下:


void Replica::apply_committed_entries_from_task() {OperGate::Guard guard;if (oper_gate_.acquire(guard).fail()) {return;//...}if (ensure_local_state_running().fail() || !applier_) {return;//...}//...{std::lock_guard lock(raft_core_mutex_);if (raft_core_->commit_index() <= raft_core_->last_applied()) {return;}}//...std::lock_guard lock(state_machine_mutex_);Status status = fault_if_fail(applier_->apply_committed_entries());//...
}

这个后台任务执行时会:

  1. 获取 OperGate
  2. 确认本地仍是 RUNNING
  3. raft_core_mutex_ 下检查 commit_index 是否领先 last_applied
  4. 获取 state_machine_mutex_
  5. 调用 ReplicaApplier::apply_committed_entries

而这个 ReplicaApplier 内部其实本身也比较简单,就是每一次都提取出来已经 commit 但是还没有 applylog,逐个进行 apply


Status ReplicaApplier::apply_committed_entries() {std::vector<LogEntry> entries;{std::lock_guard lock(context_.raft_core_mutex);entries = context_.raft_core.extract_committed_entries();}for (const LogEntry& entry : entries) {RETURN_IF_INVALID_STATUS(apply_log_entry(entry))}return Status::OK();
}Status ReplicaApplier::apply_log_entry(const LogEntry& entry) {switch (entry.op_type) {case WriteOpType::PUT:case WriteOpType::DEL:case WriteOpType::NONE:return apply_kv_log_entry(entry);case WriteOpType::ADD_LEARNER:case WriteOpType::PROMOTE_VOTER:case WriteOpType::REMOVE_MEMBER:return apply_config_log_entry(entry);}return Status{StatusCode::ERROR, "unsupported raft log entry type"};
}Status ReplicaApplier::apply_kv_log_entry(const LogEntry& entry) {Status status = context_.state_machine.apply(entry);if (status.fail()) {return status;}{std::lock_guard lock(context_.raft_core_mutex);context_.raft_core.advance_last_applied(entry.index);}return Status::OK();
}

这里最终就是交给 state_machine 去进行 apply,而目前代码里面只有个 MAP_ENGINE 为底座的 kv_state_machine,到这里后面的内容其实就很简单了,就是一个 map 去进行 PUT 了。

由于这个项目我觉得应该不需要再在这个存储引擎里面下功夫了,复杂一点可以接入 RocksDB,或者说写个简单的跳表、B+ 等这种都是没有问题的,只不过我觉得这个项目的主旨其实是分布式 KV,所以目前就先用 MAP 了,后续再打算进一步去考虑。

写在结尾

我靠东西实在太多(虽然这个项目也就那样吧,感觉垃圾的地方太多了,突然有一种感觉当做教学项目都不太够格的感觉)。

下一篇打算讲一下 Raft 的内部细节,然后再把 Replica 这次没有提到的链路模块再展开讲讲。