TiKV batch-system 组件深度解析:raftstore 底层的 FSM 批处理执行框架
【免费下载链接】tikvDistributed transactional key-value database, originally created to complement TiDB项目地址: https://gitcode.com/GitHub_Trending/ti/tikv
导读
components/batch-system是 TiKV 中位于 raftstore 之下的通用 FSM(有限状态机)执行框架,负责 mailbox 消息路由、轮询、批处理、重新调度与线程池级指标采集。本文以官方维护指南为主体,结合组件源码与 raftstore 调用点,系统讲解它的架构模型、所有权契约、生命周期、配置参数、可观测性与变更风险,帮助读者理解这个"小而关键"的 crate 如何在 TiKV 的写入与 Apply 路径上保持正确性、公平性与低延迟。
组件定位:raftstore 脚下的通用执行框架
在 TiKV 的分层架构中,raftstore 负责 Raft 共识与 Region 状态机的驱动,而 batch-system 正是这套状态机得以并发、批量执行的运行骨架。官方维护指南对它的定位非常明确:
batch-systemis the generic FSM execution framework underneath raftstore. It owns mailbox routing, polling, batching, rescheduling, and pool-level metrics.
它由更高层的运行时拥有者(主要是 raftstore)实例化,本身不感知具体业务,只提供一套抽象的"FSM + Mailbox + Router + Poller"执行模型。指南中特别强调:这个 crate 很小,但它位于关键路径上,细微改动可能改变 raftstore 的公平性、延迟、关闭行为与队列背压。
主要使用者
- components/raftstore/src/store/fsm/store.rs 中的
create_raft_batch_system:创建 store 批处理系统(RaftBatchSystem),其中 store 批处理系统调用batch_system::create_system时传入None(不做优先级调度),同时内部再创建 apply 批处理系统; - components/raftstore/src/store/fsm/apply.rs 中的
create_apply_batch_system:创建 apply 批处理系统,处理 Apply FSM。
架构视图与核心概念
官方指南给出的执行视图是四个步骤:mailbox 投递消息 → router 定位目标 FSM → batch 将 normal 与 control FSM 聚合进同一轮轮询 → poller 与 handler 执行工作。
对应到源码(见 lib.rs 的导出),核心抽象为:
| 概念 | 职责 | 源码位置 |
|---|---|---|
Fsm | 可执行的状态机抽象,带Message、FSM_TYPE、is_stopped、set/take_mailbox、get_priority | fsm.rs |
FsmScheduler | 将 FSM 调度给 poller 的抽象(schedule/shutdown/consume_msg_resource) | fsm.rs |
BasicMailbox | 消息队列 + FSM 所有权交接容器 | mailbox.rs |
Batch | 一轮轮询中的 normal + control FSM 集合 | batch.rs |
Router | 按地址查找 mailbox 并投递消息 | router.rs |
Poller/PollHandler | raftstore 使用的执行钩子 | batch.rs |
Normal FSM 与 Control FSM 的语义差异
router 的注释(router.rs)明确区分了两种 FSM:
- Normal FSM:完成常规工作,例如 raftstore 模型中的 peer、apply 模型中的 apply delegate;
- Control FSM:完成需要全局视图的工作(如创建缺失的 FSM、聚合指标),每个系统只有一个 control FSM 和多个 normal FSM。
两种 FSM 可以拥有不同的 scheduler(虽然并不强制),这解释了Router<N, C, Ns, Cs>泛型里两个独立 scheduler 的设计,源码注释也承认这是出于 Rust 类型系统限制(无法写一个同时实现FsmScheduler<Fsm=C>与FsmScheduler<Fsm=N>的 trait)的权宜之计。
数据模型与所有权契约
FsmState:三态所有权状态机
FsmState是整个框架的核心所有权契约(fsm.rs)。它由AtomicUsize状态 +AtomicPtr<N>数据指针构成,状态定义如下:
NOTIFYSTATE_NOTIFIED (0):FSM 已被外部执行者取走,data持有空指针;NOTIFYSTATE_IDLE (1):没有参与者在用 FSM,data拥有 FSM;NOTIFYSTATE_DROP (2):FSM 已被丢弃,data持有空指针。
关键操作:
take_fsm:用 CAS 把IDLE原子地换成NOTIFIED,成功后再用swap取走数据指针——只有IDLE才能被取走,从源头杜绝两个 poller 同时取走同一个 FSM;notify:先take_fsm,若取成功则set_mailbox后交给 scheduler;若已被取走则什么都不做,避免重复调度;release:把 FSM 归还,只有之前状态是NOTIFIED才允许转回IDLE,若期间收到DROP则直接释放 Box——这种严格的状态机"很小但极易被看似无害的重构破坏",指南因此把它列为重点审查对象;clear:原子地切换为DROP并释放数据。
BasicMailbox:消息入队与"空闲即调度"的耦合
BasicMailbox持有发送端mpsc::LooseBoundedSender<Owner::Message>与Arc<FsmState<Owner>>(mailbox.rs)。文档注释说明了核心设计:生产者投递消息后,mailbox 会检查 FSM 是否空闲(未被取走和调度),若空闲则立即调度,从而把 FSM 所有权临时转移给 poller,poller 处理完后必须通过release归还。
force_send与try_send的代码顺序(mailbox.rs)是正确性敏感的:
pub fn force_send<S: FsmScheduler<Fsm = Owner>>( &self, msg: Owner::Message, scheduler: &S, ) -> Result<(), SendError<Owner::Message>> { scheduler.consume_msg_resource(&msg); // 1. 资源记账 self.sender.force_send(msg)?; // 2. 消息入队 self.state.notify(scheduler, Cow::Borrowed(self)); // 3. 空闲则调度 Ok(()) }指南警告:任何拆分或重排这三步的改动都可能导致丢失唤醒(lost wakeups)或重复调度(duplicate scheduling)。此外close()同时执行sender.close_sender()与state.clear(),也是关闭路径的核心。
Batch::release:依赖队列长度检查的消息可见性保证
Batch::release(batch.rs)的语义值得细读:
fn release(&mut self, mut fsm: NormalFsm<N>, expected_len: usize) -> Option<NormalFsm<N>> { let mailbox = fsm.take_mailbox().unwrap(); mailbox.release(fsm.fsm); if mailbox.len() == expected_len { None } else { // 有新消息到达:重新在本 poller 中调度,或已被其他线程调度 match mailbox.take_fsm() { ... } } }当 FSM 的待处理消息数与expected_len不一致时,说明 release 之后又有新消息入队,需要尝试重新取回 FSM 继续处理。指南明确指出:这段逻辑依赖队列长度检查和 mailbox 的重新取回(re-taking)来保持消息可见性,评审时应视为正确性敏感逻辑,而非仅性能敏感逻辑。
调度与轮询:BatchSystem 如何运行
create_system:运行时构造入口
官方指南指出的主构造路径是batch.rs::create_system(batch.rs),它完成:
- 用
state_cnt计数器与 control FSM 构造control_box; - 创建两条无界调度队列(normal 队列带
ResourceController,low 优先级队列不带资源控制); - 构造
NormalScheduler(内含 normal/low 两个 sender)与ControlScheduler; - 返回
(BatchRouter, BatchSystem)对。
raftstore 侧再通过create_raft_batch_system与create_apply_batch_system做子系统级装配(见 store.rs)。
PollHandler:一轮轮询的生命周期钩子
PollHandlertrait(batch.rs)定义了每轮的固定流程,源码注释给出伪代码:
loop { begin if control is ready: handle_control foreach ready normal: handle_normal light_end end }各钩子职责:
begin(batch_size, update_cfg):每轮最开头调用,可借此在线更新配置(如max_batch_size);handle_control:返回Some(len)表示直到 control FSM 待处理消息超过len才再次调用;返回None则下一轮继续处理同一 FSM;handle_normal:返回HandleResult——KeepProcessing表示下一轮继续,StopAt { progress, skip_end }表示已处理progress条消息后释放(除非有新消息);light_end/end:轻量收尾与整轮收尾;pause:批系统即将休眠时调用;get_priority:handler 的优先级。
每个 poll 线程拥有自己的 handler(PollHandler无需Sync)。
Poller::poll:批处理主循环
Poller::poll(batch.rs)是框架的心脏,几个值得注意的设计:
- 防饥饿:每轮结束后重新
fetch_fsm,max_batch_size会取max(自身配置, 当前 batch 的 normal 数),保证"即使有热点 Region 也不会让其他 Region 饿死"; - 热点 FSM 限流重调度:若某 normal FSM 被连续轮询超过
reschedule_duration,会被标记为热点,且只把一半的热点 FSM 重新调度(hot_fsm_count % 2 == 0才Schedule),避免下一轮又把所有热点一次性拉满; - 优先级错配重调度:当
p.get_priority() != self.handler.get_priority()时,FSM 被标记为ReschedulePolicy::Schedule,送回对应优先级的队列; - 关闭信号:
FsmTypes::Empty是 scheduler 关闭的哨兵,push返回false时主循环退出;退出前会把手头剩余的 control/normal FSM 全部调度回去,并打印poller will exit日志。
重新调度策略(ReschedulePolicy)
batch.rs内部枚举(batch.rs):
Release(usize):按expected_len释放(Batch::release);Remove:FSM 已停止时移除(Batch::remove,要求 mailbox 已空,否则返回Some让调用者继续轮询以消费完剩余消息);Schedule:直接送回 scheduler(计数器FSM_RESCHEDULE_COUNTER递增)。
Batch::schedule+swap_reclaim配合使用:先从 batch 槽位取出 FSM 执行对应策略,若槽位腾空则用swap_remove回收,且必须从大索引向小索引逆序遍历,避免swap_remove移动元素时影响后续索引。
配置参数详解
Config定义在 config.rs,并通过OnlineConfig支持在线热更新(其中reschedule_duration与low_priority_pool_size标记为online_config(skip),不可在线修改):
| 参数 | 默认值 | 说明 |
|---|---|---|
max_batch_size | None(读取时回退到 256) | 每轮批处理的最大 FSM 数量上限;Config::validate未调用时(如测试环境)取 256 |
pool_size | 2 | normal 优先级 poller 线程数 |
reschedule_duration | 5s(ReadableDuration::secs(5)) | FSM 连续被轮询超过该时长即被视为热点,参与热点重调度 |
low_priority_pool_size | 1 | low 优先级 poller 线程数 |
max_batch_size会在每轮begin钩子中通过update_cfg闭包在线更新self.max_batch_size(batch.rs)。BatchSystem::spawn(batch.rs)会按pool_size启动{name_prefix}-{i}线程、按low_priority_pool_size启动{name_prefix}-low-{i}线程,线程名由thd_name!生成。
生命周期与启动/关闭时序
官方指南强调的生命周期要点:
- 所有权必须先于运行建立:FSM、mailbox、router 的所有权必须在 poller 启动前完全建立;
- 启动:主构造路径为
batch.rs::create_system,随后 raftstore 的create_raft_batch_system/create_apply_batch_system完成子系统装配; - IO 分类:poller 线程在
start_poller中以set_io_type(IoType::ForegroundWrite)显式声明前台写入 IO(batch.rs),若改动 poller 同步执行的内容,需验证前台/后台 IO 假设是否仍然成立; - 关闭:
BatchSystem::shutdown(batch.rs)调用router.broadcast_shutdown():置shutdown标志、逐个close所有 normal mailbox、清空 registry、关闭 control mailbox,并调用两个 scheduler 的shutdown; - 哨兵机制:
NormalScheduler::shutdown与ControlScheduler::shutdown(scheduler.rs)通过向队列发送 256 个FsmTypes::Empty(源码注释说明:任意大于 poll 池大小的数字即可)唤醒所有 poller 退出,配合 mailbox 关闭与 FSM 状态清理,确保没有任何 FSM 被遗留。
指南提醒:改动停止路径逻辑时,必须把 batch.rs、mailbox.rs、fsm.rs 放在一起评审。
资源控制集成
指南指出资源控制是契约的一部分:Fsm::Message: ResourceMetered且FsmScheduler::consume_msg_resource影响调度公平性与记账。
从源码可见三处落点:
force_send/try_send在消息入队前调用scheduler.consume_msg_resource(&msg)完成资源记账(mailbox.rs);NormalScheduler::consume_msg_resource转发给 resource_control channel 的 sender,而ControlScheduler::consume_msg_resource为空操作(control FSM 不做资源计量,scheduler.rs);create_system中 normal 队列使用resource_ctl构造无界 channel,low 队列显式传None(无资源控制,batch.rs)。
raftstore 侧通过ResourceGroupManager::derive_controller为 store 批系统与 apply 批系统派生控制器(store.rs),而 store 批系统本身传None(注释说明 store 批系统不做优先级调度)。因此,涉及优先级或资源记账的改动,应同时对照 components/resource_control 评审,而不仅限于 raftstore。
可观测性:关键指标解读
指标定义位于 metrics.rs,核心信号按type标签区分store/apply两种 FSM 类型:
| 指标 | 含义 | 桶配置(max) |
|---|---|---|
tikv_batch_system_fsm_reschedule_total | FSM 重新调度总次数(计数器) | — |
tikv_batch_system_fsm_schedule_wait_seconds | FSM 等待被轮询的时长(直方图) | 0.001 起 20 桶,约 10s |
tikv_batch_system_fsm_poll_seconds | FSM 处理完所有消息的总时长(可能跨多轮,直方图) | 约 10s |
tikv_batch_system_fsm_poll_rounds | FSM 处理完所有消息所需轮询轮数(直方图) | 1 起 20 桶 |
tikv_batch_system_fsm_count_per_poll | 单轮轮询的 FSM 数量(直方图) | 1 起 20 桶 |
tikv_channel_full_total | channel 满错误总数(计数器,标签normal/control) | — |
tikv_broadcast_normal_duration_seconds | 向所有 normal FSM 广播的耗时(直方图) | 约 10s |
MetricsCollector(batch.rs)在 FSM 的Drop时上报poll_rounds与poll_duration,每轮tick_round上报count_per_poll;FsmState::new/Drop通过state_cnt维护活跃 FSM 计数,Router::trace据其计算泄漏量(leak = total - alive - 1,减 1 代表 control FSM,router.rs)。
官方指南给出运维建议:许多用户可见的症状首先出现在 raftstore 指标而非这里,因此 batch-system 信号通常应与 raftstore 的队列、proposal、apply 延迟信号一起解读。
变更管理与审查要点
变更影响矩阵(来自官方指南)
| 变更类型 | 需要检查的文件 |
|---|---|
| FSM 所有权或唤醒变更 | fsm.rs、mailbox.rs、batch.rs |
| Router 或投递变更 | router.rs、mailbox 语义、raftstore router 调用方 |
| 轮询、重调度或批形状变更 | batch.rs、scheduler.rs、raftstore poll handlers |
| 优先级或资源记账变更 | fsm.rs、scheduler.rs、components/resource_control、raftstore apply/store pollers |
关键不变量(Critical Invariants)
- 同一时刻至多一个 poller 拥有某个 FSM(由
FsmState的 CAS 保证); - release/remove/reschedule 路径必须保持待处理消息的可见性(不得丢失唤醒、不得重复拥有);
- Control FSM 与 Normal FSM 的调度语义必须保持区分;
- 指标不应扭曲热路径行为;
- 队列所有权交接必须保持锁安全与竞态安全。
Review Checklist(来自官方指南)
- 该变更是否改变了 mailbox 所有权或 release 行为?
- 是否改变批重调度策略或轮询轮次形状?
- 是否在 poller 热路径上增加了工作?
- 是否改变关闭信号或 empty 哨兵处理?
- 是否改变与资源控制的优先级集成?
若变更影响FsmState、mailbox 关闭或 release/remove 语义,必须明确论证"为何不会丢失唤醒、为何不会出现重复拥有";若影响优先级或资源记账,应以resource_control的视角评审,而非仅考虑 raftstore。
常见故障模式
官方指南归纳的四类高频故障,与源码一一对应:
- FSM 重复拥有(duplicate ownership):多线程同时
take_fsm——正常情况下被FsmState的compare_exchange挡住;一旦出现,说明状态机被破坏; - release 后丢失重调度(lost reschedule after release):
Batch::release的expected_len检查被改动或notify顺序被拆分,导致新消息到达时 FSM 未被唤醒; - mailbox 状态漂移导致 FSM 卡死:
FsmState的NOTIFIED/IDLE/DROP状态在release/clear/close路径上不一致,panic 提示invalid release state(fsm.rs); - 负载下队列公平性退化:批处理形状、热点重调度策略或资源记账的回归,先表现为 raftstore 队列延迟异常。
测试、基准与配套阅读
- 测试:
tests/cases/router.rs覆盖消息投递、mailbox 注册、容量限制与force_send语义(例如"发送应尊重容量限制而 force_send 不必");tests/cases/batch.rs覆盖批处理与重调度; - 基准:
benches/router.rs与benches/batch-system.rs提供路由与批系统的性能基准; - 配套阅读:官方指南建议按
fsm.rs → mailbox.rs → scheduler.rs → batch.rs → router.rs的顺序阅读源码,并以 components/raftstore/src/store/fsm/store.rs 与 components/raftstore/src/store/fsm/apply.rs 作为实际使用范例;配套文档见 repo-overview.md 与 components/raftstore.md。
许多回归会先以 raftstore 延迟或队列异常的形式暴露,因此验证往往还需要 raftstore 级别的测试配合。总体来说,任何改动都应被当作横切变更处理:mailbox、所有权或调度语义一旦变化,本指南与 raftstore 指南需要同步更新,并避免在热路径上随意增加埋点、分配或同步原语——因为该 crate 同时位于 store 与 apply FSM 执行之下。
【免费下载链接】tikvDistributed transactional key-value database, originally created to complement TiDB项目地址: https://gitcode.com/GitHub_Trending/ti/tikv
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考