MongoDB 分片集群 DISTINCT_SCAN 多 chunk 场景全解析:孤儿文档下的去重扫描与执行计划
【免费下载链接】mongoThe MongoDB Database项目地址: https://gitcode.com/GitHub_Trending/mo/mongo
本篇技术指南以 MongoDB 官方仓库(jstests/query_golden_sharding/expected_output/sbeRestricted/distinct_scan_multi_chunk.md)的 golden 测试输出为分析蓝本,系统拆解distinct命令与$group聚合在"每个分片持有多个(含孤儿)chunk"的分布下如何生成执行计划、如何在各分片间去重合并。读完本文,你将掌握DISTINCT_SCAN、SHARD_MERGE、SINGLE_SHARD、SHARDING_FILTER、$groupByDistinctScan等关键 stage 的触发条件与组合规律,并学会如何复现、验证这些场景。
一、测试背景:为什么要在"多 chunk + 孤儿文档"下验证 DISTINCT_SCAN
DISTINCT_SCAN是 MongoDB 查询引擎中一种特殊的索引扫描方式:它利用索引键有序的特性,跳过值重复的相邻条目,从而在不做内存内去重的情况下直接产出去重后的结果。它通常被distinct命令和部分$group聚合(如_id去重、$first/$last累加器)复用。
在分片集群中,去重语义被进一步复杂化:
- 同一逻辑文档可能因为 chunk 迁移、split 等原因在多个分片上残留孤儿文档(orphaned docs),这些文档不属于任何当前 chunk 的归属范围,但物理上仍存在于分片数据文件中;
- 因此,每个分片内去重还不够,还需要按分片间进行全局去重(
SHARD_MERGE); - 若谓词能通过 shard key 精确锁定唯一分片,则可退化为单分片执行(
SINGLE_SHARD),省去跨分片合并。
本测试的驱动文件 jstests/query_golden_sharding/distinct_scan_multi_chunk.js 在其文件头注释中明确说明了测试目的:
Tests the results and explain for a DISTINCT_SCAN in case of a sharded collection where each shard contains multiple (orphaned) chunks.
该测试同时带有三个运行前置标记(tags):featureFlagGetExecutorDeferredEngineChoice、featureFlagShardFilteringDistinctScan与requires_fcv_82,说明它验证的是基于分片过滤的DISTINCT_SCAN新特性(featureFlagShardFilteringDistinctScan),并需要 FCV 8.2 环境支持。本文分析的sbeRestricted输出目录对应 SBE 受限模式下的执行计划形态。
二、测试环境与数据构造:每个分片三个 chunk、各带孤儿文档
分片环境由工具函数setupShardedCollectionWithOrphans搭建,实现在 jstests/libs/query/golden_sharding_utils.js:
- 启动一个2 分片的
ShardingTest; - 集合
test.distinct_scan_multi_chunk以{shardKey: 1}作为分片键; - 初始为每个分片规划 3 个 chunk:
shard0持有chunk1_s0、chunk2_s0、chunk3_s0,shard1持有chunk1_s1、chunk2_s1、chunk3_s1; - 每个 chunk 写入 3 个不同的
shardKey值(${chunk}_0、${chunk}_1、${chunk}_2),且每个shardKey对应 3 条文档,其notShardKey分别以1notShardKey_、2notShardKey_、3notShardKey_为前缀——这保证了shardKey 值唯一、notShardKey 值不唯一且分组键种类丰富的数据形态; - 数据写入完成后执行
shardCollection,再按 chunk 边界split,并将shard1的三个 chunkmoveChunk迁移到非主分片; - 最后故意向两个分片插入孤儿文档:在
shard0直接插入属于shard1chunk 的chunk*_s1_*_orphan文档,在shard1直接插入属于shard0chunk 的chunk*_s0_*_orphan文档(测试顶部设置了TestData.skipCheckOrphans = true,以便绕过孤儿检查)。
集合初始拥有索引shardKey_1与shardKey_1_notShardKey_1;随着测试推进,还会动态创建/删除notShardKey_1、notShardKey_1_shardKey_1等索引,文档中每个用例都会打印Total indexes on the collection,便于核对查询规划器实际可用的索引集合。
结果与 explain 的输出由 jstests/libs/query/golden_test_utils.js 中的outputDistinctPlanAndResults(对distinct命令)与outputAggregationPlanAndResults(对聚合管道)统一生成,最终落盘为 Markdown golden 文件,供回归比对。
三、对 shard key 执行 distinct:五种谓词的五种执行形态
测试文件第 1 节对"shardKey"依次执行了五种distinct,覆盖无谓词、等值、非 shard key 等值、shard key 范围、非 shard key 范围五种组合。这是理解"谓词如何决定分片间执行方式"的核心对照实验。
3.1 无谓词:DISTINCT_SCAN + SHARD_MERGE
distinct("shardKey", {})返回 18 个值(3 chunk × 3 值 × 2 分片,无重复),执行计划为:
"mergerPart": [ { "stage": "SHARD_MERGE" } ], "shardsPart": { "distinct_scan_multi_chunk-rs0": { "rejectedPlans": [], "winningPlan": [ { "stage": "PROJECTION_COVERED", "transformBy": { "_id": 0, "shardKey": 1 } }, { "direction": "forward", "indexBounds": { "shardKey": [ "[MinKey, MaxKey]" ] }, "indexName": "shardKey_1", "isFetching": false, "isShardFiltering": true, "stage": "DISTINCT_SCAN" } ] } }要点解读:
- mongos 侧的合并 stage 为
SHARD_MERGE:各分片各自产出有序的去重结果,路由节点负责把跨分片可能重复的值(本例中正好无重复)合并去重; - 分片侧使用
shardKey_1索引上的DISTINCT_SCAN,indexBounds为[MinKey, MaxKey]全区间; isShardFiltering: true是本文档中最值得关注的字段:它表示该DISTINCT_SCAN内部集成了分片过滤(shard filtering)能力,即扫描索引条目时会剔除不属于本分片所属 chunk 的孤儿条目,因此不再需要额外的SHARDING_FILTERstage,也不需要FETCH(isFetching: false,索引已覆盖所需字段);PROJECTION_COVERED说明结果可直接由索引键投影得出,全程无回表。
3.2 shard key 等值谓词:SINGLE_SHARD 单分片直达
distinct("shardKey", {shardKey: {$eq: "chunk1_s0_1"}})的合并 stage 变为SINGLE_SHARD:因为等值谓词命中的 shard key 唯一属于某个 chunk,mongos 可直接将请求只发给distinct_scan_multi_chunk-rs0一个分片,结果仅[ "chunk1_s0_1" ]。
分片侧最终胜出的计划与 3.1 相同(shardKey_1+DISTINCT_SCAN,区间收窄为["chunk1_s0_1", "chunk1_s0_1"],isShardFiltering: false——因为精确等值下不需要再过滤孤儿)。值得注意的是其rejectedPlans中被淘汰的备选计划:查询规划器曾考虑过复合索引shardKey_1_notShardKey_1(索引边界为notShardKey: [MinKey, MaxKey]与shardKey: ["chunk1_s0_1", "chunk1_s0_1"]),最终因单字段索引shardKey_1已足够覆盖而落选。这直观展示了规划器在不同索引间的代价选择。
3.3 非 shard key 等值谓词:被迫退化为 COLLSCAN + SHARDING_FILTER
distinct("shardKey", {notShardKey: {$eq: "1notShardKey_chunk1_s0_1"}})是最有代表性的"负面"用例:谓词打在非 shard key 上,无法裁剪目标分片,于是两个分片都被访问,合并 stage 恢复为SHARD_MERGE;而在每个分片上,由于过滤字段与 distinct 字段不在同一个可用索引前缀上,规划器最终选择了全表扫描:
"winningPlan": [ { "stage": "SHARDING_FILTER" }, { "direction": "forward", "filter": { "notShardKey": { "$eq": "1notShardKey_chunk1_s0_1" } }, "nss": "test.distinct_scan_multi_chunk", "stage": "COLLSCAN" } ]此时SHARDING_FILTER作为独立 stage 挂在COLLSCAN之上,显式剔除孤儿文档。结果值[ "chunk1_s0_1" ]虽然正确,但代价是扫描全部数据。这说明:想让 DISTINCT_SCAN 在分片场景发挥优势,过滤谓词最好能落在 shard key 上,或与 distinct 字段构成索引前缀。
3.4 shard key 范围谓词:区间裁剪 + DISTINCT_SCAN
distinct("shardKey", {shardKey: {$gte: "chunk1_s0_1"}})与 3.1 形态一致(SHARD_MERGE+shardKey_1上的DISTINCT_SCAN),区别仅在于索引边界收窄为["chunk1_s0_1", {}),结果去掉了小于边界的 5 个值。被淘汰的备选计划同样包含基于shardKey_1_notShardKey_1的DISTINCT_SCAN方案,胜出者仍是覆盖最精简的shardKey_1。范围谓词同样具备 shard 裁剪能力,且isShardFiltering保持true。
3.5 非 shard key 范围谓词:与 3.3 相同的退化路径
distinct("shardKey", {notShardKey: {$gte: "1notShardKey_chunk1_s0_1"}})在两个分片上均采用SHARDING_FILTER+COLLSCAN,结果集与无谓词时完全一致(因为该谓词几乎覆盖全部notShardKey值域),说明范围谓词打在非 shard key 上同样无法获得索引加速。
小结:shard key 上 distinct 的执行规律
| 谓词 | 合并 stage | 分片侧执行 | isShardFiltering |
|---|---|---|---|
| 无 | SHARD_MERGE | DISTINCT_SCAN(shardKey_1) | true |
| shardKey $eq | SINGLE_SHARD | DISTINCT_SCAN(shardKey_1) | false |
| notShardKey $eq | SHARD_MERGE | COLLSCAN + SHARDING_FILTER | - |
| shardKey $gte | SHARD_MERGE | DISTINCT_SCAN(shardKey_1) | true |
| notShardKey $gte | SHARD_MERGE | COLLSCAN + SHARDING_FILTER | - |
四、对非 shard key 字段执行 distinct:isFetching 的代价
测试第 2 节在coll.createIndex({notShardKey: 1})之后,对"notShardKey"重复了与第 3 节相同的五种谓词实验。由于 distinct 目标字段本身是notShardKey_1索引的键,DISTINCT_SCAN仍然成立,但isFetching变为true:
{ "direction": "forward", "indexBounds": { "notShardKey": [ "[MinKey, MaxKey]" ] }, "indexName": "notShardKey_1", "isFetching": true, "isShardFiltering": true, "stage": "DISTINCT_SCAN" }原因在于分片过滤需要比对文档的 shard key 归属,而notShardKey_1单键索引不含 shard key 字段,因此DISTINCT_SCAN必须回表取文档才能完成孤儿判定(isFetching: true)。这与第 3 节中shardKey_1索引(键即分片键、天然携带过滤信息、isFetching: false)形成鲜明对照,是理解"分片过滤能否被索引覆盖"的关键实例。
四种组合的规律如下:
- 无谓词:两个分片各自
DISTINCT_SCAN(notShardKey_1),合并SHARD_MERGE,输出 54 个值(3 前缀 × 3 chunk × 3 值 × 2 分片); - shardKey 等值:合并 stage 变为
SINGLE_SHARD,且分片侧胜出计划使用复合索引shardKey_1_notShardKey_1上的DISTINCT_SCAN(边界notShardKey: [MinKey, MaxKey]+shardKey: ["chunk1_s0_1", "chunk1_s0_1"],isFetching: false),因为该索引前缀恰好覆盖谓词与 distinct 字段,无需回表;被淘汰的备选计划是shardKey_1上的IXSCAN+FETCH; - notShardKey 等值:
notShardKey_1上DISTINCT_SCAN,边界精确为["1notShardKey_chunk1_s0_1", "1notShardKey_chunk1_s0_1"],isFetching: true; - shardKey 范围:胜出计划使用
shardKey_1_notShardKey_1复合索引(isFetching: false),备选方案为shardKey_1索引IXSCAN+SHARDING_FILTER+FETCH; - notShardKey 范围:
notShardKey_1上DISTINCT_SCAN,边界["1notShardKey_chunk1_s0_1", {}),isFetching: true。
由此可提炼出分片distinct的两条经验法则:
- distinct 字段 + 过滤字段若能落在同一索引前缀(如复合索引
shardKey_1_notShardKey_1),就能得到"覆盖式 DISTINCT_SCAN",isFetching: false; - shard key 等值谓词可触发 SINGLE_SHARD,将全局去重降级为单分片执行,是最经济的形态。
五、$group 对非 shard key 字段分组:$groupByDistinctScan 优化
测试第 3 节验证 MongoDB 将{$group: {_id: "$notShardKey"}}这类"按字段去重"的聚合优化为基于DISTINCT_SCAN的分片内预聚合。其 explain 呈现出一种分片特有的三段式结构(见 distinct_scan_multi_chunk.md):
"distinct_scan_multi_chunk-rs0": [ { "$cursor": { "winningPlan": [ { "indexName": "notShardKey_1", "isFetching": true, "isShardFiltering": true, "stage": "DISTINCT_SCAN", "indexBounds": { "notShardKey": [ "[MinKey, MaxKey]" ] } } ] } }, { "$groupByDistinctScan": { "newRoot": { "_id": "$notShardKey" } } } ], "mergeType": "router", "mergerPart": [ { "$mergeCursors": { "nss": "test.distinct_scan_multi_chunk", "allowPartialResults": false, "tailableMode": "normal" } }, { "$group": { "$doingMerge": true, "_id": "$$ROOT._id" } } ], "shardsPart": [ { "$group": { "_id": "$notShardKey" } } ]结构解读:
shardsPart是逻辑上发给各分片的聚合片段:$group按notShardKey分组;- 分片内实际执行被拆成
$cursor(内层DISTINCT_SCAN)与$groupByDistinctScan(把 distinct 扫描的键重新构造成_id字段),从而让分片侧先做一轮"索引级去重"; mergerPart在 mongos 上执行:$mergeCursors汇聚各分片游标后,再做一次$doingMerge: true的$group完成跨分片去重;mergeType: "router"表示合并发生在路由节点。
5.1 带 $first/$last 累加器的分组
当$group附加$first/$last累加器时(对应测试文件 distinct_scan_multi_chunk.js 中的四个用例),$groupByDistinctScan的newRoot会把累加器引用的字段一并带入:
{ "$groupByDistinctScan": { "newRoot": { "_id": "$notShardKey", "accum": "$notShardKey" } } }而 mongos 侧的 merge$group会相应地对accum执行$first/$last。两个关键观察:
$first对应 forward 扫描,$last对应 backward 扫描:explain 中$first用例的DISTINCT_SCAN方向为"forward"(索引边界[MinKey, MaxKey]),$last用例的方向为"backward"(索引边界写为[MaxKey, MinKey]),这正对应"取分组内第一/最后一个元素"的语义;$first: "$shardKey"与$first: "$notShardKey"结果值不同但计划相同:由于分组键本身唯一(每个notShardKey值只对应一个shardKey),两种累加器输出的accum值恰好一致,但索引选择仍是notShardKey_1。
5.2 前置 $match:能否下推索引决定成败
$match打在 shard key 上({$match: {shardKey: {$gte: "chunk1_s0_1"}}}):分片侧最终胜出计划为GROUP→PROJECTION_COVERED(_id: false, notShardKey: true)→SHARDING_FILTER→IXSCAN使用复合索引shardKey_1_notShardKey_1(边界notShardKey: [MinKey, MaxKey]+shardKey: ["chunk1_s0_1", {}))。注意此时走的是IXSCAN而非DISTINCT_SCAN,因为前置$match引入了额外的谓词求值,去重被上移到GROUPstage 完成;被淘汰的备选计划是基于shardKey_1索引的IXSCAN+SHARDING_FILTER+FETCH;$match打在 notShardKey 上({$match: {notShardKey: {$gte: "1notShardKey_chunk1_s0_1"}}}):则保持DISTINCT_SCAN(notShardKey_1,边界["1notShardKey_chunk1_s0_1", {}),isFetching: true)+$groupByDistinctScan的形态,$match被吸收进索引边界。
5.3 前置 $sort:SORT_KEY_GENERATOR 介入
当管道为{$sort: {notShardKey: 1, shardKey: 1}}+$group且临时创建了notShardKey_1_shardKey_1索引时(distinct_scan_multi_chunk.js),分片侧胜出计划变为:
[ { "stage": "PROJECTION_COVERED", "transformBy": { "_id": false, "notShardKey": true } }, { "stage": "SORT_KEY_GENERATOR" }, { "stage": "SHARDING_FILTER" }, { "indexName": "notShardKey_1_shardKey_1", "stage": "IXSCAN", "indexBounds": { "notShardKey": [ "[MinKey, MaxKey]" ], "shardKey": [ "[MinKey, MaxKey]" ] } } ]mongos 侧的$mergeCursors携带了sort: {notShardKey: 1, shardKey: 1},说明排序信息通过合并游标传递;由于各分片按全局排序键输出,路由侧无需再做$doingMerge的$group去重($willBeMerged: false),可以直接按顺序聚合。这组用例展示了$sort + $group + $first/$last场景下优化器如何利用复合索引一次完成"排序 + 分组取首尾"。
六、$group 对 shard key 分组:$top/$bottom 与 $first/$last 的对照
第 4、5 节把分组键换成 shard key 本身,重点验证$top/$bottom与$first/$last两类累加器。
6.1 $top/$bottom:按 shard key 排序取极值
$group使用$top/$bottom累加器时(distinct_scan_multi_chunk.js),累加器自带sortBy排序规范。文档覆盖了四种组合:
- sortBy 只有 shard key:
$top对应sortBy: {shardKey: 1}(正序),$bottom对应sortBy: {shardKey: -1}(逆序),且反向 sortBy 与$top/$bottom的四种排列组合都做了验证; - sortBy 为 shard key + 另一字段:如
{$top: {sortBy: {shardKey: 1, notShardKey: 1}, output: "$shardKey"}}及对应的$bottom、output: "$notShardKey"变体; - sortBy 只含非 shard key 字段:测试注释(TODO SERVER-95198)明确指出,此类用例当前退化为 COLLSCAN 计划,但在理论上可以通过
shardKey_1_notShardKey_1上的DISTINCT_SCAN回答——这是文档为我们留下的"已知优化空间"信号。
6.2 $first/$last:依赖前置 $sort 与否
第 5 节系统对照了"有/无前置$sort"对$first/$last的影响:
- 有前置
$sort(如{$sort: {shardKey: -1}}):排序信息可下推,分片按序输出,路由侧直接聚合; - 无前置
$sort:$first/$last依赖数据自然序或额外的排序代价,对应用例出现在 "without preceding $sort" 小节; - 前置 $sort + 插入 $match:如
{$sort: {shardKey: 1, notShardKey: 1}}+{$match: {shardKey: {$gt: "chunk1_s0"}}}+$group,验证排序、过滤、分组三者叠加时$first/$last(含$$ROOT输出)的行为; - 输出非 shard key 字段:
$first: "$notShardKey"、$last: "$notShardKey"作为对照组。
6.3 一个重要的"理论可行、当前未实现"案例
在 "sort by non-shard key field" 两个小节中,文档反复出现如下注释(源码见 distinct_scan_multi_chunk.js):
TODO SERVER-95198: Currently results in a COLLSCAN plan. However, in theory it could also be answered by a DISTINCT_SCAN on 'shardKey_1_notShardKey_1'.
也就是说:按非 shard key 排序再对 shard key 做$top/$bottom分组时,规划器目前选择 COLLSCAN,但理论上复合索引shardKey_1_notShardKey_1存在通过DISTINCT_SCAN回答的可能。这类 TODO 是 golden 测试文档的价值所在——它不仅记录"现状是什么",也标注了"未来可以做什么"。
七、multikey 字段上的 distinct:被迫回表
测试第 6 节(distinct_scan_multi_chunk.js)额外插入了数组字段文档:
coll.insertMany([ {_id: "mk1", notShardKey: [1, 2, 3]}, {_id: "mk2", notShardKey: [2, 3, 4]}, ]);随后对 multikey 的notShardKey字段执行distinct(无谓词、{shardKey: null}、{notShardKey: 3}三种过滤)。explain 中isFetching: true普遍存在,源码注释(TODO SERVER-97235)进一步说明:multikey 场景下当前无法避免 fetching,测试也因此特意保留 multikey 值以覆盖该行为,防止将来无意中"优化"出错误的覆盖计划。这提醒我们:multikey 索引条目在分片过滤语义下不能直接信任索引键,必须回表核对文档归属。
八、queryShapeHash 与 explain 摘要的稳定性价值
文档中每一个用例的 explain 都带有queryShapeHash字段(如3E08037AA6B94982B95D1C3F310D7433C831031FBDB14536694772A401A17D74)。它是对查询形状(query shape)的稳定哈希,用于查询分析、计划缓存与遥测归因。在 golden 测试中,这些哈希与执行计划一起被固化为基线,任何会导致查询形状或计划变化的代码改动都会使输出 diff,从而在 CI 中暴露回归。rejectedPlans数组的完整保留(而非仅记录胜出计划)则进一步放大了这种敏感度——规划器候选集的任何增减都会被记录在案。
九、如何复现与进一步研究
本文所述全部输出均由真实测试生成,你可以在本仓库中按如下方式继续研究:
- 测试驱动脚本:jstests/query_golden_sharding/distinct_scan_multi_chunk.js 定义了全部用例与执行顺序,可通过 resmoke 测试框架运行;
- Golden 基线文件:jstests/query_golden_sharding/expected_output/sbeRestricted/distinct_scan_multi_chunk.md(本文分析对象,SBE 受限模式),同目录下还有 sbeFull、sbeDisabled 与 featureFlagSbeFull 等变体,可横向对比不同执行引擎下的计划差异;
- 分片环境构造:jstests/libs/query/golden_sharding_utils.js 中的
setupShardedCollectionWithOrphans负责孤儿数据注入; - 输出格式化与断言:jstests/libs/query/golden_test_utils.js 中的
outputDistinctPlanAndResults、outputAggregationPlanAndResults解释了每个 Markdown 小节(Pipeline、Results、Total indexes on the collection、Summarized explain)是如何生成的; - 相关对比用例:jstests/query_golden_sharding/distinct_chunk_skipping.js 及其 expected_output 覆盖了 distinct 过程中的 chunk 跳过优化,可与本文的多 chunk 场景互为补充。
十、核心结论
DISTINCT_SCAN在分片场景的价值取决于索引是否携带分片过滤所需信息:以 shard key 为键的索引可实现isFetching: false的覆盖式去重扫描,而非 shard key 单键索引则必须isFetching: true回表判定孤儿归属;- shard key 等值谓词是最优路径,可将全局
SHARD_MERGE降级为SINGLE_SHARD单分片执行; - 谓词打在非 shard key 上时,
distinct与$group可能退化为SHARDING_FILTER+COLLSCAN,或依赖复合索引前缀(如shardKey_1_notShardKey_1)挽救覆盖能力; $group去重可被优化为$groupByDistinctScan,配合$first/$last时自动选择 forward/backward 扫描方向;$top/$bottom场景则按 sortBy 推导方向,并存在(SERVER-95198 跟踪的)DISTINCT_SCAN优化空间;- golden 测试以"结果 + 完整 explain(含 rejectedPlans 与 queryShapeHash)"固化了这些行为,是排查分片查询性能与规划回归的第一手资料。
【免费下载链接】mongoThe MongoDB Database项目地址: https://gitcode.com/GitHub_Trending/mo/mongo
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考