What is a chunk in a sharded MongoDB collection, and what does the balancer do with chunks?
answer
- contiguous span of shard-key values
- half-open interval, MinKey to MaxKey
- exactly one owning shard, not copies
- config servers hold the range map
- 6.0: balances by data size, no auto-splitter
basics
~20 sA 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 sSharding 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 linessh.status()
sh.getBalancerState() // is balancing enabled?
sh.isBalancerRunning() // is a round active right now?
db.adminCommand({ balancerCollectionStatus: "app.events" })go deeper
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.
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.
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.
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