MongoDB 系列 · 第八篇 · 完结篇

分片策略:一条查询是怎么知道该去问哪个 shard 的

这是 MongoDB 系列的最后一篇。前面七篇都在单个副本集内打转——这一篇往上再加一层:数据量大到一个副本集放不下时,MongoDB 怎么把它拆到多个"分片"(shard)上,而 mongos 这个路由层,又是怎么在不知道数据具体在哪的情况下,决定一条查询该往哪几个分片发。

本机真实搭了一个最小可用的分片集群:1 个配置服务器、2 个分片、1 个 mongos,手动把一个 collection 的数据真实分布到两个分片上,再用真实的 explain() 看不同写法的查询到底问了哪几个分片。

4 个真实进程组成的集群
1 配置服务器 + 2 分片 + 1 mongos,本机真实跑起来,真实做过一次 chunk 迁移
3 处真实源码
mongos 判断查询该发给几个分片的核心函数,逐行核对过路由逻辑

真实实验:搭一个最小分片集群

配置服务器和每个分片本身也是副本集(这里都用单节点简化),mongos 只连配置服务器,不直接管理数据。

真实实测
$ mongod --configsvr --replSet cfgrs --port 27101 --dbpath ./cfg &
$ mongod --shardsvr  --replSet shard1rs --port 27102 --dbpath ./shard1 &
$ mongod --shardsvr  --replSet shard2rs --port 27103 --dbpath ./shard2 &
$ mongosh --port 27101 --eval 'rs.initiate({_id:"cfgrs", configsvr:true, members:[{_id:0,host:"127.0.0.1:27101"}]})'
$ mongosh --port 27102 --eval 'rs.initiate({_id:"shard1rs", members:[{_id:0,host:"127.0.0.1:27102"}]})'
$ mongosh --port 27103 --eval 'rs.initiate({_id:"shard2rs", members:[{_id:0,host:"127.0.0.1:27103"}]})'

$ mongos --configdb cfgrs/127.0.0.1:27101 --port 27100 &
$ mongosh --port 27100 --eval '
  sh.addShard("shard1rs/127.0.0.1:27102");
  sh.addShard("shard2rs/127.0.0.1:27103");
'
[ { _id: 'shard1rs', host: 'shard1rs/127.0.0.1:27102' }, { _id: 'shard2rs', host: 'shard2rs/127.0.0.1:27103' } ]

这跟前两篇搭多节点副本集的手法完全一样——本机随便起几个 mongod/mongos 进程,不需要容器、不需要网络,组一个真实的分片集群。mongos 全程不存一份用户数据,它手里只有"分片元数据"(哪个 collection、哪段 key 范围、在哪个分片上),这份元数据来自配置服务器。

真实实验:把一个 collection 真实拆到两个分片上

shopdb.orders{region:1, _id:1} 分片,插入 400 篇文档(4 个地区各 100 篇),手动把一半地区迁到 shard2。

真实实测
$ mongosh --port 27100 --eval '
  sh.enableSharding("shopdb");
  db.getSiblingDB("shopdb").orders.createIndex({region:1, _id:1});
  sh.shardCollection("shopdb.orders", {region:1, _id:1});
  sh.splitAt("shopdb.orders", {region:"south", _id:MinKey});
  sh.moveChunk("shopdb.orders", {region:"south", _id:0}, "shard2rs");
'
{ millis: 304, ok: 1, ... } // chunk 迁移真实完成,耗时 304ms # 绕过 mongos,直接连每个分片确认数据真的分开了: $ mongosh --port 27102 shopdb --eval 'db.orders.aggregate([{$group:{_id:"$region",n:{$sum:1}}}])' // shard1 [ { _id: 'east', n: 100 }, { _id: 'north', n: 100 } ] $ mongosh --port 27103 shopdb --eval 'db.orders.aggregate([{$group:{_id:"$region",n:{$sum:1}}}])' // shard2 [ { _id: 'south', n: 100 }, { _id: 'west', n: 100 } ]

分片键 {region:1, _id:1}region 字母序切了一刀:region < "south"(eastnorth)留在 shard1,region >= "south"(southwest)真实搬到了 shard2——这不是配置,是磁盘上真实发生的数据搬迁。

真实源码:mongos 怎么决定发给哪个分片

核心是一个"快速路径 + 兜底路径"的两段逻辑,都在同一个函数里。

src/mongo/s/query/shard_key_pattern_query_util.cpp · mongodb/mongo @ r8.3.7, L535 // Fast path for targeting equalities on the shard key. auto shardKeyToFind = extractShardKeyFromQuery(cm.getShardKeyPattern(), query); if (!shardKeyToFind.isEmpty()) { try { auto chunk = cm.findIntersectingChunk(shardKeyToFind, commandCollation, ...); shardIds->insert(chunk.getShardId()); return; } catch (const DBException&) { // The query uses multiple shards } }
src/mongo/s/query/shard_key_pattern_query_util.cpp · mongodb/mongo @ r8.3.7, L235 BSONObj extractShardKeyFromQuery(const ShardKeyPattern& shardKeyPattern, const CanonicalQuery& query) { // We only care about extracting the full key pattern paths - if they don't exist // (or are conflicting), we don't contain the shard key. ... }

只有查询条件里对分片键每一个字段都给出精确相等值,才走"快速路径":直接算出这个分片键值落在哪个 chunk 范围里,只问那一个分片。少了任意一个字段(哪怕只差 _id),这条快速路径就直接不走了——extractShardKeyFromQuery 返回空,退回下面的兜底逻辑。

真实源码:没有完整分片键时,兜底逻辑怎么算范围

src/mongo/s/query/shard_key_pattern_query_util.cpp · mongodb/mongo @ r8.3.7, L556 // Transforms query into bounds for each field in the shard key // for example : // Key { a: 1, b: 1 }, // Query { a : { $gte : 1, $lt : 2 }, b : { $gte : 3, $lt : 4 } } // => Bounds { a : [1, 2), b : [3, 4) } auto bounds = getIndexBoundsForQuery(cm.getShardKeyPattern().toBSON(), query); ... cm.getShardIdsForRange(min, max, shardIds, ...);

走不了快速路径,mongos 会把分片键里每一个字段各自的查询条件转成一段范围——查询里完全没提到的字段,范围就是无约束的 [MinKey, MaxKey]。把各字段的范围拼起来,就是一段完整分片键范围,拿这段范围去跟每个 chunk 的边界比对,重叠的 chunk 归属哪个分片,就往哪个分片发。如果这段范围覆盖了所有 chunk,那就是每个分片都要问一遍——也就是常说的 scatter-gather。

交互演示:3 种查询,3 种路由结果

同一个分片集合,只是查询条件不一样,mongos 问的分片数量就完全不同。

分片路由实录未开始
点击"下一步"或"播放"开始。

路由判断由 Python 脚本按上面两段真实源码转写:快速路径要求分片键每个字段都有相等值,兜底路径把缺失字段当成 [MinKey,MaxKey] 无约束范围。3 个场景的模型判断结果,和真实 mongos explain()winningPlan.shards 列出的分片名逐一断言相等;chunk 定位函数还有一套完全独立实现的遍历算法交叉核对过。

真实实验的一个意外收获:直连分片会看到"幽灵数据"

chunk 迁移之后,原来那个分片上的旧数据是立刻删除,还是留了一手?

真实实测
$ mongosh --port 27102 --eval 'db.getSiblingDB("config").rangeDeletions.find()'
[ { nss: 'shopdb.orders', donorShardId: 'shard1rs', range: { min: {region:'south',...}, max: {region:MaxKey,...} }, whenToClean: 'delayed', numOrphanDocs: 200 } ] # 通过 mongos 查,数量是对的: $ mongosh --port 27100 shopdb --eval 'db.orders.countDocuments({region:"south"})' 100 # 但直接连 shard1(旧主人)查,south 的"幽灵数据"还在: $ mongosh --port 27102 shopdb --eval 'db.orders.countDocuments({region:"south"})' 100 // 这 100 条本该属于 shard2 的文档,物理上还没被删

chunk 迁移完成后,数据的归属权(元数据)立刻切换,但物理删除是延迟执行的(whenToClean: 'delayed'),给还在读旧快照的请求留出缓冲时间。mongos 靠元数据过滤,查出来的数量永远是对的;但如果绕开 mongos 直接连某个分片查,就可能看到本不属于它、还没来得及清理的"孤儿文档"——这是分片集群里一条很具体的运维原则的来源:永远通过 mongos 访问数据,不要直连分片。

参考与说明

  • 本文源码引用(getShardIdsAndChunksForCanonicalQuery 快速路径与兜底路径、extractShardKeyFromQuery)均取自 mongodb/mongo 仓库 r8.3.7 标签,与本机安装的 MongoDB Community 8.3.7 版本一致,直接从 GitHub 拉取源文件核对过函数签名、关键调用和行号。
  • 集群搭建、chunk 迁移、explain() 输出、直连分片查询结果均为本机真实跑起来的 1 配置服务器 + 2 分片 + 1 mongos 集群产生,未做删改;实验结束后已关闭全部 4 个测试用的进程。
  • 演示数据的自检:3 个场景的分片路由结果,和真实 explain() 捕获的 shards 列表逐一断言相等;chunk 定位用两种完全不同的遍历算法独立算过,结果一致;"没有分片键字段的查询必须命中全部分片"这条断言也验证过。
  • 没有涉及:哈希分片(hashed sharding)与范围分片的选型权衡、balancer 自动均衡的触发条件、分片键选择不当导致的"热点分片"问题、跨分片事务与聚合的具体执行方式(第四篇聚合管道提到过管道下推,分片场景下 mongos 还要再做一层跨分片合并)。
到这里,MongoDB 系列按照最初那条招聘要求的顺序全部写完了:整体架构 → WiredTiger 存储引擎 → 索引优化 → 聚合管道 → oplog → 副本集选举 → 读写一致性 → 分片策略。下一个系列会回到之前搁置的 Kubernetes,从控制平面往下继续。
☕ 如果这篇文章帮到你,可以请作者喝杯咖啡 · 爱发电