MongoDB 聚合 Pipeline 重写与优化:从启发式改写机制到规则注册实践
【免费下载链接】mongoThe MongoDB Database项目地址: https://gitcode.com/GitHub_Trending/mo/mongo
本文以 Pipeline 优化模块文档 为主体,深入讲解 MongoDB 聚合管道在被解析为Pipeline后如何经历两阶段启发式重写:跨阶段优化(阶段交换、合并、插入)与阶段内优化(常量折叠、无操作阶段消除)。读完本文,你将掌握optimizePipeline()的完整调用链、依赖追踪的五大依赖类别、如何基于REGISTER_RULES宏注册新的重写规则,以及如何用disablePipelineOptimization故障点验证改写语义正确性。
一、总览:Pipeline 重写的两步走
用户在aggregate命令中提交的管道,解析为Pipeline之后,会经历启发式重写(heuristic rewrites),把整条管道以及管道内的各个阶段改写为更高效的等价形式。优化入口是 optimizePipeline() 函数,它包含两个主要步骤:
- 跨阶段优化(Inter-stage Optimization):优化整个
Pipeline对象。它在内部表现为DocumentSource的容器,优化过程通过合并(combining)、交换(swapping)、删除(dropping)和插入(inserting)阶段来修改这个容器。 - 阶段内优化(Stage-specific Optimization):逐个优化每个阶段,即每个
DocumentSource。
从源码实现看,optimize.cpp 中的optimizePipeline()在确认管道未被冻结(isFrozen())且未被故障点禁用后,按顺序调用两轮规则引擎:
applyRuleBasedRewrites(rbr::PipelineRewriteContext(pipeline), Tags::Reordering); applyRuleBasedRewrites(rbr::PipelineRewriteContext(pipeline), Tags::InPlace);两轮分别只运行带Reordering标签和InPlace标签的规则,且每次重写引擎都受internalQueryMaxPipelineRewrites查询旋钮的改写次数上限保护(见 optimize.cpp#L19-L30)。
禁用优化以便验证语义
如果你新增了某个重写规则,想验证其语义正确性,可以将开启优化的结果与未优化形态的执行结果做对比。仓库提供了两个手段:
打开
disablePipelineOptimization故障点,阻止单个DocumentSource被优化。在 mongo shell 中执行:db.adminCommand({configureFailPoint: "disablePipelineOptimization", mode: "alwaysOn"})对应源码即 optimize.cpp#L41-L43 中对
MONGO_unlikely(disablePipelineOptimization.shouldFail())的短路检查。在管道中每个阶段之前插入
{$_internalInhibitOptimation: {}}阶段,确保各阶段不参与整条Pipeline级的优化(如阶段下推或交换)。
二、跨阶段优化:交换、合并与插入阶段
跨阶段优化的入口是 optimizeContainer(),它调用基于规则的重写引擎执行所有可能组合或重排相邻阶段、或以其他形式修改管道结构的规则。目前除$match、$sample、$project、$redact的下推外,全部跨阶段重写都实现在各DocumentSource的公开optimizeAt()方法中,并注册为无条件规则(unconditional rules)。例如 qo_rules_to_move.cpp 中可以看到为DocumentSourceSkip、DocumentSourceLimit、DocumentSourceGroup、DocumentSourceUnionWith、DocumentSourceUnwind、DocumentSourceSort等逐一注册OPTIMIZE_AT_RULE与OPTIMIZE_IN_PLACE_RULE。
跨阶段优化大致分为三类:
2.1 交换阶段(Swapping stages)
典型例子是$match 下推。一般来说,希望在文档进入计算更重的阶段之前先过滤,以最小化工作量。如果用户把$match写在$sort之后,优化器会尽可能将其下推,以最小化需要排序的文档数量。例如用户写出:
{ $sort: { age : -1 } }, { $match: { status: 'A' } }优化器会将其改写为:
{ $match: { status: 'A' } }, { $sort: { age : -1 } }这一能力的落地是 match_rules.cpp 中注册的MATCH_PUSHDOWN规则(match_rules.cpp#L490-L499):前置条件为matchCanSwapWithPrecedingStage,变换函数为pushMatchBeforePrecedingStage,优先级为kDefaultPushdownPriority(100.0,所有规则中的最高档),标签为Reordering。值得注意的是,$match 并非总能整段交换:pushdownMatch()会先调用DocumentSourceMatch::splitMatchByModifiedFields()把谓词按“可被前级阶段改名/修改的字段”拆成两部分(match_rules.cpp#L340-L368),再通过Transforms::partialPushdown()把可下推部分移到前级阶段之前、不可下推部分保留在原地。另外源码中还能看到一些精细的守卫逻辑,例如groupMatchSwapVerified()会阻止与$group桶化语义冲突的谓词($exists、$type、部分$expr)下推(match_rules.cpp#L178-L250),以及文本搜索谓词的$match本就要求位于管道首位因而无需下推。
2.2 合并阶段(Coalescing stages)
在一系列重排优化之后,优化器尽可能把某阶段合并进其前驱阶段。例如当$sort位于$limit之前、且中间没有会改变文档数量的阶段(如$unwind或$group)时,优化器可以把$limit吸收进$sort。给定:
{ $sort : { age : -1 } }, { $project : { age : 1, status : 1, name : 1 } }, { $limit: 5 }优化器会改写为:
{ "$sort" : { "sortKey" : { "age" : -1 }, "limit" : NumberLong(5) } }, { "$project" : { "age" : 1, "status" : 1, "name" : 1 } }这样排序阶段在推进过程中只需保留前 N 条结果(N 为 limit 值),减少了需要驻留内存的文档数量。
两个连续的同名阶段也可以合并(等价于删除其中一个)。例如相邻的两个$limit可以合并为其中较小的一个:
{ $limit: 100 }, { $limit: 10 }改写为:
{ $limit: 10 }2.3 插入阶段(Inserting stages)
虽然看似反直觉,但向管道中插入阶段有时是有益的:它可以约束传递给下游阶段的文档流,并利用更优的索引。例如当一个$redact后面紧跟$match时,可以把$match中的一部分“上插”到$redact之前。给定:
{ $redact: { $cond: { if: { $gte: [ "$sensitivity", 3 ] }, then: "$$PRUNE", else: "$$DESCEND" } } }, { $match: { "status": "active", "sensitivity": { $lt: 5 } } }优化器可以推断出status字段与$redact阶段无关,而sensitivity字段与之相关,于是把$match拆分为独立部分与依赖部分——前者推到$redact之前,后者保留在之后:
{ $match: { "status": "active" } }, { $redact: { $cond: { if: { $gte: [ "$sensitivity", 3 ] }, then: "$$PRUNE", else: "$$DESCEND" } } }, { $match: "sensitivity": { $lt: 5 } }改写后的管道阶段数变多了,但执行更优,因为它同时(1)减少了进入资源密集$redact阶段的文档量,(2)让第一个$match可以利用status上的索引。
除了交换与合并,仓库中还存在“删除冗余阶段”这类重写。例如 sort_rules.cpp 注册的REDUNDANT_SORT_REMOVAL规则:当$sort的排序模式已被上游某阶段建立的排序模式所扩展(isExtensionOf)、且中间各阶段满足preservesOrderAndMetadata约束、不会改名或重算排序键字段时,当前$sort被Transforms::eraseCurrent直接删除。前置条件sortIsRedundantGivenPrecedingStages()还会特意排除已吸收$limit的$sort与时序集_timeSorter等带有额外执行语义的形态(sort_rules.cpp#L67-L112)。
2.4 依赖追踪与分析
为了判断管道是否合法、阶段能否互相下推、或某个阶段的部分谓词能否裁剪,优化器依赖依赖追踪与分析(dependency tracking and analysis):识别每个阶段执行所依赖的字段或变量。若某阶段依赖更早的阶段,它在整个重写过程中必须保持在该阶段之后;反之,若某阶段与前一阶段独立且下推有利,则可以安全地下推。主要依赖类别有:
- 文档字段依赖(Document field dependencies):阶段所需的特定文档字段。例如
{$project: {name: 1}}依赖name字段。 - 计算字段(Computed fields):由
$addFields、$group等阶段生成的新字段,可能成为后续阶段的依赖。例如{$addFields: {total: {$sum: ["price", "$tax"]}}}计算出total,可被后续阶段引用。 - 改名(Renames):改变字段名的映射关系,后续阶段必须感知这些变换才能解析正确字段。例如
{$project: {newField: "$oldField"}}把oldField改名为newField。 - 变量引用(Variable references):对作用域变量的依赖,如用户自定义变量或系统变量(
$$CURRENT、$$ROOT)。例如:
{ "$project": { "adjustedValue": { "$let": { "vars": { "discount": 0.1 }, "in": { "$multiply": ["$price", { "$subtract": [1, "$$discount"] }] } } } } }其中$$discount就是在$let内定义并使用的变量引用。
- 元数据(Metadata):不属于文档核心字段的附加信息,常出现在依赖上下文信息的操作中,如文本搜索得分、地理邻近度或保留字段。例如
{$project: {score: {$meta: "textScore"}}}依赖的是textScore元数据而非文档字段。
依赖追踪与校验的实现细节可参阅 expression_algo.h、semantic_analysis.h 和 dependencies.h。从源码结构看,重写上下文还维护了DependencyGraph(阶段间顺序依赖图),见 rule_based_rewriter.h 中的getDependencyGraph()与 rule_based_rewriter.cpp 中postTransform()的图失效逻辑:每次变换之后,上下文会保守地把从上一位置起的依赖图标记为失效,从而保证后续规则看到的依赖信息是最新的。
三、阶段内优化:单阶段重写与常量折叠
确定阶段最终顺序后,系统会再次调用规则引擎,但这次使用只运行阶段内优化规则的配置。这些规则目前实现为DocumentSource子类的公开optimize()方法,并注册为无条件规则。每条规则要么返回一个语义等价的优化后DocumentSource,要么在当前阶段是 no-op 时将其删除。例如 no-op 阶段{$match: {}}会被直接移除。
$match阶段中的MatchExpression包含专门的重写逻辑,详见 MatchExpression 说明。
此外,含有ExpressionConstant值的Expression的阶段可能符合常量折叠条件。例如:
{ $project: { a: { $sum: [ 4, 5, 1 ] } } }其中的常量可以折叠为一个:
{ $project: { a: { $literal: 10 } } }常量折叠(Constant folding):对包含常量或可解析为常量的表达式求值,并用计算结果替换原表达式。在查询规划期简化表达式,从而降低执行期的计算开销。
表达式(Expression):查询中解析为某个值的组件。它是无状态的,即返回一个值而不改变用于构建表达式的任何值。例如表达式
{$add: [3, "$inventory.total"]}由$add运算符与两个输入表达式构成:常量3和字段路径表达式"$inventory.total",它返回输入文档在路径inventory.total处取值加 3 的结果。
这些优化看起来显然,但聚合管道往往由计算机生成,应用层通常不会(也不应该)做这类分析。而且原始查询可能更复杂,经过前面的启发式重写后,可能发现比最初预期更多的值可以合并折叠。
整个优化流程如下:
四、注册新的重写规则
所有管道重写都经由基于规则的重写引擎触发。虽然多数重写目前仍实现在DocumentSource::optimizeAt()与optimize()中并注册为无条件规则(即前置条件恒为真),但新的重写应当实现并注册为独立的规则。一条规则由名称、前置条件与变换函数、优先级(priority)以及一组标签(tags)定义(见 Rule 结构体):
precondition:决定是否执行transform;transform:应用规则本身,返回值指示引擎是否需要重排/重放当前位置;priority:数值越大优先级越高,同位置多条规则命中时按此排序执行;tags:允许引擎只运行某一子集规则。
引擎的核心循环在 RewriteEngine::applyRules():对每个元素先向上下文收集可应用规则,然后按优先级尝试,根据变换是否改变了当前位置来决定重放(Requeue)还是前进(Advance)。
4.1 规则注册表与注册宏
规则注册表(Rule registry) 是DocumentSource子类型与可应用规则之间的映射,它作为ServiceContext的 decoration 存在(rule_based_rewriter.cpp#L73-L84)。这意味着规则注册在创建服务上下文时(即新 mongod/mongos 进程启动时)被调用。每当重写引擎推进到新元素,就会调用 PipelineRewriteContext::enqueueRules(),其中按当前DocumentSource的类型查表,并检查规则关联的 feature flag 是否启用(未启用时以kLastLTS作为回退 FCV 判断),把适用规则入队。
规则可通过REGISTER_RULES宏注册:第一个参数是DocumentSource子类,随后是逗号分隔的规则列表。以DocumentSourceMatch的注册为例(match_rules.cpp#L490-L499):
REGISTER_RULES(DocumentSourceMatch, OPTIMIZE_AT_RULE(DocumentSourceMatch), OPTIMIZE_IN_PLACE_RULE(DocumentSourceMatch), { .name = "MATCH_PUSHDOWN", .precondition = matchCanSwapWithPrecedingStage, .transform = pushMatchBeforePrecedingStage, .priority = kDefaultPushdownPriority, .tags = PipelineRewriteContext::Tags::Reordering, });宏内部展开为ServiceContext::ConstructorActionRegisterer静态注册器(rule_based_rewriter.h#L48-L70),OPTIMIZE_AT_RULE(DS)与OPTIMIZE_IN_PLACE_RULE(DS)两个辅助宏分别把DocumentSource的optimizeAt()与optimize()包装成无条件规则。仓库为不同用途的规则约定了默认优先级(rule_based_rewriter.h#L96-L103):
| 常量 | 值 | 用途 |
|---|---|---|
kDefaultPushdownPriority | 100.0 | 尽量早地推入$match等高优先级下推 |
kDefaultHoistPriority | 50.0 | 条件性地提升计算以促成更多$match下推 |
kDefaultOptimizeAtPriority | 10.0 | 与相邻阶段交换或吸收 |
kDefaultOptimizeInPlacePriority | 1.0 | 原地优化阶段内部 |
若需让规则受 feature flag 门控,使用REGISTER_RULES_WITH_FEATURE_FLAG宏,用法与REGISTER_RULES类似,但第二个参数为 feature flag。例如 match_rules.cpp#L513-L521 中PUSH_MATCH_BEFORE_SINGLE_DOC_TRANSFORMATION规则就挂在gFeatureFlagImprovedDepsAnalysis之下。
另一种让规则被条件调用的方式,是在另一条规则的前置条件或变换函数中调用RewriteContext::addRule()动态入队。需要注意:如果引擎被配置为只运行某一组规则,动态入队的规则只有在同属该组时才会执行。例如规则PUSH_MATCH_BEFORE_CHANGE_STREAMS就是在canPushMatchBefore()检测到前级阶段是 mongos 内部 change stream 阶段时通过ctx.addRule()入队的(match_rules.cpp#L291-L297)。
4.2 标签(Tags)
当前部分重写依赖一个假设:所有跨阶段优化都在任何原地优化之前对整条管道完成。若此假设被破坏,重写之间可能互相干扰。因此现有管道重写被划分为两类标签:
Reordering:可能改变当前阶段之外其他阶段的规则(如重排、合并、删除阶段);InPlace:只优化阶段内部、从不触碰相邻阶段的规则;- 另有
Testing标签的规则,仅在enablePipelineOptimizationAdditionalTestingRules查询旋钮开启时应用(见 optimize.cpp#L52-L55)。
两组成分是分别对管道执行的。编写新规则时的选择原则很简单:如果你的规则可能改变当前阶段以外的任何阶段,就给它Reordering标签;否则给它InPlace标签。
五、小结与延伸阅读
Pipeline 优化机制可以概括为:optimizePipeline()先用Reordering轮驱动跨阶段的结构改写(下推、合并、插入、冗余删除),再用InPlace轮驱动各阶段内部的等价重写(无操作消除、常量折叠等),而全部规则通过ServiceContext级别的注册表 + feature flag 体系集中管理,支持按标签分组、按优先级调度、按条件动态入队。
希望继续深入时,建议按以下路径阅读:
- 优化入口与两轮驱动:optimize.cpp
- 管道重写上下文、标签与
Transforms原语:rule_based_rewriter.h、rule_based_rewriter.cpp - 通用规则引擎设计:query/compiler/rewrites/README.md、rule_based_rewriter.h
$match下推与拆分:match_rules.cpp;冗余$sort删除:sort_rules.cpp;批量optimizeAt注册:qo_rules_to_move.cpp- 端到端测试:pipeline_rewriter_test.cpp
- 各阶段
MatchExpression重写细节:matcher/README.md
【免费下载链接】mongoThe MongoDB Database项目地址: https://gitcode.com/GitHub_Trending/mo/mongo
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考