skip to content

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

level: seniorimportance: should knowfreq 55%

answer

  1. A boundary must sit between two distinct values
  2. Five possible values means at most five groups
  3. One dominant value fails the same way
  4. The balancer will not move an oversized indivisible range
  5. Only changing the key actually helps

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.

solid answer

~50 s

Range boundaries are shard key values, so a split is only possible *between* distinct values. If a shard key has few distinct values — an order `status`, a `countryCode`, a boolean flag — or if one value dominates, a range ends up containing exactly one distinct value and has nowhere to be cut. It keeps growing as documents with that value arrive, exceeds the configured range size, and gets marked jumbo; the balancer will not move such a range under normal settings, so the shard owning it grows without limit while the others stay level. You see it as permanent skew in `db.collection.getShardDistribution()` and as jumbo markers in `sh.status()`. Manual splitting cannot help, because there is no boundary to split on. The real fixes are `refineCollectionShardKey` to append a high-cardinality suffix, or `reshardCollection` if the leading field itself is wrong.

code

javascript · 9 lines
javascript
// Only a handful of distinct values: at most that many indivisible groups
sh.shardCollection("shop.orders", { status: 1 })

// Give the range map a splittable suffix (MongoDB 4.4+)
db.orders.createIndex({ status: 1, orderId: 1 })
db.adminCommand({
  refineCollectionShardKey: "shop.orders",
  key: { status: 1, orderId: 1 }
})

go deeper

for a junior

Know that MongoDB divides sharded data into ranges of key values and that it needs two different values to cut between. A field with only a few possible values therefore limits how far data can be divided.

for a middle

Be able to explain the jumbo flag: a range holding one distinct value cannot be split, grows past the size limit, and is skipped by the balancer, leaving one shard permanently larger.

for a senior

Interviewers expect the diagnosis path — getShardDistribution and sh.status — plus the judgment that neither extra shards nor manual splits help, and a correct pick between refineCollectionShardKey and reshardCollection.

for a principal

Own the prevention side: set the standard that candidate keys are evaluated on worst-case value frequency before sharding, since the remedy after the fact is a full collection rewrite with real disk, I/O and cut-over cost.

## Ranges are intervals of shard key values A sharded collection's data is described by a map of half-open intervals over the shard key: one range covers keys from value X up to but not including value Y, and one shard owns it. To split a range, MongoDB must pick a new boundary value that falls strictly inside the interval and actually separates documents. That is only possible where two distinct shard-key values exist. This is the mechanical reason cardinality is a shard-key criterion rather than a style preference: the number of distinct values is a hard ceiling on how many pieces the collection can be divided into. ## How a range becomes indivisible Suppose you shard an `orders` collection on `{ status: 1 }` and `status` takes five values. The key space can be divided into at most five groups, one per value. As the collection grows to hundreds of gigabytes, each group grows too, and MongoDB has no boundary available inside a group — every document in it carries the identical key. The range keeps growing past the configured range size, and MongoDB flags it as a jumbo range. Because a migration must move a whole range, and moving an enormous one is expensive and disruptive, the balancer will not move a jumbo range under default settings. So the shard owning `status: "shipped"` keeps accumulating data forever, while shards owning smaller statuses sit half empty and nothing rebalances. The same thing happens with high-cardinality keys that have a dominant value. Shard on `{ tenantId: 1 }` where one enterprise tenant is 60% of the data, and every other tenant is fine while the whale's documents form a single indivisible mass. Interviewers like this case because it shows that cardinality and frequency are separate criteria that fail the same way. ## Symptoms in production The visible symptoms are a shard whose storage keeps growing while the balancer reports nothing to do, migrations that are attempted and abandoned, and disk-pressure alerts on one node. `db.collection.getShardDistribution()` gives the per-shard document and byte counts that show the skew. `sh.status()` shows the range map and flags jumbo ranges. In recent MongoDB versions the balancer distributes based on data size per shard rather than counting ranges, but that changes nothing here — the indivisible range still cannot be broken up or moved, so the imbalance persists no matter what the balancer is optimising. ## Why manual intervention does not help A reasonable first instinct is to split the range by hand with `sh.splitAt()` or `sh.splitFind()`. It does not work: those commands need a boundary value, and there is no value strictly between the single value the range contains and itself. Adding shards does not help either — new shards are capacity that no piece of this range can ever be assigned to. Raising the range size just delays the flag. The cause is the key, so the fix has to change the key. ## Fix 1: refine the shard key Since MongoDB 4.4, `refineCollectionShardKey` lets you append one or more fields to an existing shard key, keeping the current key as a prefix. Going from `{ status: 1 }` to `{ status: 1, orderId: 1 }` gives the range map a suffix it can split on: the previously indivisible `status: "shipped"` group can now be divided at `orderId` boundaries. Two things to understand about refinement. First, you must create an index supporting the new, longer key before running the command. Second, it is essentially a metadata change: it does not rewrite documents or immediately redistribute anything. Existing range boundaries stay as they are, and the collection converges over time as new splits and migrations use the finer key. That makes it cheap, but slow to take effect, and it cannot help if the *prefix* is what is wrong. ## Fix 2: reshard the collection Since MongoDB 5.0, `reshardCollection` can change the shard key entirely, including its leading field. It works by rewriting the collection under the new key while the workload continues, then cutting over. That means it needs spare disk on the shards and I/O headroom for the copy, it takes time proportional to the collection size, and there is a brief write-blocking critical section near the end of the operation. It is the right tool when the whole key was misjudged, and the wrong tool when appending a suffix would do. ## The lesson interviewers are listening for The candidate who answers well connects three things: split boundaries are shard-key values; therefore low cardinality *or* a dominant value produces an indivisible range; therefore the balancer, extra shards and manual splits are all powerless and the only real remedies change the key. Naming `refineCollectionShardKey` and `reshardCollection`, and knowing which one applies when, is what separates a textbook answer from an operational one.

  • Can you fix an indivisible range by running sh.splitAt() manually?
    No. A manual split still needs a boundary value that separates documents, and inside a range where every document carries the same shard-key value there is none. Manual splitting only helps when distinct values exist but MongoDB has not split there yet. The real remedies are appending a high-cardinality suffix with `refineCollectionShardKey`, or `reshardCollection` when the leading field is the problem.
  • After refining { status: 1 } to { status: 1, orderId: 1 }, does the skew disappear immediately?
    No. Refinement is a metadata change: existing range boundaries are untouched and no documents move as a result of the command. What changes is that future splits can use the `orderId` suffix, so the previously indivisible group can now be divided and migrated over time. If you need the imbalance gone quickly, resharding is the tool that actually rewrites placement.
  • Does the newer data-size-based balancing behaviour make this problem go away?
    No. Balancing on data size rather than range counts changes what the balancer optimises for, not what it is able to move. An indivisible range is still a single unit that cannot be broken up, and moving it wholesale is exactly what the jumbo flag exists to prevent. The imbalance persists until the shard key changes.

saying these in an interview costs you the question

  • Suggests adding shards to relieve an indivisible oversized range
  • Thinks a manual split can divide a single-value range
  • Confuses the range size limit with the 16 MB document limit
  • Believes the balancer eventually moves jumbo ranges on its own
  • Reaches for resharding when appending a suffix would suffice

context