MongoDB 分片键选错,8.2 亿条数据全压在一个分片上:我的诊断和补救全过程

🔑 关键词:MongoDB分片键,MongoDB写热点,hash分片,chunk迁移,refineCollectionShardKey

📖 摘要:记录一次真实的分片键选型失误:8.2 亿条气象数据因为分片键单调递增导致写热点,三个分片 CPU 90/15/15。包含诊断命令、四个候选方案对比、迁移实操细节,以及分片键选型检查清单。

先说背景。2023 年 4 月我接手一个气象物联网平台,8.2 亿条文档,三个分片的 MongoDB 集群,版本 5.0,后来升到 6.0。上线半年之后运维找过来,说 shard01 的 CPU 长期挂在 90% 往上,另外两个分片在 15% 左右摸鱼。我第一反应是慢查询,explain('executionStats') 跑了三十多条,totalDocsExamined 和 nReturned 的比例都挺健康,索引看着没毛病。那天下午我盯着 mongostat 看了二十分钟才反应过来——问题不在读,在写。三个分片的 insert 速率是 1:1:1,但写入排队的时间全堆在一个分片上。

图片

分片键是我自己定的:{deviceId: 1, ts: 1}。当时的想法特别朴素,业务查询基本都是按设备查一段时间,这个复合键天然支持前缀匹配,ts 放在第二位刚好能走范围扫描。问题出在 deviceId 上——它是雪花算法生成的自增 ID,写入天然单调递增。MongoDB 的范围分片是按 chunk 的上下界切的,新写入永远落在最后一个 chunk 上,而那个 chunk 永远在同一个 shard。于是不管我加多少个分片,写入永远是一个分片在扛。这就是所谓写热点,教科书上讲的时候轻飘飘一句话,真撞上就是三个月。

一、先把现状看清楚,别急着改

sh.status() 的输出又长又乱,嵌套好几层,我基本不看了。MongoDB 7.0 开始有个更好使的聚合阶段:

db.getSiblingDB('admin').aggregate([
  { $shardedDataDistribution: {} }
])

它会把每个分片上的 numChunks、numOwnedDocuments、numOrphanedDocs 直接拍在你脸上。我当时跑出来大概是 shard01 有 132 个 chunk,shard02 21 个,shard03 19 个;8.2 亿文档里 shard01 占了 6.3 亿。那一刻我心里就有数了。

顺带说一句,6.0 及以前的集群只能靠 sh.status() 硬看,或者去 config 库翻 config.chunks 集合自己统计,那个集合几万条记录,count 一次都要等半天。

二、为什么不能简单换个分片键

图片

这是最多人栽的地方。

分片键在 MongoDB 里几乎是终身的。5.0 之前完全不可改,要改只能把数据全量导出重建集合。5.0 引入了 refineCollectionShardKey,但注意它的限制:

db.adminCommand({
  refineCollectionShardKey: 'station.observation',
  key: { deviceId: 1, ts: 1, siteId: 1 }
})

它只能在原有分片键后面追加后缀字段,不能动前缀。也就是说,我原来 {deviceId: 1, ts: 1} 这个结构,deviceId 这个罪魁祸首是删不掉的。

还有个更隐蔽的连带影响:分片集合上的唯一索引,必须把分片键作为索引前缀,否则建索引的时候直接抛错。我们项目里有个 {deviceId: 1, ts: 1} 的唯一索引,我后来想加个 {sn: 1} 的唯一约束,折腾半天才发现这个规则。另外 upsert 在分片集合上的行为也变了——如果查询条件里不含完整分片键,mongos 会把 upsert 广播到所有分片去执行,同一条数据可能被多个分片尝试插入,靠 _id 唯一性去重,性能很难看。

三、我比较过的四个方案

方案 A:{siteId: 1, ts: 1}。站点总数几百个,分布算均匀,比 deviceId 强得多。但同一个站点下的所有设备写入还是往一起挤,热点没根治,只是从一个分片热变成几个分片轮流热。而且查某个设备最近 24 小时数据这类查询用不上分片键前缀,得广播。

图片

方案 B:{_id: 'hashed'}。写入被打得稀碎,分布完美。但范围查询彻底废了,每个按时间的查询都要 scatter-gather 到全部分片,我们的看板查询 90% 是时间范围扫描,这个代价我承受不起。哈希分片这个东西,适合随机点查,不适合时序。

方案 C:{siteId: 1, ts: 1} 加上 zone sharding。手动把不同 site 的 chunk tag 到指定分片,灵活度最高。代价是 tag range 要自己维护,站点增删的时候你得起个流程去改,运维成本实打实。

方案 D:不改分片键,在写入前面加一层 Kafka 或者内存缓冲,做批量写入。这是治标,写吞吐上去了但 chunk 分布一点没变,而且引入了新的故障点。我排除掉了。

最后我选的是 A 加半个 C:分片键换成 {siteId: 1, ts: 1},siteId 上做了粗粒度的 zone 配置。迁移花了大概三周,全在凌晨窗口跑。

四、迁移这件事的一些实操细节

balancer 的时间窗口必须限制住,默认它是全天候跑的:

use config
db.settings.updateOne(
  { _id: 'balancer' },
  { $set: { activeWindow: { start: '02:00', stop: '05:00' } } },
  { upsert: true }
)

图片

chunk size 默认 64MB,我一开始想改成 32MB 让迁移粒度更细,后来放弃了——chunk 太小会让 config server 的元数据操作量暴涨,而且 balancer 更频繁地触发迁移,反而更吵。这个值在 MongoDB 5.0 之后可以设到 1MB 到 1024MB 之间,但不是说越小越好。

迁移过程中监控:

db.currentOp({ 'command.moveChunk': { $exists: true } })

如果看到 moveChunk 卡着不动,八成是目标分片的磁盘 IO 到顶了。

还有一个东西必须检查——孤儿文档。迁移过程中如果有客户端还带着老路由信息往旧分片写,就会留下 orphaned documents。7.0 的 $shardedDataDistribution 会直接给你 numOrphanedDocs,我迁移完之后看是 0,纯属运气好。如果不清掉,得用:

db.adminCommand({ cleanupOrphaned: 'station.observation' })

图片

这个命令不能对整张表跑,只能对着某个 chunk 的范围跑,而且会短暂阻塞,我没实际用过,只是备着。

五、顺便说说 PostgreSQL

我们同期还有一套 PostgreSQL 15 的分区表系统,量级差不多。说句公道话,在数据分布不均匀这件事上,PG 的原生分区表处理得更优雅——你随时可以 ATTACH 一个新分区,或者 DETACH 掉老的,不用跟 balancer、chunk 边界、孤儿文档这些东西较劲,分区裁剪的逻辑也更好预测。

那为什么这个项目我最后还是留在 MongoDB?两个原因。一是聚合管道,$facet、$bucketAuto、$setWindowFields 写多维度看板是真的省事,同样的东西在 SQL 里要写一堆 CTE 和窗口函数,调试成本高不少。二是 document 模型对气象数据这种每个设备上报字段还不一样的场景,省掉了大量 schema 演进的工作。

我的观点可能有点偏:MongoDB 的分片键不是一个配置项,它是一个架构决策。你在选键的那一刻,等于隐式规定了未来三年哪些查询可以快、哪些查询注定要慢。这个代价在选型阶段基本没人跟你讲,大家聊的都是 MongoDB 支持水平扩展。

六、如果你现在正要选分片键,拿这几条去对

第一,先算基数。跑一下:

图片

db.observation.aggregate([
  { $group: { _id: '$候选字段', count: { $sum: 1 } } },
  { $sort: { count: -1 } },
  { $limit: 20 }
])

看 Top20 占比。如果前 20 个值占了 40% 以上的数据量,这个字段做分片键基本就是制造热点。

第二,单调递增的字段单独做分片键,等于自杀。自增 ID、ObjectId、纯时间戳,都属于这一类。ObjectId 前四个字节是时间戳,所以它整体也是递增的,很多人拿 _id 直接当分片键,然后就在写热点里出不来了。

第三,分片键最好是你大部分高频查询的前缀。做不到全部,至少覆盖 80%。

第四,提前想好退路。5.0 的 refineCollectionShardKey 只能追加后缀,不能改前缀,所以选键的时候就要考虑后面还能不能补字段。

第五,唯一索引、upsert、$lookup 这几个地方的行为都会受分片键影响,选完键之后先在测试环境把这几个场景跑一遍。

现在这套集群,三个分片的 chunk 分布大概在 40/32/28 这个区间浮动,写不再集中了。P99 写入延迟我没做精确的前后对比,但运维再也没在群里 at 过我,这大概就是最好的指标。

🏷️ 标签: