skip to content

How would you use MongoDB zone sharding to keep recent data on fast shards and older data on cheap shards?

level: principalimportance: should knowfreq 28%

answer

  1. the leading key field must be time-ordered
  2. label the tiers, then map ranges to them
  3. the boundary does not move itself
  4. one period per roll must fit the window
  5. new inserts all land on the newest range

basics

~20 s

Shard on a time-ordered leading field, label fast shards as a hot zone and cheap ones as an archive zone, and map recent key ranges to hot and older ranges to archive. A scheduled job rolls the boundary forward, and the balancer does the moving.

solid answer

~50 s

The design has four parts. **Shard key**: the leading field must be time-ordered — a date or a bucketed period such as `{month: 1, deviceId: 1}` — because zone ranges are shard-key ranges. **Zones**: label the NVMe shards `hot` and the large-disk shards `archive`. **Ranges**: map everything from the cutoff upward to `hot` and everything below it to `archive`. **A roll-forward job**: on a schedule, remove the old boundary ranges and re-map them so the cutoff advances; the balancer then drains aged data down to the archive tier. What you must size is the *migration budget*: one period's worth of data has to cross the network every roll, one migration per shard at a time, inside whatever balancer window you allow. If a month of data cannot move in a month of windows, the design fails quietly by never converging. Accept the trade-off that a time-ordered leading key concentrates inserts on the newest range, so the hot tier must absorb the whole write rate.

code

javascript · 10 lines
javascript
// shard key leads with a bucketed period, e.g. "2026-08"
sh.shardCollection("app.readings", { month: 1, deviceId: 1 })
sh.addShardToZone("shard-nvme-0", "hot")
sh.addShardToZone("shard-hdd-0", "archive")

const cut = "2026-06";
sh.updateZoneKeyRange("app.readings",
  { month: MinKey, deviceId: MinKey }, { month: cut, deviceId: MinKey }, "archive")
sh.updateZoneKeyRange("app.readings",
  { month: cut, deviceId: MinKey }, { month: MaxKey, deviceId: MaxKey }, "hot")

go deeper

for a junior

Know that MongoDB can pin ranges of a sharded collection to particular shards, so a cluster can mix fast and cheap hardware for recent and old data.

for a middle

Be able to explain why the shard key's leading field must be time-ordered for age-based zoning, and that the balancer performs the actual movement.

for a senior

Show that you would schedule and monitor the boundary roll, size the per-roll migration volume against the balancer window, and watch the orphan and lag signals it generates.

for a principal

Own the whole trade: cost per terabyte against permanent migration load, insert concentration on the newest range, archive latency on user-facing reads, and the honest alternative of moving cold data out of the cluster entirely.

## Why zones can express tiering at all Zone sharding constrains which shards may hold which shard-key ranges. If the leading field of the shard key is time-ordered, then ranges *are* time windows, and a rule about ranges becomes a rule about age. That is the whole trick: the same mechanism used for geographic residency becomes hot/cold tiering purely by changing what the leading field means. ## The design **Shard key.** Something like `{month: 1, deviceId: 1}` where `month` is a bucketed period such as `"2026-08"`, or a date field truncated to a period. Buckets are usually better than raw timestamps: they give you a small, predictable number of boundaries to manage, and each roll moves one well-understood unit of data. **Zones.** `sh.addShardToZone("shard-nvme-0", "hot")` for the fast shards, `sh.addShardToZone("shard-hdd-0", "archive")` for the cheap ones. Nothing prevents the two tiers having different instance sizes, disk types or even member counts. **Ranges.** Everything from the cutoff up to `MaxKey` maps to `hot`; everything from `MinKey` to the cutoff maps to `archive`. Cover the full key space so no range escapes onto the wrong tier. **The roll-forward job.** Nothing about the cutoff is automatic. A scheduled job must remove the boundary ranges and re-declare them with the new cutoff, after which the balancer drains the newly-aged period down to archive: ```javascript const oldCut = "2026-05", newCut = "2026-06"; sh.removeRangeFromZone("app.readings", { month: oldCut, deviceId: MinKey }, { month: MaxKey, deviceId: MaxKey }); sh.updateZoneKeyRange("app.readings", { month: MinKey, deviceId: MinKey }, { month: newCut, deviceId: MinKey }, "archive"); sh.updateZoneKeyRange("app.readings", { month: newCut, deviceId: MinKey }, { month: MaxKey, deviceId: MaxKey }, "hot"); ``` That job is production code: it needs idempotency, alerting when it does not run, and a check that the previous roll actually converged before the next one starts. ## The capacity question that decides whether this works Every roll pushes one period of data across the network. The cluster can run at most about *n*/2 concurrent migrations for *n* shards, each move copies a range's worth of bytes, and migrations only start inside the balancer's active window. Multiply it out before committing: if a month is 4 TB and your nightly window realistically moves 200 GB, the tier never converges and you are permanently mid-migration, paying the IO cost without ever getting the placement. Levers if the arithmetic fails: shorter buckets so each roll is smaller and more frequent; a wider or continuous balancer window with `_secondaryThrottle` to protect replication instead of a hard time limit; more shards in each tier to raise migration parallelism; or abandoning migration-based tiering in favour of writing cold data to a separate cluster or an object store. ## The costs you are signing up for **Insert concentration.** A time-ordered leading shard-key field means all new writes land in the newest range on one shard. The hot tier must be sized for the entire write rate, not for its share of the data, and adding hot shards does not spread inserts. This is an accepted consequence of making tiering expressible, and you should say so out loud rather than pretend the balancer will fix it. **Archive-tier read latency.** Queries touching old periods land on slow disks. That is the point, but it must be a product decision: if reporting regularly scans two years, the archive tier is on the critical path of a user-visible feature. **Scatter-gather on non-time queries.** A query without the leading field fans out across both tiers, so its latency is set by the slowest tier. Access patterns that ignore time undermine the whole arrangement. **Permanent background load.** Unlike a one-off rebalance, tiering means the cluster is always moving data. Range deletion on the donors is continuous too, and the orphan backlog needs watching. ## When not to do this If old data is only ever deleted, a TTL index and a shard key chosen for query locality are simpler and cheaper than moving data you are about to discard. If old data is read a few times a year, exporting it out of the operational cluster entirely usually beats keeping it online on slow shards. Migration-based tiering earns its keep when cold data must stay queryable through the same collection, with the same code paths, and the cost difference between tiers is large enough to pay for the movement. ## What to monitor Whether the roll job ran; whether the previous roll converged (per-shard byte totals per tier); migration throughput from `config.changelog`; range-deletion backlog on donors; replication lag on recipients; and the write latency on the newest range, which is the single most likely thing to degrade first.

  • What breaks first if the roll-forward job silently stops running?
    Nothing alarms immediately — that is the danger. The hot tier keeps absorbing new periods that were never mapped away, so it fills up while the archive tier sits idle. You discover it as a disk-space alert on the fast shards weeks later, and the catch-up migration is then several periods of data instead of one.
  • Why not just add more hot-tier shards instead of tiering?
    Because a time-ordered leading shard key sends all inserts to the newest range on one shard, extra hot shards add capacity but not insert throughput. They also multiply the cost of the expensive tier for data nobody reads. Tiering is worth it when the cold data must stay queryable in place and the per-terabyte price gap justifies the migration traffic.
  • How do you verify a tier has actually converged rather than assuming it has?
    Compare measured bytes per shard per tier against what the zone ranges say should be there, and check config.changelog for migrations still in flight for that collection. sh.status() shows the zone definitions and current range placement; a range listed on a shard outside its zone means the balancer still has work queued or has no window to do it in.

saying these in an interview costs you the question

  • Assumes the age boundary advances by itself with no scheduled job
  • Ignores whether a period of data can actually move in the available window
  • Claims adding hot shards spreads inserts from a time-ordered shard key
  • Uses tiering where a TTL index and deletion would do
  • Forgets queries without the time field fan out to the slow tier

context