news 2026/9/14 13:40:37

MongoDB 内部客户端:副本集监控、主机定向与重试策略的源码解析

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
MongoDB 内部客户端:副本集监控、主机定向与重试策略的源码解析

MongoDB 内部客户端:副本集监控、主机定向与重试策略的源码解析

【免费下载链接】mongoThe MongoDB Database项目地址: https://gitcode.com/GitHub_Trending/mo/mongo

MongoDB 服务进程(mongod/mongos)本身也是一个“客户端”:当 mongos 需要向某个分片副本集发送命令时,必须先回答“发往哪个节点”的问题。本文基于src/mongo/client模块的 README 与对应源码,系统讲解副本集监控(Replica Set Monitoring)、主机定向(Host Targeting)以及RetryStrategy重试机制三者的协作方式,帮助读者理解内部客户端如何做节点选择、如何感知拓扑变化(stepdown、选举、reconfig),以及在节点过载时如何通过退避与令牌桶控制重试行为。

一、内部客户端的定位

外部驱动(Driver)与服务器内部代码共用一套 SDAM(Server Discovery And Monitoring)语义:根据$readPreference和节点网络延迟(RTT)在拓扑中挑选可用节点。src/mongo/client目录承载的就是这套“内部客户端”实现,其核心职责包括:

  • 维持对目标副本集的本地拓扑视图(TopologyDescription),持续刷新;
  • 依据读偏好(Read Preference)与 RTT 做主机定向;
  • 通过RemoteCommandTargeter对选出的主机做失败记账(not primary / 网络错误 / shutdown);
  • 通过RetryStrategy接口族统一管理“何时重试、等多久重试、重试时避开哪些节点”。

二、主机定向:读偏好、RTT 与 TargetingMetadata

2.1 节点资格判定

将一条命令路由到副本集时,内部客户端需要找到满足$readPreference的成员。节点的“资格”取决于两个维度:

  1. 节点类型:primary、secondary 等(来自监控学到的拓扑状态);
  2. 平均网络往返时延(RTT):例如{ $readPreference: secondary }要求客户端先知道哪些节点是 secondary,再在这些节点中选出 RTT 处于最小 RTT 一定 delta 范围内的节点。

当“主机选择超时”且没有任何合格节点时,会抛出FailedToSatisfyReadPreference错误——这一点在 ReplicaSetMonitorInterface 的getHostOrRefresh注释中被明确记录为已知错误。

2.2 TargetingMetadata:定向提示

TargetingMetadata用于给定向过程提供“提示”。最典型的用法是:报告自身过载、返回带SystemOverloadedError标签错误的节点,会被加入降权服务器列表(deprioritizedServers);定向时优先选择其他符合条件的节点,仅在没有其他可选节点时才回落到降权节点。结构定义见 targeting_metadata.h:

struct TargetingMetadata { struct Stats { // 定向时成功避开降权节点的次数(无降权节点时不计) Atomic<int64_t> numTargetingAvoidedDeprioritized; }; // 应尽量避免定向的服务器列表;若无其他合格节点,仍可能被选中 std::vector<HostAndPort> deprioritizedServers; std::shared_ptr<Stats> stats = {}; };

从源码结构看,该结构随RetryStrategy一起流转:DefaultRetryStrategy持有一个TargetingMetadata成员(见 retry_strategy.h),并在返回SystemOverloadedError时把对应服务器追加进列表,供下一轮定向使用。

2.3 RemoteCommandTargeter 的失败记账

RemoteCommandTargeter 是“给定副本集或独立节点”之上的定向接口。它的findHost会阻塞最长gDefaultFindReplicaSetHostTimeoutMS毫秒以等待合格节点出现;请求返回后,调用方应通过updateHostWithStatus把结果反馈给 targeter(见 remote_command_targeter.h):

void updateHostWithStatus(const HostAndPort& host, const Status& status) { if (status.isOK()) return; if (ErrorCodes::isNotPrimaryError(status.code())) markHostNotPrimary(host, status); else if (ErrorCodes::isNetworkError(status.code())) markHostUnreachable(host, status); else if (ErrorCodes::isShutdownError(status.code())) markHostShuttingDown(host, status); };

这三类记账(非 primary、不可达、正在关闭)都会让 monitor 更新本地视图,避免后续请求继续选中问题节点。defaultFindReplicaSetHostTimeoutMS是一个默认 15000 毫秒的服务器参数(当前标记为 test-only),定义于 replica_set_monitor_server_parameters.idl。

三、副本集监控:从扫描到 Streamable 协议

3.1 每个副本集一个 ReplicaSetMonitor

副本集监控的本质是周期性刷新本地拓扑视图。客户端为它需要定向的每个副本集维护一个ReplicaSetMonitor:如果 mongos 需要为一个查询定向到 2 个分片,它就会持有(或创建)这 2 个分片副本集对应的 monitor。monitor 的获取入口在 ReplicaSetMonitorManager,例如getMonitor(setName)getMonitorForHost(host)

ReplicaSetMonitorInterface(源码)暴露了监控侧的关键能力:

方法作用
init()/drop()调度首次刷新 / 结束所有进行中的刷新
getHostOrRefresh(readPref, targetingMetadata, cancelToken)按读偏好返回一个主机;无合格节点时报FailedToSatisfyReadPreference
getHostsOrRefresh(...)返回满足条件的主机集合(用于多目标场景)
getPrimaryOrUassert()返回当前视图中的 primary,找不到则 uassert
failedHost / failedHostPreHandshake / failedHostPostHandshake上报主机失败,区分握手前/后的失败(SDAM 语义)
isPrimary / isHostUp / contains基于本地缓存视图的快速查询(可能过期)
pingTime(server)返回某节点的 ping 时延
getMinWireVersion / getMaxWireVersion副本集整体支持的最低/最高 wire 版本

3.2 awaitable hello 与两种协议版本

SDAM 规范支持awaitable helloisMaster/hello命令可以等待一次“重要拓扑变化”或超时后再返回。启用它的客户端能比固定间隔轮询更早感知 stepdown、选举、reconfig 等事件。

仓库当前实现支持两个协议版本(枚举ReplicaSetMonitorProtocol,见 replica_set_monitor_server_parameters.h):

协议特点
sdam符合 SDAM 规范,但不支持awaitable hello + exhaust
streamable(默认)支持 awaitable hello 与 exhaust,拓扑事件以“流”的形式推给客户端

通过启动参数replicaSetMonitorProtocol选择(仅 startup 可设),定义于 replica_set_monitor_server_parameters.idl。从 replica_set_monitor_server_parameters.cpp 可见全局默认值为kStreamable

ReplicaSetMonitorProtocol gReplicaSetMonitorProtocol{ReplicaSetMonitorProtocol::kStreamable};

(关于为何同时存在 hello 与 isMaster,可延伸阅读 Replication Architecture Guide。)

3.3 StreamableReplicaSetMonitor 的内部结构

在 streamable 协议下,StreamableReplicaSetMonitor 负责收集并维护客户端的本地拓扑描述(TopologyDescription),其中保存了副本集每个成员的已学习状态。它组合了以下组件(见 成员声明):

  • sdam::TopologyManager:拓扑状态机,产出TopologyDescription
  • sdam::ServerSelector:按读偏好从拓扑中筛选节点;
  • ServerDiscoveryMonitor:发送(awaitable)hello 的心跳/发现通道;
  • ServerPingMonitor因为支持 exhaust,RTT 不再依赖 hello 响应时延,而是以固定频率向拓扑中每个节点发送ping测量(README 中的关键设计点,对应事件回调onServerPingSucceededEvent/onServerPingFailedEvent,见 streamable_replica_set_monitor.h);
  • StreamableReplicaSetMonitorQueryProcessor/DiscoveryTimeProcessor:当存在未决的主机查询(_outstandingQueries)时注册为TopologyListener,拓扑一变就尽快唤醒等待者,配合kExpeditedRefreshPeriod = 500ms的加速刷新常量(L78)。

因此,满足读偏好所需的信息由两条通道异步汇聚:RTT 来自独立 ping,其余(成员角色、拓扑变化)来自异步发送到各节点的 awaitable hello。

四、RetryStrategy:重试条件、退避与令牌桶

4.1 接口契约

RetryStrategy(retry_strategy.h)是一个标准接口,定义“何种失败应当重试、两次重试之间等多久”,同时追踪可用于重试定向的TargetingMetadata,并预留重试指标回调。核心虚函数:

  • recordFailureAndEvaluateShouldRetry(status, origin, errorLabels, baseBackoffMS):每次失败请求结束时调用,返回是否应重试;
  • recordSuccess(origin):即使没有发生重试,成功时也要调用;
  • getNextRetryDelay():返回带抖动的退避时长;
  • getTargetingMetadata():拿到降权节点列表;
  • recordBackoff(delay):记录实际发生的退避时间。

头文件给出的参考用法(retry_strategy.h):

Status status = ...; while (!status.isOK() && strategy.recordFailureAndEvaluateShouldRetry(status, target, labels)) { wait_for(strategy.getNextRetryDelay()); status = ...; } if (status.isOK()) { strategy.recordSuccess(); }

调用方通常不直接调用这些方法,而是把RetryStrategy传给支持重试的客户端(如async_rpcAsyncTry)。此外还提供一个独立函数runWithRetryStrategy(retry_strategy.h),它会循环执行传入的可调用对象,直到成功、被中断或策略判定不再重试;可调用对象必须接收const TargetingMetadata&并返回RetryStrategy::Result<T>(该 Result 类型同时封装值、错误标签、来源节点和服务器下发的 base backoff)。更底层的客户端(如TaskExecutor)不直接对接RetryStrategy,因为它们服务于 ARS 等高层客户端——真正用到策略的是高层。

4.2 DefaultRetryStrategy

DefaultRetryStrategy(retry_strategy.h)是其他策略的基座:

  • 默认重试条件defaultRetryCriteria检查错误是否属于RetriableError类别,或错误标签中包含可重试标签(如 retryable write 相关 label),实现见 retry_strategy.cpp;
  • 退避:对带SystemOverloadedError标签的错误使用带抖动的指数退避。底层是BackoffWithJitter(backoff_with_jitter.h),公式为j * min(maxBackoff, baseBackoff * 2^i)(i 为重试次数,j 为 [0,1) 的随机抖动);
  • 定向提示:任何返回SystemOverloadedError标签错误的服务器都会被加入策略TargetingMetadata中的降权列表;
  • 支持自定义重试条件回调(RetryCriteria)。

RetryParameters由三个全局服务器参数驱动,定义于 retry_strategy_server_parameters.idl:

服务器参数默认值可设置时机说明
defaultClientBaseBackoffMillis100startup / runtime指数退避的基数;若大于defaultClientMaxBackoffMillis则改用后者
defaultClientMaxBackoffMillis10000startup / runtime单次退避的时长上限
defaultClientMaxRetryAttempts3startup / runtime遇到可重试错误时的最大重试次数

另外,当服务器在过载错误响应中携带了不同的 base backoff 时,会通过ResultbaseBackoffMS传入recordFailureAndEvaluateShouldRetry,覆盖策略自身的默认基数(仅在错误标签含SystemOverloadedError时生效,见 接口注释)。

4.3 AdaptiveRetryStrategy 与令牌桶 RetryBudget

AdaptiveRetryStrategy(retry_strategy.h)包装另一个策略(默认包装DefaultRetryStrategy),引入RetryBudget概念,用令牌桶给重试次数设上限:

  • 每个以SystemOverloadedError失败的请求都会先尝试获取一个令牌;桶耗尽则直接不重试(tryAcquireToken返回 false);
  • 其他任何结果(包括普通错误)都会按returnRate回补令牌——这被视为系统在恢复的信号;
  • 若之前为某次重试消耗过令牌、而该请求最终成功,recordSuccess会在回补之外再额外归还一枚完整令牌(recordNotOverloaded(returnExtraToken))。

效果是:大量重试都失败时,重试会“短路”(不再放大压力);而重试大多成功时则继续执行。RetryBudget的关键字段是capacity(令牌上限)与returnRate(每个成功请求回补的令牌量),并支持通过updateRateParameters在线调整(见 RetryBudget 定义)。

4.4 Shard 场景与其他实现

  • Shard::RetryStrategy:面向分片场景的AdaptiveRetryStrategy变体,其预算按“整个分片的所有节点”共享,并且用Shard::RetryPolicy(而非默认条件)判定哪些错误可重试。mongos 的AsyncRequestsSender(ARS)接口本身不接收RetryStrategy,但内部基于调用方传入的Shard::RetryPolicy构造了这个分片级策略(可参考 async_requests_sender.h 中的_retryPolicy成员)。
  • NoRetryStrategy:当前代码中还有一个“永不再试”的实现,所有评估方法返回不重试(retry_strategy.h),用于需要彻底禁用重试的场景——这是 README 之外、当前源码已具备的能力。
  • RetryStrategyWithFailureRetryHook:一个模板包装器,在底层策略判定重试时触发传入的回调(retry_strategy.h),便于在不改变策略语义的前提下挂载钩子逻辑。

单元测试覆盖了上述行为,可参考 retry_strategy_test.cpp、backoff_with_jitter_test.cpp 与 remote_command_retry_scheduler_test.cpp;协议层行为另有 replica_set_monitor_integration_test.cpp 与 streamable_replica_set_monitor_heartbeat_test.cpp。

五、小结

src/mongo/client内部客户端的设计可以概括为三层协作:

  1. 监控层:每个目标副本集一个ReplicaSetMonitor,streamable 协议下由TopologyManager+ServerDiscoveryMonitor(awaitable hello)+ServerPingMonitor(独立 ping 测 RTT)组成,快速学习 stepdown/选举/reconfig;
  2. 定向层RemoteCommandTargeter按读偏好 + RTT 选出节点,TargetingMetadata把过载节点降权,失败结果通过markHostNotPrimary/Unreachable/ShuttingDown反哺本地视图;
  3. 重试层RetryStrategy统一“是否重试 + 等多久”,DefaultRetryStrategy提供可重试判定与指数抖动退避(参数可由服务器参数在线调整),AdaptiveRetryStrategy用令牌桶在系统大面积过载时自动收缩重试规模,分片场景则以整个分片为共享预算单元。

这套机制让 mongos/mongod 在服务端内部复用与外部驱动一致的 SDAM 语义,同时叠加了过载感知与预算控制,是理解 MongoDB 分片路由与高可用行为的关键入口。相关源码入口:README、retry_strategy.h、replica_set_monitor_interface.h、streamable_replica_set_monitor.h、targeting_metadata.h、remote_command_targeter.h。

【免费下载链接】mongoThe MongoDB Database项目地址: https://gitcode.com/GitHub_Trending/mo/mongo

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/14 13:35:39

153、MLIR的Stream(流)编程模型与FIFO通信

MLIR的Stream(流)编程模型与FIFO通信 一个让我熬夜三天的bug 去年做AI加速器编译器的时候,遇到一个诡异的死锁问题。硬件仿真跑得好好的,一到FPGA上板就卡死,波形抓出来一看,两个加速器核之间的FIFO读指针永远停在0x3F,写指针却已经跳到0x80——溢出了。查了三天,最后…

作者头像 李华
网站建设 2026/9/14 13:34:36

手部几何识别实战:基于Matlab的图像预处理与BP神经网络实现

简介&#xff1a;基于Matlab实现的手部几何特征识别系统&#xff0c;面向生物识别、人机交互与无障碍技术开发者&#xff0c;提供从图像捕获、预处理、特征提取到匹配识别的完整算法流程&#xff0c;可作为课程设计或科研入门参考。压缩包共30个文件&#xff0c;约1.25MB&#…

作者头像 李华
网站建设 2026/9/14 13:32:38

SNMP协议栈选型:Net-SNMP与国产自研SDK的信创适配之道

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

作者头像 李华
网站建设 2026/9/14 13:31:30

合并两个有序数组:从暴力排序到双指针原地归并的保姆级教程

合并两个有序数组这道题&#xff0c;在 LeetCode 上挂着 Easy 的标签&#xff0c;但真到了面试现场&#xff0c;它能淘汰的人远比想象中多。我印象很深&#xff0c;有一次候选人把“先合并再排序”写出来&#xff0c;然后理直气壮说这就是最优解&#xff0c;我追问了一句“那如…

作者头像 李华