skip to content

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

level: middleimportance: must knowfreq 78%

answer

  1. Three data properties plus one workload property
  2. How many pieces can the data split into
  3. When one value holds most documents
  4. Values that only ever increase
  5. Does the filter carry the key at all

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.

solid answer

~50 s

MongoDB judges a shard key on three data properties plus one workload property. **Cardinality** bounds how finely the key space can be divided: a range holding a single shard-key value can never be split, so few distinct values caps how many pieces the collection can break into. **Frequency** is how evenly documents spread across those values — a key with a million distinct values is still bad if most documents carry one of them. **Monotonicity** matters because ranged sharding hands contiguous value ranges to shards, so an always-increasing key such as an ObjectId `_id` or a timestamp sends every new insert into the one range at the top. The fourth criterion is targeting: if the key is absent from common filters, `mongos` must broadcast. You want high cardinality, even spread, no monotonic growth, and presence in the read path.

code

javascript · 8 lines
javascript
// Poor: few distinct values, ranges cannot be split further
sh.shardCollection("shop.orders", { status: 1 })

// Poor: ObjectId increases over time, all inserts hit the top range
sh.shardCollection("shop.orders", { _id: 1 })

// Better: high cardinality, splittable suffix, carried by common filters
sh.shardCollection("shop.orders", { customerId: 1, orderId: 1 })

go deeper

for a junior

Know that a sharded collection places documents by one chosen key, and that the choice decides which shard stores and serves a document. Be able to say why a field with only a few values is a poor candidate.

for a middle

Be ready to define cardinality, frequency and monotonicity in your own words and give a concrete bad key for each, plus the index requirement that the shard key must be a prefix of its supporting index.

for a senior

Interviewers expect you to start from access patterns, reason about the worst-case value frequency rather than averages, and explain why unique constraints and key mutability constrain the choice in production.

for a principal

Own the framing that shard key choice is a one-way door relative to its cost: be able to explain how you would evaluate candidates against real data and traffic before committing, and what the exit path costs if you get it wrong.

## What the shard key actually controls In a sharded cluster, one collection's documents are divided into contiguous ranges of shard key values, and each range is owned by exactly one shard. The shard key is the only input to that placement decision. It determines where a document lands, which shards a query must visit, and how evenly write traffic spreads across the cluster. The router process, `mongos`, reads the range-to-shard map from the config servers; a query whose filter carries the shard key can be sent only to the shards owning the matching ranges, while a query without it must be broadcast to all of them. Because placement derives from this one key, a bad choice is not something more hardware, a bigger balancer, or better indexes can repair. ## Cardinality — how finely the space can be divided Cardinality is the number of distinct values the key takes. MongoDB can only cut the key space *between* two distinct values, so cardinality is a hard ceiling on how many pieces the data can be split into. If you shard on a `status` field with five values, the collection can never be divided into more than five indivisible groups no matter how many terabytes it holds or how many shards you add. A range that ends up containing exactly one distinct shard-key value cannot be split at all; if it grows large, it becomes a permanently oversized range that the balancer will not move, and one shard grows without bound. ## Frequency — how evenly documents spread over the values Cardinality alone is not enough. Frequency asks how many documents sit on each distinct value. Consider a `customerId` key in a business where one enterprise account produces 70% of all orders: cardinality is in the millions, but that one value's documents all share a single key value, so they cannot be divided among shards. The result is the same skew a low-cardinality key produces, just concentrated on one value instead of spread over five. In practice you look at the distribution, not just the count: the worst case is the *most frequent* value, because its documents are inseparable. ## Monotonicity — the direction values move over time Ranged sharding assigns ordered, contiguous ranges. The topmost range has an upper bound of `MaxKey` and is owned by exactly one shard. If the shard key increases with every insert — a timestamp, an auto-incrementing counter, or a default `ObjectId`, whose leading bytes encode the creation time — then every new document falls into that top range, and one shard absorbs 100% of the insert load while the rest sit idle. The balancer reacts by continuously migrating ranges off that shard, which costs I/O without fixing the cause. Monotonically *decreasing* keys have the same problem at the bottom of the range. Hashed sharding and non-monotonic compound prefixes are the standard answers. ## Query targeting — the workload half of the decision The first three properties are about the data; the fourth is about your queries. A key that spreads perfectly but never appears in a filter forces every read to fan out to every shard, so cluster-wide read capacity stops scaling with shard count and tail latency becomes the slowest shard's latency. That is why shard key selection starts with a list of the collection's real access patterns, not with the schema. Frequently the best key is a compromise: something coarse that queries always carry, plus something fine that makes ranges splittable. ## Mechanical requirements A few rules constrain the choice regardless of the data. A supporting index must exist: a ranged key needs an index with the shard key as a prefix, and a hashed key needs a hashed index. MongoDB creates it for you when you shard an empty collection, but on a non-empty collection you must build it first. Any unique index on a sharded collection must be prefixed by the shard key, because each shard can only enforce uniqueness over the documents it holds — so you cannot add an arbitrary unique constraint after sharding. Shard key fields are also effectively business-immutable: changing one is possible in currently supported versions, but only inside a transaction or as a retryable write, so schemas that mutate the key field constantly are a bad fit. ## Putting it together For an `orders` collection queried mostly as "all orders for this customer" and occasionally as "this specific order", `{ status: 1 }` fails on cardinality, `{ createdAt: 1 }` fails on monotonicity, and `{ orderId: 1 }` spreads well but leaves customer queries broadcasting. `{ customerId: 1, orderId: 1 }` gives high cardinality, keeps a large customer's documents splittable by the `orderId` suffix, is non-monotonic in its prefix, and lets customer-scoped queries target a narrow set of ranges. That reasoning — access patterns first, then cardinality, frequency and monotonicity — is what an interviewer is listening for.

  • Does the shard key have to be indexed, and what happens to unique indexes on a sharded collection?
    Yes. A ranged shard key needs an index with the key as a prefix; a hashed key needs a hashed index. MongoDB creates it when you shard an empty collection but requires it up front on a non-empty one. Unique indexes are constrained too: on a sharded collection a unique index must be prefixed by the shard key, because each shard can only enforce uniqueness over the documents it stores.
  • Can you change a document's shard key value after it is inserted?
    Yes in currently supported versions, but deliberately awkwardly: an update that changes shard key values must run inside a multi-document transaction or as a retryable write, because the document may have to move to a different shard. Treat shard key fields as effectively immutable business data — a field that changes routinely is a poor shard key even if its distribution looks good.
  • Why can't you fix a bad shard key by just adding more shards?
    Placement is a pure function of the shard key. If a key has five distinct values, or one value holds most documents, adding shards creates capacity that no range can ever be assigned to — the indivisible groups still live on one shard each. More shards only help when the key space can actually be cut into more pieces.

Think of the shard key as the label you sort mail by. Sorting by country puts most letters in one bin (low cardinality); sorting by postmark time means every new letter lands in today's bin (monotonic); sorting by recipient plus letter number spreads the work and still lets you find one person's mail quickly.

saying these in an interview costs you the question

  • Assumes any unique field, like _id, is automatically a good shard key
  • Treats high cardinality alone as sufficient and ignores value frequency
  • Believes the balancer will eventually fix a badly chosen shard key
  • Confuses the shard key with the primary key _id
  • Picks the key from the schema without looking at query filters

context