skip to content

Sharding & Scalability

Scaling horizontally: shard keys, chunk distribution and migration, mongos routing, and zones. Interviewers ask because a bad shard key is effectively unfixable without a migration, which makes it the highest-stakes design question in MongoDB.

part ofMongoDBoverview, primer and where to startread it →
on this pageshow

questions

18

What is a chunk in a sharded MongoDB collection, and what does the balancer do with chunks?

level: middleimportance: must knowfreq 70%

answer

  1. contiguous span of shard-key values
  2. half-open interval, MinKey to MaxKey
  3. exactly one owning shard, not copies
  4. config servers hold the range map
  5. 6.0: balances by data size, no auto-splitter

basics

~20 s

A chunk (range) is a contiguous span of shard-key values owned by exactly one shard and recorded in config-server metadata. The balancer migrates ranges from the shard holding the most data for a collection to the shard holding the least.

solid answer

~50 s

Sharding splits a collection along its shard key into **ranges** (historically called chunks): half-open intervals `[min, max)` of shard-key values. A newly sharded collection starts as a single range from `MinKey` to `MaxKey`, and that space is subdivided as data grows. Each range is owned by exactly one shard — redundancy comes from each shard being a replica set, not from placing a range on two shards. The authoritative range-to-shard map lives in the `config` database on the config server replica set (`config.chunks`). `mongos` caches it and refreshes when a shard replies that its cached version is stale. The **balancer** runs on the config-server primary. Per collection, it compares how much data each shard holds and issues `moveRange` to even things out; a shard takes part in at most one migration at a time. Since MongoDB 6.0 it balances by data size rather than chunk count, and the old mongos-side auto-splitter is gone.

code

javascript · 4 lines
javascript
sh.status()
sh.getBalancerState()      // is balancing enabled?
sh.isBalancerRunning()     // is a round active right now?
db.adminCommand({ balancerCollectionStatus: "app.events" })

go deeper

for a junior

Know that sharding spreads one collection's documents across shards by shard-key range, and that a background process called the balancer keeps the shards roughly even without you asking.

for a middle

Be ready to describe range bounds as a half-open interval, say that the map lives on the config servers and is cached by mongos, and explain that the balancer compares per-shard data size for one collection.

for a senior

Expect to read sh.status() and argue whether a collection is genuinely balanced, and to explain what a migration costs a running cluster and when you would let it happen.

for a principal

Own the framing that the balancer only equalizes stored bytes. Traffic skew, capacity headroom and migration bandwidth are separate design inputs, and no amount of balancing rescues a shard key that concentrates writes.

## The unit of distribution When you shard a collection, MongoDB partitions the *value space* of the shard key into contiguous, non-overlapping intervals. Current documentation calls them **ranges**; interviews and older tooling call them **chunks**. Each interval is half-open — `[min, max)`, lower bound included, upper excluded — so the intervals tile the whole space with no gaps and no overlap. A freshly sharded collection has one range spanning `MinKey` to `MaxKey`, two special BSON types that sort below and above every other value. Ownership is exclusive: a range belongs to exactly one shard at a time. This is the point people most often get wrong. Copies of data exist *within* a shard, because every shard is a replica set; they do not exist *across* shards. Leftover documents on a donor shard after a migration are **orphans**, and shards filter them out of query results using their own copy of the ownership metadata. ## Where the map lives The config server replica set (CSRS) stores the routing table in the `config` database: `config.shards`, `config.databases`, `config.collections`, `config.chunks` (one document per range, with `min`, `max`, the owning `shard`, and a version stamp) and `config.tags` for zones. A `mongos` router caches this table in memory. When a router sends an operation to a shard using a stale version, the shard rejects it with a stale-config error; the router refreshes from the config servers and retries. That is why a migration produces a brief burst of metadata refreshes rather than wrong answers. ## What the balancer is and how it decides The balancer is a background task on the **config-server primary** (it moved there in MongoDB 3.4; before that it ran in `mongos`). It works in rounds and evaluates each sharded collection independently — one collection can be perfectly even while another is skewed. Starting in MongoDB 6.0 the balancer compares **data size per shard for that collection**, not the number of chunks. That matters because ranges are not uniform: one 128 MB range and one 2 KB range used to count the same, so a cluster could report an even chunk count while one shard carried most of the bytes. The documented trigger is a difference between the most- and least-loaded shard larger than three times the configured range size; below that the collection is considered balanced and nothing moves. When it decides to act, the balancer sends `moveRange` to the donor shard's primary. A shard participates in at most one migration at a time, so the number of concurrent migrations in an *n*-shard cluster is bounded by roughly *n*/2 disjoint donor/recipient pairs. ## Splitting and merging Before MongoDB 6.0, `mongos` tracked bytes written and triggered an auto-split once a chunk was estimated to exceed the configured chunk size. That auto-splitter was removed in 6.0: ranges are now split only when the balancer needs to carve off data to move, plus whatever an operator does explicitly with `sh.splitAt()` or `sh.splitFind()`. In the other direction, recent versions merge contiguous ranges owned by the same shard automatically, and `configureCollectionBalancing` with `defragmentCollection: true` runs that consolidation on demand. Fewer, larger ranges mean less routing metadata for every router and shard to cache. ## What you actually look at `sh.status()` prints, per sharded collection, the shards, their ranges and any zones. `sh.getBalancerState()` says whether the balancer is enabled; `sh.isBalancerRunning()` says whether a round is active right now; `db.adminCommand({balancerCollectionStatus: "db.coll"})` says whether the balancer thinks a specific collection still has work to do. Completed moves are logged in `config.changelog` as `moveChunk.start` / `moveChunk.commit` entries, which is the cheapest way to see how much migration traffic a cluster is generating. ## What the balancer does not fix The balancer equalizes *stored bytes*, not *traffic*. If every insert lands at the top of a monotonically increasing shard key, the newest range is always hot and moving it elsewhere just moves the hotspot. Distribution problems caused by shard-key design are not balancer problems.

  • How does mongos know which shard owns a range, and what happens when its cache is out of date?
    It caches the routing table read from the config servers. If it targets a shard with an outdated metadata version, the shard rejects the operation with a stale-config error; mongos refreshes from the config servers and retries transparently. Correctness never depends on the cache being fresh — only latency does, which is why a migration causes a short refresh burst.
  • Does the balancer look at the cluster as a whole or at one collection at a time?
    One sharded collection at a time. Each collection's ranges are balanced independently, so a cluster can have one perfectly even collection and one badly skewed one. sh.status() reports distribution per collection, and balancerCollectionStatus tells you whether the balancer still considers a given collection unbalanced.
  • Why did MongoDB 6.0 stop balancing on chunk count?
    Chunk counts are a poor proxy for load. Ranges vary hugely in size — an indivisible range full of one shard-key value can be enormous while a freshly split neighbour is nearly empty — so an even chunk count could hide a very uneven byte distribution. Balancing on measured data size per shard reflects what actually consumes disk and cache.

saying these in an interview costs you the question

  • Says every shard keeps a copy of each chunk
  • Claims mongos still auto-splits chunks on write in current versions
  • Thinks the balancer equalizes chunk counts in MongoDB 6.0+
  • Says the balancer runs inside mongos rather than on the config-server primary
  • Expects the balancer to fix a hotspot caused by shard-key design

context

open as a page

Which query filters let mongos target one shard, and which force a broadcast to every shard?

level: middleimportance: must knowfreq 72%

basics

~20 s

A filter containing the shard key — or a leading prefix of a compound shard key — lets mongos map values to chunks and contact only the owning shards. Any other filter, including one on _id when _id is not the shard key, is broadcast to all shards.

open as a page

What makes a good MongoDB shard key, and how do cardinality, frequency and monotonicity affect it?

level: middleimportance: must knowfreq 78%

basics

~10 s

A good MongoDB shard key has many distinct values, no single value that dominates the data, values that do not grow monotonically, and it appears in common query filters so requests reach one shard.

open as a page

During a MongoDB chunk migration, what happens on the donor and recipient shards, and how does live traffic feel it?

level: seniorimportance: must knowfreq 55%

basics

~20 s

The recipient clones the range's documents and then catches up on changes; the donor enters a short critical section that blocks writes to that range while the config metadata is committed; afterwards the donor deletes the moved documents in the background as orphans.

open as a page

After sharding, p99 read latency got worse and every shard is busy — how do you confirm broadcast queries are the cause?

level: seniorimportance: must knowfreq 58%

basics

~20 s

Run explain() through mongos on the hot query shapes: a SHARD_MERGE stage with every shard named in the winning plan means scatter-gather. Compare per-shard keysExamined against nReturned, then fix the filter to carry the shard key rather than adding shards.

open as a page

Inserts into a collection sharded on { createdAt: 1 } all land on one shard — why, and how do you fix it?

level: seniorimportance: must knowfreq 72%

basics

~20 s

Ranged 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.

open as a page

In a MongoDB sharded cluster, what does mongos do with a find() the driver sends it?

level: juniorimportance: should knowfreq 45%

basics

~20 s

mongos is the query router. It uses the chunk map cached from the config servers to pick which shards can hold matches — one when the filter names the shard key, all shards otherwise — and merges their replies into a single cursor.

open as a page

How does MongoDB's configured chunk (range) size affect balancing, and when would you change the default?

level: middleimportance: should knowfreq 45%

basics

~20 s

The default range size is 128 MB in MongoDB 6.0 and later (64 MB before), settable from 1 to 1024 MB. Smaller ranges balance more finely but produce more migrations and more routing metadata; larger ranges migrate less often but distribute more coarsely.

open as a page

How does mongos handle sort(), limit() and skip() when a find() fans out to several shards?

level: middleimportance: should knowfreq 50%

basics

~20 s

Each shard sorts and limits its own results; mongos merge-sorts the pre-sorted streams and re-applies the limit. It cannot push skip down — it fetches unskipped results and skips while assembling, passing skip plus limit to the shards when both are present.

open as a page

When should you choose a hashed shard key in MongoDB instead of a ranged one, and what do you give up?

level: middleimportance: should knowfreq 68%

basics

~20 s

Choose hashed sharding when the natural key is monotonic or unevenly distributed and the workload is mostly equality lookups. You give up range targeting: a range filter on that field must be broadcast to every shard.

open as a page

If a MongoDB sharded cluster's config server replica set loses its majority, what stops working?

level: seniorimportance: should knowfreq 38%

basics

~20 s

Metadata writes stop: no range splits, no migrations or balancing, no sharding a new collection, no adding or removing shards, no zone changes. Reads and writes to existing sharded collections keep working while routers still hold a usable cached routing table.

open as a page

How would you pin European customers' documents to shards in Frankfurt using MongoDB zone sharding?

level: seniorimportance: should knowfreq 42%

basics

~20 s

Tag the Frankfurt shards with a zone name using sh.addShardToZone, then map the shard-key ranges that hold European rows to that zone with sh.updateZoneKeyRange. The balancer then migrates those ranges onto the zoned shards and keeps them there.

open as a page

What must an updateOne or deleteOne on a sharded collection include in its filter, and what changed in MongoDB 7.0?

level: seniorimportance: should knowfreq 38%

basics

~20 s

Before MongoDB 7.0, a single-document update, delete or findAndModify on a sharded collection had to include the shard key or _id in its filter, or it errored. MongoDB 7.0 lifted that: the cluster now locates the document itself, at the cost of a broadcast. Upserts still need the full shard key.

open as a page

How can a low-cardinality MongoDB shard key create data ranges that can never be split or migrated?

level: seniorimportance: should knowfreq 55%

basics

~20 s

MongoDB can only cut the key space between two distinct shard-key values. A range holding just one distinct value has no split point, so it grows without bound, is flagged jumbo, and the balancer normally skips it.

open as a page

How do refineCollectionShardKey and reshardCollection differ when a MongoDB shard key proves wrong?

level: seniorimportance: should knowfreq 45%

basics

~20 s

refineCollectionShardKey only appends suffix fields, keeping the current key as a prefix, and is a metadata change that moves no data. reshardCollection replaces the key entirely by rewriting the collection, at real disk, I/O and cut-over cost.

open as a page

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%

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.

open as a page

How would you choose a shard key for a multi-tenant MongoDB collection where one tenant holds most of the data?

level: principalimportance: should knowfreq 40%

basics

~20 s

A bare tenantId fails on frequency: the dominant tenant's documents all share one key value and cannot be divided. Use a compound key such as tenantId plus a high-cardinality field, and validate the candidate against real data first.

open as a page

How is an aggregation pipeline split between the shards and a merging node in a sharded cluster?

level: seniorimportance: nice to knowfreq 33%

basics

~20 s

MongoDB splits the pipeline into a shards part that runs on every participating shard and a merger part that runs on one node. explain() shows splitPipeline with shardsPart, mergerPart and mergeType; stages needing an unsharded collection, such as $out, force the merge onto the database's primary shard.

open as a page