
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_priorityfsm.rsFsmScheduler将 FSM 调度给 poller 的抽象schedule/shutdown/consume_msg_resourcefsm.rsBasicMailbox消息队列 FSM 所有权交接容器mailbox.rsBatch一轮轮询中的 normal control FSM 集合batch.rsRouter按地址查找 mailbox 并投递消息router.rsPoller/PollHandlerraftstore 使用的执行钩子batch.rsNormal FSM 与 Control FSM 的语义差异router 的注释router.rs明确区分了两种 FSMNormal FSM完成常规工作例如 raftstore 模型中的 peer、apply 模型中的 apply delegateControl FSM完成需要全局视图的工作如创建缺失的 FSM、聚合指标每个系统只有一个 control FSM 和多个 normal FSM。两种 FSM 可以拥有不同的 scheduler虽然并不强制这解释了RouterN, C, Ns, Cs泛型里两个独立 scheduler 的设计源码注释也承认这是出于 Rust 类型系统限制无法写一个同时实现FsmSchedulerFsmC与FsmSchedulerFsmN的 trait的权宜之计。数据模型与所有权契约FsmState三态所有权状态机FsmState是整个框架的核心所有权契约fsm.rs。它由AtomicUsize状态 AtomicPtrN数据指针构成状态定义如下NOTIFYSTATE_NOTIFIED (0)FSM 已被外部执行者取走data持有空指针NOTIFYSTATE_IDLE (1)没有参与者在用 FSMdata拥有 FSMNOTIFYSTATE_DROP (2)FSM 已被丢弃data持有空指针。关键操作take_fsm用 CAS 把IDLE原子地换成NOTIFIED成功后再用swap取走数据指针——只有IDLE才能被取走从源头杜绝两个 poller 同时取走同一个 FSMnotify先take_fsm若取成功则set_mailbox后交给 scheduler若已被取走则什么都不做避免重复调度release把 FSM 归还只有之前状态是NOTIFIED才允许转回IDLE若期间收到DROP则直接释放 Box——这种严格的状态机很小但极易被看似无害的重构破坏指南因此把它列为重点审查对象clear原子地切换为DROP并释放数据。BasicMailbox消息入队与空闲即调度的耦合BasicMailbox持有发送端mpsc::LooseBoundedSenderOwner::Message与ArcFsmStateOwnermailbox.rs。文档注释说明了核心设计生产者投递消息后mailbox 会检查 FSM 是否空闲未被取走和调度若空闲则立即调度从而把 FSM 所有权临时转移给 pollerpoller 处理完后必须通过release归还。force_send与try_send的代码顺序mailbox.rs是正确性敏感的pub fn force_sendS: FsmSchedulerFsm Owner( self, msg: Owner::Message, scheduler: S, ) - Result(), SendErrorOwner::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::releasebatch.rs的语义值得细读fn release(mut self, mut fsm: NormalFsmN, expected_len: usize) - OptionNormalFsmN { 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_systembatch.rs它完成用state_cnt计数器与 control FSM 构造control_box创建两条无界调度队列normal 队列带ResourceControllerlow 优先级队列不带资源控制构造NormalScheduler内含 normal/low 两个 sender与ControlScheduler返回(BatchRouter, BatchSystem)对。raftstore 侧再通过create_raft_batch_system与create_apply_batch_system做子系统级装配见 store.rs。PollHandler一轮轮询的生命周期钩子PollHandlertraitbatch.rs定义了每轮的固定流程源码注释给出伪代码loop { begin if control is ready: handle_control foreach ready normal: handle_normal light_end end }各钩子职责begin(batch_size, update_cfg)每轮最开头调用可借此在线更新配置如max_batch_sizehandle_control返回Some(len)表示直到 control FSM 待处理消息超过len才再次调用返回None则下一轮继续处理同一 FSMhandle_normal返回HandleResult——KeepProcessing表示下一轮继续StopAt { progress, skip_end }表示已处理progress条消息后释放除非有新消息light_end/end轻量收尾与整轮收尾pause批系统即将休眠时调用get_priorityhandler 的优先级。每个 poll 线程拥有自己的 handlerPollHandler无需Sync。Poller::poll批处理主循环Poller::pollbatch.rs是框架的心脏几个值得注意的设计防饥饿每轮结束后重新fetch_fsmmax_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日志。重新调度策略ReschedulePolicybatch.rs内部枚举batch.rsRelease(usize)按expected_len释放Batch::releaseRemoveFSM 已停止时移除Batch::remove要求 mailbox 已空否则返回Some让调用者继续轮询以消费完剩余消息Schedule直接送回 scheduler计数器FSM_RESCHEDULE_COUNTER递增。Batch::scheduleswap_reclaim配合使用先从 batch 槽位取出 FSM 执行对应策略若槽位腾空则用swap_remove回收且必须从大索引向小索引逆序遍历避免swap_remove移动元素时影响后续索引。配置参数详解Config定义在 config.rs并通过OnlineConfig支持在线热更新其中reschedule_duration与low_priority_pool_size标记为online_config(skip)不可在线修改参数默认值说明max_batch_sizeNone读取时回退到 256每轮批处理的最大 FSM 数量上限Config::validate未调用时如测试环境取 256pool_size2normal 优先级 poller 线程数reschedule_duration5sReadableDuration::secs(5)FSM 连续被轮询超过该时长即被视为热点参与热点重调度low_priority_pool_size1low 优先级 poller 线程数max_batch_size会在每轮begin钩子中通过update_cfg闭包在线更新self.max_batch_sizebatch.rs。BatchSystem::spawnbatch.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)显式声明前台写入 IObatch.rs若改动 poller 同步执行的内容需验证前台/后台 IO 假设是否仍然成立关闭BatchSystem::shutdownbatch.rs调用router.broadcast_shutdown()置shutdown标志、逐个close所有 normal mailbox、清空 registry、关闭 control mailbox并调用两个 scheduler 的shutdown哨兵机制NormalScheduler::shutdown与ControlScheduler::shutdownscheduler.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.rsNormalScheduler::consume_msg_resource转发给 resource_control channel 的 sender而ControlScheduler::consume_msg_resource为空操作control FSM 不做资源计量scheduler.rscreate_system中 normal 队列使用resource_ctl构造无界 channellow 队列显式传None无资源控制batch.rs。raftstore 侧通过ResourceGroupManager::derive_controller为 store 批系统与 apply 批系统派生控制器store.rs而 store 批系统本身传None注释说明 store 批系统不做优先级调度。因此涉及优先级或资源记账的改动应同时对照 components/resource_control 评审而不仅限于 raftstore。可观测性关键指标解读指标定义位于 metrics.rs核心信号按type标签区分store/apply两种 FSM 类型指标含义桶配置maxtikv_batch_system_fsm_reschedule_totalFSM 重新调度总次数计数器—tikv_batch_system_fsm_schedule_wait_secondsFSM 等待被轮询的时长直方图0.001 起 20 桶约 10stikv_batch_system_fsm_poll_secondsFSM 处理完所有消息的总时长可能跨多轮直方图约 10stikv_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 广播的耗时直方图约 10sMetricsCollectorbatch.rs在 FSM 的Drop时上报poll_rounds与poll_duration每轮tick_round上报count_per_pollFsmState::new/Drop通过state_cnt维护活跃 FSM 计数Router::trace据其计算泄漏量leak total - alive - 1减 1 代表 control FSMrouter.rs。官方指南给出运维建议许多用户可见的症状首先出现在 raftstore 指标而非这里因此 batch-system 信号通常应与 raftstore 的队列、proposal、apply 延迟信号一起解读。变更管理与审查要点变更影响矩阵来自官方指南变更类型需要检查的文件FSM 所有权或唤醒变更fsm.rs、mailbox.rs、batch.rsRouter 或投递变更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 releaseBatch::release的expected_len检查被改动或notify顺序被拆分导致新消息到达时 FSM 未被唤醒mailbox 状态漂移导致 FSM 卡死FsmState的NOTIFIED/IDLE/DROP状态在release/clear/close路径上不一致panic 提示invalid release statefsm.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),仅供参考