Inserts into a collection sharded on { createdAt: 1 } all land on one shard — why, and how do you fix it?
answer
- Ranged placement means ordered, contiguous intervals
- Only one shard owns the top of the key space
- Default _id values have the same shape
- Migrating ranges does not change who owns the top
- Put something non-monotonic in front of it
basics
~20 sRanged sharding gives each shard a contiguous interval of key values, and the top interval runs to MaxKey. Because createdAt only increases, every insert falls in that one interval, so a single shard absorbs all writes.
solid answer
~50 sThis is the classic monotonic shard-key hotspot. Ranged sharding orders the key space, and the highest range has an upper bound of `MaxKey` and lives on exactly one shard. Every new document has a `createdAt` greater than all previous ones, so it lands in that range — one shard takes 100% of the insert throughput while the others idle, and reads for recent data concentrate there too. The balancer notices the imbalance and continuously migrates ranges off the hot shard, which adds I/O without addressing the cause, since the next insert still goes to whoever owns the top. The same trap catches `{ _id: 1 }` with default `ObjectId` values, because an `ObjectId` encodes its creation time in its leading bytes. Fix it either by hashing (`{ createdAt: "hashed" }`) or, usually better, by prefixing with a high-cardinality non-monotonic dimension: `{ deviceId: 1, createdAt: 1 }`.
code
javascript · 8 lines// Hotspot: every insert has the largest ts so far and lands in the top range
sh.shardCollection("app.readings", { ts: 1 })
// Fix A: spread by hash, lose time-range targeting
sh.shardCollection("app.readings", { ts: "hashed" })
// Fix B: non-monotonic prefix, per-device time windows stay targeted
sh.shardCollection("app.readings", { deviceId: 1, ts: 1 })go deeper
Recall that ranged placement keeps key values in order, so a key that always grows sends every new document to the same place. Know that a default _id behaves like a timestamp.
Be able to explain the top range with the MaxKey bound, why balancer migrations do not fix it, and name both remedies: hashing the field, or prefixing with a non-monotonic high-cardinality field.
Interviewers expect a diagnosis path — spotting the hot shard and matching it to the top range — plus a defensible choice between hashing and a compound prefix based on the read pattern, and awareness of what changing the key now would cost.
Own the trade: this failure moves cost between write concentration and read fan-out. Be ready to argue which the system can absorb, and to set a standard that time-shaped collections never take a bare timestamp or ObjectId as a leading shard key field.
## Why one shard takes everything A sharded collection's data is split into contiguous ranges of shard key values, each range owned by one shard, and the ranges tile the whole key space from `MinKey` to `MaxKey`. Ranged sharding preserves value order, which is what makes range queries targetable — and what creates this failure mode. There is exactly one range whose upper bound is `MaxKey`, and it lives on exactly one shard at any moment. A monotonically increasing shard key guarantees that every newly inserted document's key is greater than every existing key, so it always falls into that top range. Writes therefore serialise onto one shard's storage engine, one shard's cache, and one shard's journal, regardless of how many shards the cluster has. ## The keys that do this without you noticing `createdAt` is the obvious offender, but the more common one is the default `_id`. An `ObjectId` is a 12-byte value whose leading bytes are a seconds-resolution timestamp, so `ObjectId` values generated over time increase. Sharding on `{ _id: 1 }` therefore behaves almost exactly like sharding on a timestamp. Auto-incrementing counters carried over from a relational design have the same shape, as do sequence-derived invoice numbers and any key built by concatenating a date prefix. Monotonically *decreasing* keys have the identical problem at the bottom of the key space. ## Why the balancer does not save you The balancer's job is to keep data volume roughly even across shards. When one shard's share grows, it schedules migrations of ranges away from it. That is a reaction to the symptom: the ownership of the *top* range is what matters, and after a migration the new owner of the top range simply becomes the new hotspot. In the meantime, migration traffic competes with the very insert load that is already saturating the shard. You also pay in the read path if the workload reads recent data — the shard holding the newest range serves the hottest reads as well as all writes. ## Diagnosing it `db.collection.getShardDistribution()` shows how documents and bytes are spread per shard, and `sh.status()` shows the range map. The tell is not simply an uneven byte count — the balancer keeps that roughly even — but a shard whose write throughput or ticket usage dominates while it owns the top range. Correlating the hot shard with the highest range in `sh.status()` confirms the diagnosis quickly. ## Fix 1: hash the key `sh.shardCollection("app.events", { createdAt: "hashed" })` partitions on the hash of the timestamp instead of on the timestamp. Adjacent instants hash to unrelated values, so consecutive inserts scatter uniformly. The write hotspot disappears immediately. The cost is that value order is gone: a query for the last hour can no longer be narrowed to a few shards and must be broadcast and merged. For an append-only log that is read rarely or by a different key entirely, this is a fine trade. For a dashboard whose entire workload is recent-window queries, it moves the pain rather than removing it. ## Fix 2: put a non-monotonic prefix in front The usual production answer for time-shaped data is a compound ranged key whose prefix has high cardinality and no time trend: `{ deviceId: 1, createdAt: 1 }`, `{ tenantId: 1, createdAt: 1 }`, `{ userId: 1, createdAt: 1 }`. Now the range a document lands in is determined first by the prefix, so concurrent writes from many devices spread across all shards. Within one device the timestamp suffix still gives the range map somewhere to split, so no single device's data becomes indivisible as it grows. Crucially, the common query — one device's recent readings — carries both fields and remains targetable. What you lose is the cross-device query "everything in the last hour", which now fans out; whether that matters depends on whether it is a user-facing path or an analytics job. ## Fix 3: compound hashed `{ deviceId: "hashed", createdAt: 1 }` is a middle path available in MongoDB 4.4 and later, where a compound shard key may contain exactly one hashed field. It guarantees an even device spread even if device ids are themselves clumpy, while keeping an ordered suffix. ## Doing it to an already-sharded collection If the collection is already live on the bad key, changing the *prefix* is not something `refineCollectionShardKey` can do — that command can only append suffix fields. Changing the leading field requires `reshardCollection`, available since MongoDB 5.0, which rewrites the collection under the new key. Plan for the disk and I/O headroom that copy needs and for a short write-blocking cut-over near the end. That expense is precisely why interviewers ask this question: the cheap moment to think about monotonicity is before the first `sh.shardCollection` call.
- Does sharding on { _id: 1 } avoid the problem because _id values are unique?No. Uniqueness gives cardinality, not distribution over time. A default `ObjectId` encodes a timestamp in its leading bytes, so `ObjectId` values increase and behave like a timestamp key: every insert lands in the top range. Uniqueness only means the range map *could* be split finely; it says nothing about where new documents go.
- Your workload both inserts by time and reads the last hour. Which fix do you pick?Prefer a compound ranged key such as `{ deviceId: 1, createdAt: 1 }` if the reads are scoped to a device or tenant — writes spread by prefix and scoped time-window reads stay targeted. Hashing the timestamp would cure the write hotspot but make every last-hour read a broadcast, which is usually the worse trade for a user-facing dashboard.
- If the collection is already sharded on { createdAt: 1 }, can you just refine the key to { createdAt: 1, deviceId: 1 }?You can refine it, since refinement appends suffix fields, but it will not fix the hotspot. Placement is still driven by the leading `createdAt` value, so new documents keep landing at the top of the key space. The refinement only makes the top range splittable. Changing the leading field requires `reshardCollection`.
saying these in an interview costs you the question
- Claims the balancer will redistribute the incoming writes
- Thinks a unique _id shard key spreads inserts evenly
- Suggests adding shards to relieve the write hotspot
- Proposes appending a field to fix a monotonic prefix
- Assumes hashing the timestamp is free with no read-side cost