Learn
MongoDB/17-sharding

分片集群

复制集解决了「机器会挂」,但解决不了「一台机器装不下」和「一台机器写不动」。分片(sharding)是 MongoDB 的水平扩展方案:把数据按某个键切分到多个复制集上。它的威力很大,坑也很深——分片键选错了几乎无法挽回。

1. 什么时候需要分片

1.1 三个信号

信号说明
存储超过单机容量数据量增长到单机磁盘装不下
工作集超过内存热数据装不进 RAM,命中率暴跌
写入超过单机 IOPSprimary 的写入已经打满磁盘

注意读吞吐不在这个列表里——读压力大应该先加 secondary 或者上缓存,那比分片便宜得多。

1.2 分片的代价

单复制集:3 台机器,1 套运维
                 ↓ 分片成 3 片
分片集群:3 × 3 = 9 台数据节点
        + 3 台 config server
        + 2 台 mongos(或和应用同机部署)
        = 14 个进程要监控、备份、升级
 
  运维复杂度增加 4 到 5 倍

其他代价:

  • 不带分片键的查询会广播到所有分片
  • $lookup 跨分片时性能更差
  • 唯一索引只能建在分片键上(或包含分片键的前缀)
  • 备份恢复要考虑分片间一致性
  • 分片键选错后修正成本极高
⚠️不要过早分片

「以后数据会很大」不是现在分片的理由。一个配置得当的复制集能轻松处理几 TB 数据和几万 QPS。先把索引优化好、把建模改对、把缓存加上,这些的收益通常比分片大得多,而成本低一个数量级。真正需要分片时,你会从监控指标上看得很清楚。

2. 架构回顾

                  应用(driver)
                        │
              ┌─────────┴─────────┐
              ▼                   ▼
        ┌──────────┐        ┌──────────┐
        │  mongos  │        │  mongos  │      无状态路由,可水平扩展
        └────┬─────┘        └────┬─────┘
             │                   │
             └────────┬──────────┘
                      │  查/缓存元数据
             ┌────────▼─────────┐
             │  Config Servers  │  复制集,存 chunk 分布
             │  (csrs)          │
             └────────┬─────────┘
                      │
       ┌──────────────┼──────────────┐
       ▼              ▼              ▼
  ┌─────────┐    ┌─────────┐    ┌─────────┐
  │ Shard A │    │ Shard B │    │ Shard C │   每片是完整复制集
  │ rs-a    │    │ rs-b    │    │ rs-c    │
  │ P S S   │    │ P S S   │    │ P S S   │
  └─────────┘    └─────────┘    └─────────┘

2.1 启用分片

// 连接 mongos
sh.enableSharding("community")
 
// 先建索引(分片键必须有索引)
db.posts.createIndex({ authorId: 1, _id: 1 })
 
// 对集合分片
sh.shardCollection("community.posts", { authorId: 1, _id: 1 })
{ "collectionsharded": "community.posts", "ok": 1 }
sh.status()
shards:
  { "_id": "rs-a", "host": "rs-a/a1:27017,a2:27017,a3:27017", "state": 1 }
  { "_id": "rs-b", "host": "rs-b/b1:27017,b2:27017,b3:27017", "state": 1 }
 
databases:
  { "_id": "community", "primary": "rs-a", "partitioned": true }
    community.posts
      shard key: { "authorId": 1, "_id": 1 }
      chunks:
        rs-a  42
        rs-b  41

3. chunk 与均衡器

3.1 chunk 是什么

MongoDB 不是按文档分片,而是按 chunk(数据块)。每个 chunk 是分片键的一个连续区间。

分片键 authorId 的取值空间被切成若干区间:
 
  [MinKey,  "carol")  → Shard A
  ["carol", "frank")  → Shard A
  ["frank", "linda")  → Shard B
  ["linda", MaxKey)   → Shard B
 
  一个 chunk 默认最大 128MB(6.0 前是 64MB)
  超过就分裂(split)成两个
// 查看某个集合的 chunk 分布
db.getSiblingDB("config").chunks.find({ ns: "community.posts" }).limit(3)
[
  { "_id": "...", "min": { "authorId": MinKey }, "max": { "authorId": "carol" }, "shard": "rs-a" },
  { "_id": "...", "min": { "authorId": "carol" }, "max": { "authorId": "frank" }, "shard": "rs-a" },
  { "_id": "...", "min": { "authorId": "frank" }, "max": { "authorId": "linda" }, "shard": "rs-b" }
]

3.2 均衡器

后台进程,负责把 chunk 从「多的分片」搬到「少的分片」。

sh.getBalancerState()
sh.startBalancer()
sh.stopBalancer()
 
// 设置均衡窗口,只在低峰期搬迁
db.getSiblingDB("config").settings.updateOne(
  { _id: "balancer" },
  { $set: { activeWindow: { start: "02:00", stop: "06:00" } } },
  { upsert: true }
)
⚠️chunk 迁移会消耗大量资源

迁移一个 chunk 要:从源分片读取全部文档、写入目标分片、更新 config server 元数据、删除源分片数据。这个过程占用磁盘 IO、网络和 CPU。给繁忙的集群设置均衡窗口是必备操作,否则白天业务高峰期突然开始迁移,延迟会明显抖动。

3.3 6.0+ 的改进

6.0 起均衡器改成按数据量均衡(而不是 chunk 数量),并支持自动 chunk 合并,减少了因 chunk 数量相同但大小悬殊导致的不均衡。

4. 分片键选择(最重要)

分片键决定了一切:数据怎么分布、查询走不走定向路由、写入会不会集中在一个分片。

4.1 三个评价维度

维度含义不满足的后果
基数(Cardinality)不同取值的数量取值太少 → chunk 无法继续分裂
频率(Frequency)取值分布是否均匀少数值占大头 → 数据倾斜
变化率(Rate of change)是否单调递增单调递增 → 写入全打在一个分片

4.2 反例逐个分析

反例 1:用低基数字段

sh.shardCollection("community.posts", { status: 1 })
status 只有 3 个值 → 最多只能有 3 个 chunk → 最多用 3 个分片
而且 published 占 92% → 一个分片承载 92% 的数据
 
  ┌─────────┐  ┌─────────┐  ┌─────────┐
  │ Shard A │  │ Shard B │  │ Shard C │
  │ 9.2 TB  │  │ 0.7 TB  │  │ 0.1 TB  │   ← 完全没有扩展效果
  └─────────┘  └─────────┘  └─────────┘

反例 2:用单调递增字段

sh.shardCollection("community.posts", { _id: 1 })      // ObjectId 递增
sh.shardCollection("logs.events", { createdAt: 1 })    // 时间递增
  chunk 区间:
  [MinKey, id_1000)  Shard A     ← 历史数据,只读
  [id_1000, id_2000) Shard B     ← 历史数据,只读
  [id_2000, MaxKey)  Shard C     ← 所有新写入都落这里!
 
  写热点:Shard C 承担 100% 的写入
  其他分片的写能力完全浪费

反例 3:用高频重复值

sh.shardCollection("community.comments", { postId: 1 })

一篇爆款文章有 50 万条评论,它们的 postId 完全相同 → 落在同一个 chunk → 这个 chunk 无法分裂(叫 jumbo chunk),会一直增长直到拖垮那个分片。

4.3 三种分片策略

范围分片(ranged)

sh.shardCollection("community.posts", { authorId: 1, createdAt: 1 })
  • 优点:范围查询高效(连续的键在同一分片)
  • 缺点:单调递增的键会造成写热点

哈希分片(hashed)

db.events.createIndex({ userId: "hashed" })
sh.shardCollection("community.events", { userId: "hashed" })
  • 优点:分布均匀,天然无写热点
  • 缺点:范围查询必须广播到所有分片
范围分片,查 authorId 在 c 到 f 之间:
  只需要问 Shard A  ✔
 
哈希分片,同样的查询:
  哈希值不连续 → 问所有分片  ✘

组合分片键(推荐)

用「低基数但分布均匀的前缀」+ 「高基数的后缀」:

// 前缀是用户 id(打散),后缀是时间(同用户内有序)
sh.shardCollection("community.events", { userId: 1, ts: 1 })
 
// 或者人为加一个随机前缀
sh.shardCollection("logs.events", { hourBucket: 1, _id: 1 })

4.4 内容社区的分片键设计

集合分片键理由
users{ _id: "hashed" }按 id 精确查询为主,无范围查询需求
posts{ authorId: 1, _id: 1 }「某作者的帖子」是高频查询,能定向路由
comments{ postId: 1, _id: 1 }「某帖子的评论」定向路由,_id 后缀避免 jumbo
follows{ followerId: 1, followeeId: 1 }「我关注了谁」定向;查粉丝要广播(可另建反向集合)

注意 comments 的键:只用 postId 会有 jumbo chunk 风险,加上 _id 后缀,同一个 postId 的评论也能继续分裂成多个 chunk。

💡分片键设计口诀

高基数、分布均匀、包含在高频查询条件里、不单调递增。四条同时满足的键很难找,通常需要组合两三个字段。设计完之后一定要用真实数据分布验证一遍,别只看字段名想当然。

5. 查询路由

5.1 定向查询 vs 广播查询

定向查询(targeted):查询条件包含分片键前缀
  db.posts.find({ authorId: "alice" })
 
    mongos ─── 查元数据 ──▶ "alice" 在 Shard A
           └───────────────▶ 只问 Shard A
    延迟:一次网络往返
 
广播查询(broadcast / scatter-gather):不含分片键
  db.posts.find({ tags: "mongodb" })
 
    mongos ──┬──▶ Shard A
             ├──▶ Shard B
             └──▶ Shard C
             合并结果
    延迟:取决于最慢的分片;分片越多越慢
// 用 explain 确认走的是哪种
db.posts.find({ authorId: "alice" }).explain()
{
  "queryPlanner": {
    "winningPlan": {
      "stage": "SINGLE_SHARD",
      "shards": [ { "shardName": "rs-a" } ]
    }
  }
}

广播查询的 stage 会是 SHARD_MERGE。

5.2 排序与分页的放大

db.posts.find().sort({ createdAt: -1 }).skip(1000).limit(20)
3 个分片时,mongos 必须:
  向每个分片要 skip + limit = 1020 条
  收到 3060 条
  在内存里归并排序
  丢弃 3000 条,返回 20 条
 
  分片数越多,放大越严重。这也是第 6 章强调游标分页的另一个理由。

5.3 分片下的唯一索引

唯一索引只能建在包含分片键前缀的字段上。

sh.shardCollection("community.users", { _id: "hashed" })
 
db.users.createIndex({ email: 1 }, { unique: true })
MongoServerError: Cannot create unique index over { email: 1 }
with shard key pattern { _id: "hashed" }

原因:唯一性检查需要看到全部数据,而每个分片只有一部分。要在分片集群里保证 email 唯一,有三个办法:

  1. 把 email 作为分片键(或分片键前缀)
  2. 用 email 当 _id,_id 天然唯一
  3. 建一个专门的 email_registry 集合,用 email 做 _id 和分片键,写入时先占位
⚠️这是分片对建模的硬约束

业务里往往有多个唯一字段(username、email、phone),而分片键只有一个。这个约束必须在分片之前想清楚,否则上线后发现唯一性没法保证,改造成本极高。

6. 分片键不可变与重分片

6.1 历史限制

  • 4.2 前:分片键的值不可修改
  • 4.2 起:可以修改分片键的值(文档会在分片间迁移),但仍需在事务里
  • 4.4 起:可以给分片键追加后缀字段(refine shard key)
  • 5.0 起:支持完整的 reshardCollection,可以换成完全不同的键

6.2 refineCollectionShardKey

给现有分片键加后缀,用于解决 jumbo chunk:

db.comments.createIndex({ postId: 1, _id: 1 })
 
db.adminCommand({
  refineCollectionShardKey: "community.comments",
  key: { postId: 1, _id: 1 }
})

这个操作很轻量(只改元数据,不搬数据),因为原有 chunk 的边界仍然有效,只是以后能继续分裂了。

6.3 reshardCollection

db.adminCommand({
  reshardCollection: "community.posts",
  key: { authorId: "hashed" },
  numInitialChunks: 128
})
 
// 监控进度
db.getSiblingDB("admin").aggregate([ { $currentOp: { allUsers: true } },
  { $match: { type: "op", "originatingCommand.reshardCollection": { $exists: true } } } ])

原理是:在后台建一个新的临时集合按新键分布,用 Change Streams 追平增量,最后原子切换。

  代价:
  - 需要额外的磁盘空间(约等于集合大小)
  - 持续数小时到数天
  - 切换瞬间有短暂(约 2 秒)的写阻塞
  - 期间集群负载显著升高
⚠️reshard 是补救不是常规操作

虽然 5.0 提供了 reshardCollection,但对一个 10TB 的集合执行它可能需要几天,期间集群性能受损。它是给「分片键选错了」的项目一条活路,不是让你可以随便选分片键。设计阶段多花两天想清楚,比后面花两周 reshard 划算。

7. 区域分片(Zone Sharding)

把特定的数据范围绑定到特定分片,用于数据本地化或冷热分层:

sh.addShardToZone("rs-cn", "CN")
sh.addShardToZone("rs-eu", "EU")
 
sh.updateZoneKeyRange(
  "community.users",
  { region: "CN", _id: MinKey },
  { region: "CN", _id: MaxKey },
  "CN"
)
sh.updateZoneKeyRange(
  "community.users",
  { region: "EU", _id: MinKey },
  { region: "EU", _id: MaxKey },
  "EU"
)

典型用途:

场景做法
数据合规(GDPR)欧洲用户数据只存在欧洲机房的分片
就近访问按地域分区,降低跨地域延迟
冷热分层热数据放 SSD 分片,冷数据放 HDD 分片

8. 运维要点

// 集群整体状态
sh.status()
 
// 各分片的数据量
db.posts.getShardDistribution()
Shard rs-a at rs-a/a1:27017,a2:27017,a3:27017
  data : 42.31GiB docs : 12048221 chunks : 342
  estimated data per chunk : 126.7MiB
 
Shard rs-b at rs-b/b1:27017,b2:27017,b3:27017
  data : 41.88GiB docs : 11982104 chunks : 339
  estimated data per chunk : 126.5MiB
 
Totals
  data : 84.19GiB docs : 24030325 chunks : 681
  Shard rs-a contains 50.25% data, 50.13% docs in cluster

需要盯的指标:

指标健康值异常含义
各分片数据占比相差 10% 以内分片键分布不均
jumbo chunk 数量0分片键基数不足
均衡器迁移次数稳定后应该很少频繁迁移说明写入分布不均
广播查询比例越低越好说明查询没带分片键
🎯练习

一、用 Docker 搭一个 2 分片的集群(每片单成员复制集 + 1 个 config 复制集 + 1 个 mongos);二、对 posts 用 { _id: 1 } 分片,插入 10 万条数据,用 getShardDistribution() 观察分布,解释为什么不均;三、改用 { authorId: "hashed" } 重做,对比分布;四、分别执行带分片键和不带分片键的查询,用 explain 对比 SINGLE_SHARD 与 SHARD_MERGE;五、尝试在非分片键字段上建唯一索引,复现报错,然后设计一个「注册表集合」方案绕过它;六、为内容社区的 comments 集合设计分片键,说明为什么必须加 _id 后缀。

小结

  • 分片解决存储容量和写入吞吐的上限,不解决读吞吐(那应该加从库或缓存)
  • 分片的运维复杂度是复制集的 4 到 5 倍,不要过早分片
  • 数据按 chunk 划分,均衡器负责搬迁,必须设置均衡窗口
  • 分片键三要素:高基数、分布均匀、不单调递增;还要出现在高频查询条件里
  • 范围分片利于范围查询但可能有热点,哈希分片分布均匀但范围查询要广播
  • 组合分片键(打散前缀 + 有序后缀)通常是最优解
  • 唯一索引只能建在包含分片键前缀的字段上,这是硬约束
  • 5.0 的 reshardCollection 是补救手段,代价高昂,设计阶段要一次做对
  • 下一章讲备份恢复与安全 →