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的成员。节点的“资格”取决于两个维度:
- 节点类型:primary、secondary 等(来自监控学到的拓扑状态);
- 平均网络往返时延(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 hello:isMaster/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_rpc、AsyncTry)。此外还提供一个独立函数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:
| 服务器参数 | 默认值 | 可设置时机 | 说明 |
|---|---|---|---|
defaultClientBaseBackoffMillis | 100 | startup / runtime | 指数退避的基数;若大于defaultClientMaxBackoffMillis则改用后者 |
defaultClientMaxBackoffMillis | 10000 | startup / runtime | 单次退避的时长上限 |
defaultClientMaxRetryAttempts | 3 | startup / runtime | 遇到可重试错误时的最大重试次数 |
另外,当服务器在过载错误响应中携带了不同的 base backoff 时,会通过Result的baseBackoffMS传入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内部客户端的设计可以概括为三层协作:
- 监控层:每个目标副本集一个
ReplicaSetMonitor,streamable 协议下由TopologyManager+ServerDiscoveryMonitor(awaitable hello)+ServerPingMonitor(独立 ping 测 RTT)组成,快速学习 stepdown/选举/reconfig; - 定向层:
RemoteCommandTargeter按读偏好 + RTT 选出节点,TargetingMetadata把过载节点降权,失败结果通过markHostNotPrimary/Unreachable/ShuttingDown反哺本地视图; - 重试层:
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),仅供参考