TiKV batch-system 组件深度解析:raftstore 底层的 FSM 批处理执行框架
2026/9/14 2:38:19 网站建设 项目流程

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可执行的状态机抽象,带MessageFSM_TYPEis_stoppedset/take_mailboxget_priorityfsm.rs
FsmScheduler将 FSM 调度给 poller 的抽象(schedule/shutdown/consume_msg_resourcefsm.rs
BasicMailbox消息队列 + FSM 所有权交接容器mailbox.rs
Batch一轮轮询中的 normal + control FSM 集合batch.rs
Router按地址查找 mailbox 并投递消息router.rs
Poller/PollHandlerraftstore 使用的执行钩子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_sendtry_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),它完成:

  1. state_cnt计数器与 control FSM 构造control_box
  2. 创建两条无界调度队列(normal 队列带ResourceController,low 优先级队列不带资源控制);
  3. 构造NormalScheduler(内含 normal/low 两个 sender)与ControlScheduler
  4. 返回(BatchRouter, BatchSystem)对。

raftstore 侧再通过create_raft_batch_systemcreate_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_fsmmax_batch_size会取max(自身配置, 当前 batch 的 normal 数),保证"即使有热点 Region 也不会让其他 Region 饿死";
  • 热点 FSM 限流重调度:若某 normal FSM 被连续轮询超过reschedule_duration,会被标记为热点,且只把一半的热点 FSM 重新调度hot_fsm_count % 2 == 0Schedule),避免下一轮又把所有热点一次性拉满;
  • 优先级错配重调度:当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_durationlow_priority_pool_size标记为online_config(skip),不可在线修改):

参数默认值说明
max_batch_sizeNone(读取时回退到 256)每轮批处理的最大 FSM 数量上限;Config::validate未调用时(如测试环境)取 256
pool_size2normal 优先级 poller 线程数
reschedule_duration5s(ReadableDuration::secs(5)FSM 连续被轮询超过该时长即被视为热点,参与热点重调度
low_priority_pool_size1low 优先级 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!生成。

生命周期与启动/关闭时序

官方指南强调的生命周期要点:

  1. 所有权必须先于运行建立:FSM、mailbox、router 的所有权必须在 poller 启动前完全建立;
  2. 启动:主构造路径为batch.rs::create_system,随后 raftstore 的create_raft_batch_system/create_apply_batch_system完成子系统装配;
  3. IO 分类:poller 线程在start_poller中以set_io_type(IoType::ForegroundWrite)显式声明前台写入 IO(batch.rs),若改动 poller 同步执行的内容,需验证前台/后台 IO 假设是否仍然成立;
  4. 关闭BatchSystem::shutdown(batch.rs)调用router.broadcast_shutdown():置shutdown标志、逐个close所有 normal mailbox、清空 registry、关闭 control mailbox,并调用两个 scheduler 的shutdown
  5. 哨兵机制NormalScheduler::shutdownControlScheduler::shutdown(scheduler.rs)通过向队列发送 256 个FsmTypes::Empty(源码注释说明:任意大于 poll 池大小的数字即可)唤醒所有 poller 退出,配合 mailbox 关闭与 FSM 状态清理,确保没有任何 FSM 被遗留

指南提醒:改动停止路径逻辑时,必须把 batch.rs、mailbox.rs、fsm.rs 放在一起评审。

资源控制集成

指南指出资源控制是契约的一部分:Fsm::Message: ResourceMeteredFsmScheduler::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_totalFSM 重新调度总次数(计数器)
tikv_batch_system_fsm_schedule_wait_secondsFSM 等待被轮询的时长(直方图)0.001 起 20 桶,约 10s
tikv_batch_system_fsm_poll_secondsFSM 处理完所有消息的总时长(可能跨多轮,直方图)约 10s
tikv_batch_system_fsm_poll_roundsFSM 处理完所有消息所需轮询轮数(直方图)1 起 20 桶
tikv_batch_system_fsm_count_per_poll单轮轮询的 FSM 数量(直方图)1 起 20 桶
tikv_channel_full_totalchannel 满错误总数(计数器,标签normal/control
tikv_broadcast_normal_duration_seconds向所有 normal FSM 广播的耗时(直方图)约 10s

MetricsCollector(batch.rs)在 FSM 的Drop时上报poll_roundspoll_duration,每轮tick_round上报count_per_pollFsmState::new/Drop通过state_cnt维护活跃 FSM 计数,Router::trace据其计算泄漏量(leak = total - alive - 1,减 1 代表 control FSM,router.rs)。

官方指南给出运维建议:许多用户可见的症状首先出现在 raftstore 指标而非这里,因此 batch-system 信号通常应与 raftstore 的队列、proposal、apply 延迟信号一起解读。

变更管理与审查要点

变更影响矩阵(来自官方指南)

变更类型需要检查的文件
FSM 所有权或唤醒变更fsm.rsmailbox.rsbatch.rs
Router 或投递变更router.rs、mailbox 语义、raftstore router 调用方
轮询、重调度或批形状变更batch.rsscheduler.rs、raftstore poll handlers
优先级或资源记账变更fsm.rsscheduler.rscomponents/resource_control、raftstore apply/store pollers

关键不变量(Critical Invariants)

  1. 同一时刻至多一个 poller 拥有某个 FSM(由FsmState的 CAS 保证);
  2. release/remove/reschedule 路径必须保持待处理消息的可见性(不得丢失唤醒、不得重复拥有);
  3. Control FSM 与 Normal FSM 的调度语义必须保持区分;
  4. 指标不应扭曲热路径行为;
  5. 队列所有权交接必须保持锁安全与竞态安全。

Review Checklist(来自官方指南)

  • 该变更是否改变了 mailbox 所有权或 release 行为?
  • 是否改变批重调度策略或轮询轮次形状?
  • 是否在 poller 热路径上增加了工作?
  • 是否改变关闭信号或 empty 哨兵处理?
  • 是否改变与资源控制的优先级集成?

若变更影响FsmState、mailbox 关闭或 release/remove 语义,必须明确论证"为何不会丢失唤醒、为何不会出现重复拥有";若影响优先级或资源记账,应以resource_control的视角评审,而非仅考虑 raftstore。

常见故障模式

官方指南归纳的四类高频故障,与源码一一对应:

  1. FSM 重复拥有(duplicate ownership):多线程同时take_fsm——正常情况下被FsmStatecompare_exchange挡住;一旦出现,说明状态机被破坏;
  2. release 后丢失重调度(lost reschedule after release)Batch::releaseexpected_len检查被改动或notify顺序被拆分,导致新消息到达时 FSM 未被唤醒;
  3. mailbox 状态漂移导致 FSM 卡死FsmStateNOTIFIED/IDLE/DROP状态在release/clear/close路径上不一致,panic 提示invalid release state(fsm.rs);
  4. 负载下队列公平性退化:批处理形状、热点重调度策略或资源记账的回归,先表现为 raftstore 队列延迟异常。

测试、基准与配套阅读

  • 测试tests/cases/router.rs覆盖消息投递、mailbox 注册、容量限制与force_send语义(例如"发送应尊重容量限制而 force_send 不必");tests/cases/batch.rs覆盖批处理与重调度;
  • 基准benches/router.rsbenches/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),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询