Shardingsenior8+ years

A collection is sharded on `_id` (the default `ObjectId`) to spread write load across the cluster, but one shard's disk and CPU stay pinned near 100% while the others sit idle. Why does sharding on `_id` fail to spread inserts here, and what shard key would actually fix it?

A default ObjectId is monotonically increasing — its leading bytes encode a timestamp, so every newly generated id is numerically greater than the one before it. A range-sharded collection assigns contiguous ranges of the shard key to specific shards, and a monotonically increasing key means every single insert, no matter when it happens, lands at the same end of that range — the shard currently holding the highest range of ids. Every other shard sits idle for inserts while one shard absorbs the entire write load, which is exactly the hotspot this question describes. The fix is a shard key that spreads new values across the whole key space rather than always incrementing from the same end — either hashing the key (shardKey: {_id: "hashed"}, which sends logically sequential ids to effectively random shards) or choosing a genuinely high-cardinality, non-monotonic field the application's real queries filter on.

The lesson behind it →