ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

MongoDB PrimaryOnlyService 机制深度解析:构建可跨故障转移恢复的单主任务框架

MongoDB PrimaryOnlyService 机制深度解析:构建可跨故障转移恢复的单主任务框架 MongoDB PrimaryOnlyService 机制深度解析构建可跨故障转移恢复的单主任务框架【免费下载链接】mongoThe MongoDB Database项目地址: https://gitcode.com/GitHub_Trending/mo/mongoPrimaryOnlyService下称 POS是 MongoDB 副本集主节点侧的一套任务执行框架它允许开发者注册仅在当前节点成为 Primary 时才运行、且必须跨副本集故障转移持续推进直至完成的任务其状态以单文档状态机模型持久化在副本集内新当选的 Primary 可据此在故障转移后重建任务状态、从旧主节点中断处继续执行。读完本文你将掌握 POS 的三大核心类/接口的设计、如何从零定义一个 PrimaryOnlyService含参考实现与单元测试、stepUp/stepDown 期间的行为与三种中断机制以及 Instance 生命周期管理的关键注意事项并能直接对照本仓库源码进行二次开发。设计动机与适用场景副本集中常见的后台任务如 resharding 协调、configsvr DDL 协调等需要满足一个特殊的执行约束任务只能在 Primary 上运行但任务本身不能因为一次选举切换就丢失进度。如果任务状态只存在于内存中主节点一 stepDown 进度就全部丢失如果把状态完整落库又需要一套在每次 stepUp 时扫描状态、重建内存对象、继续执行的通用机制。POS 正是为此设计的通用框架。它的核心建模方式是一个任务 一个状态机 集合中的一条状态文档state document。新当选的 Primary 通过扫描服务专属的状态文档集合来回忆起旧 Primary 未完成的任务从而做到故障转移后无缝续跑。任务只要还没完成状态文档还存在就会在新 Primary 上被重新拉起任务一旦完成状态文档被删除就再也不会被重建。从当前仓库的实现看POS 框架位于 src/mongo/db/repl/primary_only_service.h 与 src/mongo/db/repl/primary_only_service.cpp并被 ReshardingCoordinatorService、ShardingCoordinatorService、ConfigsvrCoordinatorService 等实际服务继承使用。三大核心类/接口POS 框架由三个主要类/接口组成职责划分非常清晰注册表负责生命周期管理Service 负责任务分组的定义与调度Instance 负责单个任务的状态与逻辑。PrimaryOnlyServiceRegistry进程级单例注册表PrimaryOnlyServiceRegistry是一个单例在 mongod 启动时以 ServiceContext 装饰decoration 的形式安装随 mongod 进程存活整个生命周期。具体地注册时机所有 PrimaryOnlyService 必须在 ReplicationCoordinator 启动之前注册完毕因为正是 ReplicationCoordinator 的启动流程会拉起已注册的服务。源码中注册表本身通过ReplicaSetAwareServiceRegistry::Registerer注册为名为PrimaryOnlyServiceRegistry的 ReplicaSetAwareServiceprimary_only_service.cpp#L76-L78 可以看到它声明依赖ShardingInitializationMongoDRegistry。运行时查找lookupServiceByName/lookupServiceByNamespace支持按服务名或状态文档命名空间两种方式查找返回原始指针是安全的因为已注册服务集合在运行期不会变化primary_only_service.cpp#L169-L186。lookupServiceByNamespace在未找到时返回nullptr是 OpObserver 反查服务的入口。状态通知注册表本身是一个 ReplicaSetAwareService通过onStartup、onStepUpComplete、onStepDown、onShutdown等钩子把复制状态变化广播给所有已注册服务。例如 onStepUpComplete 会校验 term 与 last applied optime 一致然后逐个调用每个服务的onStepUp(stepUpOpTime)并用slowTotalOnStepUpCompleteThresholdMS/slowServiceOnStepUpCompleteThresholdMS两个阈值记录耗时告警。重复注册防护registerService会对服务名与状态文档命名空间做唯一性 invariantprimary_only_service.cpp#L145-L167重复注册会直接触发 invariant 失败对应单元测试见PrimaryOnlyServiceTestDeathTest.DoubleRegisterServiceprimary_only_service_test.cpp#L370-L380。PrimaryOnlyService任务分组的定义与调度器PrimaryOnlyService是定义一个新的主节点专属服务所需实现的接口。一个服务本质上是一组任务Instance的分组这些任务只在节点为 Primary 时运行故障转移后在新 Primary 上恢复。每个服务必须声明唯一的服务名getServiceName()唯一的、已复制的状态文档集合getStateDocumentsNS()该集合中每条文档对应一个 Instance 的当前状态。这里有一个从源码可以确认的关键约束虽然原文档表述为most likely in the admin or config databases但当前实现强制要求状态文档集合必须位于 config 数据库。registerService中有一条硬性 invariantprimary_only_service.cpp#L147-L149invariant(ns.isConfigDB(), PrimaryOnlyServices can only register a state documents namespace in the config db);其根本原因是 PrimaryOnlyServiceOpObserver 的getNamespaceFilters()只对 config 库的 delete 操作进行过滤监听。因此定义新服务时状态文档集合必须放在 config 库。在 stepUp 时每个服务会查询自己的状态文档集合为其中每条文档创建并启动一个 Instance——这正是 POS 任务在故障转移后被恢复的方式。此外服务还可以覆盖getThreadPoolLimits()定制承载 Instance 的线程池规模默认值为minThreads1, maxThreads8, maxIdleThreadAge30s见 primary_only_service.h#L199-L203。PrimaryOnlyService::Instance / TypedInstance单个任务的执行单元PrimaryOnlyService::Instance接口封装单个任务的状态与核心逻辑其关键成员包括run(ScopedTaskExecutor, CancellationToken)纯虚主入口任务的全部工作必须调度在传入的 executor 上执行interrupt(Status)被中断时的回调用于解除阻塞例如向未决的 promise 写入错误reportForCurrentOp(connMode, sessionMode)决定该 Instance 是否出现在currentOp()输出中checkIfOptionsConflict(stateDoc)getOrCreateInstance时校验已存在实例与新状态文档是否冲突。InstanceID 就是状态文档的_id字段因此同一服务内 InstanceID 唯一primary_only_service.h#L76-L80。实现者不应该直接继承PrimaryOnlyService::Instance而应继承PrimaryOnlyService::TypedInstanceInstanceType。TypedInstance提供lookup()与getOrCreate()两个静态方法将基类返回的Instance指针安全地checked_pointer_cast为正确的派生类型指针primary_only_service.h#L158-L194。定义一个新的 PrimaryOnlyService完整步骤与参考实现定义一个新服务需要同时编写PrimaryOnlyService与PrimaryOnlyService::TypedInstance两个子类Service 子类只负责两件事指明状态文档存储于哪个集合、按需构建正确类型的 Instance。constructInstance()是必须实现的纯虚方法Instance 子类承载绝大部分业务工作运行完成任务所需的逻辑并自行管理与同步内存态与磁盘态。一个极易被误解的关键设计是POS 机制本身永远不会向状态文档集合写入任何数据。所有对状态文档的写入包括初始创建、中途更新、完成后的删除都由 Instance 实现自行完成。因此大多数服务的run()第一步就是插入初始状态文档——这一写操作确保 Instance 已持久化故障转移后才会被恢复。当 Instance 在故障转移后被恢复时它拿到的是状态文档集合中该文档的当前版本据此重建内存态从而知道现在处于什么状态、还有哪些工作要做、哪些工作旧 Primary 已完成。仓库中的 primary_only_service_test.cpp 定义了一个麻雀虽小五脏俱全的TestService是官方推荐的参考实现其骨架如下class TestService final : public PrimaryOnlyService { public: std::string_view getServiceName() const override { return kTestServiceName; } NamespaceString getStateDocumentsNS() const override { return NamespaceString::createNamespaceString_forTest(config, test_service); } std::shared_ptrPrimaryOnlyService::Instance constructInstance(BSONObj initialState) override { return std::make_sharedTestService::Instance(this, std::move(initialState)); } class Instance final : public PrimaryOnlyService::TypedInstanceInstance { public: SemiFuturevoid run(std::shared_ptrexecutor::ScopedTaskExecutor executor, const CancellationToken token) noexcept override { // 1) 注册取消逻辑stepDown 时向 _completionPromise 置错 // 2) 在 executor 上调度状态机推进_runOnce(状态A, 状态B) 等 // 3) 每个 _runOnce 内更新内存态 写/删状态文档 // 4) 返回 whenAll(cancelLogic, testLogic).semi() } void interrupt(Status status) override { /* 可留空取消逻辑接管 */ } // ... }; };TestService::Instance::_runOnce展示了状态文档读写的典型写法primary_only_service_test.cpp#L223-L272构建新状态文档保留_id、更新state字段若到达终态kDone则client.remove(...)删除文档否则client.update(..., true /*upsert*/)写回。状态机从kInitializing → kOne → kTwo → kDone逐步推进。值得注意TestService还演示了两个进阶能力_rebuildService()钩子在 stepUp 重建 Instance 之前执行额外初始化例如为状态集合创建 TTL 索引见 primary_only_service_test.cpp#L286-L307其内通过AllowOpCtxWhenServiceRebuildingBlock允许重建阶段创建 OpCtx通过MONGO_FAIL_POINT_DEFINE定义TestServiceHangDuringInitialization、TestServiceHangDuringStateOne等故障注入点用于测试挂起/中断场景。状态转换期间的行为stepUp 与 stepDownstepUp异步重建所有 InstancestepUp 时每个服务查询自己的状态文档集合为每条文档创建并启动一个PrimaryOnlyService::Instance。整个重建过程相对于复制核心的 stepUp 流程是异步的——stepUp 完成、RSTL 锁释放时并不能保证所有服务已完成 Instance 重建。从源码看_doStepUpprimary_only_service.cpp#L413-L535的完整流程是将状态置为kRebuilding创建新的CancellationSource安装新的ScopedTaskExecutor并把上一任期的 Instance 与执行器暂存起来join 上一任期先(*newThenOldScopedExecutor)-join()确保旧 executor 上的任务全部结束再对每个旧 Instance 调用waitForCompletion()从而保证同一 InstanceID 的两个 Instance 永远不会同时存在调用_onServiceInitialization()钩子通过WaitForMajorityService等待新任期第一次写操作被 majority 提交确保旧任期对状态文档的所有写入均已提交再执行_rebuildService()服务级初始化如建索引_rebuildInstances(newTerm)全量查询状态文档集合primary_only_service.cpp#L784-L867逐条constructInstance(doc)并插入_activeInstances映射最后把状态置为kRunning。若重建失败如磁盘读失败服务进入kRebuildFailed状态并记录_rebuildStatus此后lookup()/getOrCreate()会抛出该状态直到节点再次 stepDownprimary_only_service.cpp#L580-L585。stepDown中断但不立即回收stepDown 时所有 Instance 被中断但运行其工作的线程不会被 joinInstance 对象及其内存态也不会被释放直到下一次 stepUp。这样设计是为了避免在状态转换过程中阻塞、拖慢整个节点的 stepDown。这种延迟回收行为同时保证了同一服务、同一 InstanceID 的两个 Instance 永远不会在同一节点上同时运行因为旧 Instance 必须在新 Instance 创建前 join 完成见上述_doStepUp流程。onStepDownprimary_only_service.cpp#L554-L573的具体动作包括调用_onServiceTermination()钩子、_interruptInstances中断全部 Instance、将状态置为kPaused、清空_rebuildStatus。stepDown 时的三种中断机制为保证 stepDown 后不再有任何 POS 相关的工作继续执行源码通过三条路径确保 Instance 被彻底中断primary_only_service.cpp#L537-L552关闭 Instance 的执行器stepDown 时ScopedTaskExecutor被 shutdown后续任何工作都无法再以该 Instance 名义调度。服务内部存在双执行器设计进程级常驻的_executor避免每次 stepUp 都重新分配线程/连接资源与每次 stepUp 创建、stepDown 销毁的_scopedExecutor保证 stepDown 时全部未完成任务被中断见 primary_only_service.h#L564-L571。Instance 只能接触_scopedExecutor而_executor通过getInstanceCleanupExecutor()暴露用于 stepDown 后仍需执行的清理工作例如 Instance 完成通知回调。中断所有关联 OpCtx通过PrimaryOnlyServiceClientObserverprimary_only_service.cpp#L94-L132凡是在 POS 线程上创建的 OpCtx 都会被注册进该服务的_opCtxs集合stepDown 时由_interruptInstances统一killOperation。这个 ClientObserver 还会在服务不处于kRunning时让新建 OpCtx一开始就处于已中断状态markKilled(NotWritablePrimary)见 primary_only_service.cpp#L294-L316。显式中断每个 Instance调用每个 Instance 必须实现的interrupt()方法用于解除那些运行在非 POS 自有 executor 线程上、却依赖该 Instance 发信号的工作例如等待 Instance 达到某状态的命令。目前这一机制是显式interrupt()调用原文档指出未来很可能改为向 Instance 拥有的CancellationToken发信号。_interruptInstances在实现上依次完成三件事_source.cancel()取消服务级 CancellationSource级联取消所有 Instance 的 token、(*_scopedExecutor)-shutdown()、遍历_activeInstances逐个interrupt(status)并 kill 全部已注册 OpCtx。stepDown 使用的错误码是InterruptedDueToReplStateChangeshutdown 时则为InterruptedAtShutdown。Instance 生命周期与 shared_ptr 管理Instance 由父服务以shared_ptr持有。生命周期有两个关键释放时机stepDown 时服务释放它拥有的全部 Instanceshared_ptr状态文档被删除时通过 PrimaryOnlyServiceOpObserver 的onDelete回调触发——当状态文档集合中的文档被删除即任务完成后在事务提交回调中调用releaseInstance(instanceId, Status::OK())把该 Instance 从_activeInstances中摘除若整个状态文档集合被 drop非预期操作则releaseAllInstances以Interrupted状态中断并释放全部 Instance。由此产生一个对实现者极其重要的陷阱通常正是 Instancerun()里的逻辑负责删除自己的状态文档而状态文档被删除的那一刻服务就不再持有该 Instance 的 shared_ptr。如果 Instance 在删除状态文档之后还有额外的逻辑或内部状态要更新就必须在删除之前通过shared_from_this()捕获一份指向自身的 shared_ptr 来延长自己的生命周期否则 Instance 可能在其后仍被使用的过程中被提前析构。这一要求在 primary_only_service.h#L108-L123 中有明确的注释强调TestService 的实现也通过在每个回调闭包中捕获self shared_from_this()来演示正确用法。另外需要注意run()是调度在ScopedTaskExecutor上的理论上存在调度 executor 已在 stepDown 时 shutdown、任务从未真正执行的可能因此创建 Instance 并不保证run()一定会被调用实现不应依赖run()来保证析构安全。可观测性与调试手段serverStatusPrimaryOnlyServiceRegistry::reportServiceInfoForServerStatus会输出primaryOnlyServices.serviceName {state, numInstances}其中 state 取值为running/paused/rebuilding/rebuildFailed/shutdown见 primary_only_service.cpp#L906-L921。currentOp每个 Instance 通过reportForCurrentOp决定是否以及如何出现在currentOp()输出中TestService 展示了依据状态文档中reportOp字段控制是否上报的写法。FailPoint 注入框架自身暴露了PrimaryOnlyServiceHangBeforeRebuildingInstances、PrimaryOnlyServiceFailRebuildingInstances、PrimaryOnlyServiceHangBeforeLaunchingStepUpLogic、PrimaryOnlyServiceHangBeforeRunningInstance、PrimaryOnlyServiceSkipRebuildingInstances等故障注入点primary_only_service.h#L52-L55测试中用于构造 stepDown 竞态、挂起重建等场景。仓库中的真实应用POS 并非理论框架而是被 MongoDB 多个核心功能实际使用的生产级基础设施。通过搜索getServiceName() const override可以在仓库中确认以下真实使用方服务源码位置状态文档集合ReshardingCoordinatorServiceresharding_coordinator_service.hconfig.reshardingOperationsReshardingDonorService / ReshardingRecipientServicesrc/mongo/db/s/resharding/resharding 相关 config 集合ShardingCoordinatorServicesharding_coordinator_service.hconfig 库 DDL 协调集合ConfigsvrCoordinatorServiceconfigsvr_coordinator_service.hconfig 库协调集合RenameCollectionParticipantServicerename_collection_participant_service.hconfig 库集合MultiUpdateCoordinatormulti_update_coordinator.hconfig 库集合以 ReshardingCoordinatorService 为例它通过getThreadPoolLimits()定制线程池、通过checkIfConflictsWithOtherInstances实现并发冲突检查、通过getAllReshardingInstances暴露遍历能力——这些正是 POS 框架为复杂业务场景提供的扩展点。这些服务的存在也解释了为什么PrimaryOnlyServiceRegistry的注册要依赖ShardingInitializationMongoDRegistry部分服务如 Resharding、ConfigsvrCoordinator依赖分片初始化状态。总结PrimaryOnlyService 是 MongoDB 副本集架构中主节点专属、跨故障转移续跑类任务的通用答案。其精髓可概括为四条设计原则状态即文档任务状态持久化在 config 库的专属集合中_id即 InstanceID文档存在即任务未完成机制与业务解耦POS 框架只管调度、重建、中断与生命周期绝不写状态文档所有读写由 Instance 自行负责故障转移续跑stepUp 时 join 旧任期实例 → 等待 majority 提交 → 扫描状态文档 → 重建全部 Instance保证同一任务同一时刻只有一个实例存在stepDown 三管齐下关闭 scoped executor 中断全部关联 OpCtx 显式 interrupt 每个 Instance确保旧主不再执行任何 POS 工作。对于需要在此基础上开发新服务的工程师务必牢记两条最容易踩坑的约束状态文档集合必须在 config 库删除状态文档即放弃 shared_ptr 所有权后续工作需先用shared_from_this()自保。【免费下载链接】mongoThe MongoDB Database项目地址: https://gitcode.com/GitHub_Trending/mo/mongo创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表