skip to content

Chunks, Balancing & Zones

How data physically moves between shards to keep them even, and how you pin data to specific hardware or regions. Comes up in data-residency and tiered-storage design questions.

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

questions

6

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

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

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

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

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