How would you choose a shard key for a multi-tenant MongoDB collection where one tenant holds most of the data?
answer
- The obvious key is the one that fails
- Many values, one of them dominant
- Hashing keeps that tenant together too
- Add something fine-grained after the tenant
- Measure the candidate before committing
basics
~20 sA 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.
solid answer
~50 s`{ tenantId: 1 }` looks appealing because every query carries the tenant, but the whale tenant's documents all share a single shard-key value, so they form one indivisible mass on one shard — high cardinality overall does not save you when one value dominates. Hashing the tenant does not fix it either: all of that tenant's documents hash to the same value and stay together. The workable answer is a compound key, `{ tenantId: 1, _id: 1 }` or `{ tenantId: 1, orderId: 1 }`, which keeps small tenants co-located and targetable while letting the whale's data split across shards on the suffix. Before committing, measure: `analyzeShardKey` (MongoDB 7.0+) reports cardinality, frequency and monotonicity for a candidate key against your actual data, and query sampling shows how often it would be targetable. If one tenant genuinely dwarfs the cluster, consider isolating it in its own collection or cluster instead.
code
javascript · 5 lines// Fails on frequency: the whale tenant's documents share one key value
sh.shardCollection("saas.orders", { tenantId: 1 })
// Works: tenant prefix keeps queries routable, suffix lets the whale split
sh.shardCollection("saas.orders", { tenantId: 1, _id: 1 })go deeper
Know that in a multi-tenant collection the tenant identifier alone is not automatically a good shard key, because one very large customer's documents all carry the same value and cannot be spread out.
Be able to explain the frequency failure and the compound remedy: tenant identifier as the prefix for targeting, plus a high-cardinality suffix so a large tenant's data can be divided.
Interviewers expect you to reject hashing the tenant with the right reason, quantify what the compound key costs the largest tenant's queries, and describe how you would measure a candidate against real data before sharding.
Own the decision framing: this key is expensive to change, so argue for measurement over intuition, and be ready to say when the right answer is isolating the largest tenants into their own deployments rather than distributing them.
## The trap in the obvious key Every query in a multi-tenant application carries the tenant identifier, so `{ tenantId: 1 }` seems to satisfy the targeting criterion perfectly, and with tens of thousands of tenants the cardinality looks excellent. The criterion it fails is frequency. Range boundaries are shard-key values, so the only place MongoDB can cut is between two distinct values. Every document belonging to the dominant tenant carries the identical `tenantId`, so they all sit inside one range with no split point. That range grows past the configured size, is flagged jumbo, and the balancer will not migrate it. One shard's storage grows without limit while others stay level, and no amount of extra shards changes that, because there is no piece to assign to them. This is worth stating plainly in an interview: cardinality and frequency are separate criteria and they fail the same way. A million distinct values with one value holding 60% of the documents behaves, for that 60%, exactly like a boolean shard key. ## Why hashing the tenant does not help The reflex fix for skew is a hashed key, but hashing is deterministic. `{ tenantId: "hashed" }` maps a given tenant to exactly one hash value, so the whale's documents still share one shard-key value and remain one indivisible mass. Hashing spreads *tenants* across shards evenly; it does nothing about a single tenant that is larger than a shard. Candidates who reach for hashing here and stop have missed the mechanism. ## The compound answer The key that usually works is `{ tenantId: 1, <high-cardinality field>: 1 }` — most often the document `_id`, or a natural unique id such as `orderId`. The prefix keeps a tenant's data contiguous, so a query filtering on `tenantId` alone still narrows to the ranges covering that tenant rather than broadcasting. The suffix gives the range map boundaries inside a tenant, so the whale's documents can be divided into as many ranges as needed and spread across every shard. Small tenants still occupy a single range and stay co-located, which is what you want for their read latency. The price is that the whale's tenant-only queries now touch several shards instead of one, which is unavoidable — its data no longer fits on one shard by definition. If the workload has a natural sub-scope, using it as the middle field, as in `{ tenantId: 1, projectId: 1, _id: 1 }`, restores targeted access at the level the application actually queries. A compound hashed variant, `{ tenantId: 1, sessionId: "hashed" }`, is available in MongoDB 4.4 and later and is useful when the suffix itself is clumpy or monotonic. ## Measuring before committing The reason this question sits at a lead level is that the decision is expensive to revisit: changing the leading field later means `reshardCollection`, which rewrites the collection with real disk, I/O and a write-blocking cut-over. So the professional move is to measure rather than reason. Since MongoDB 7.0, `analyzeShardKey` evaluates a candidate key against the collection's actual data and reports cardinality, the frequency distribution of key values, and whether the key is monotonically changing. Running it on `{ tenantId: 1 }` surfaces the whale immediately as a most-frequent value; running it on `{ tenantId: 1, _id: 1 }` shows a flat distribution. Pairing it with query sampling, configured via `configureQueryAnalyzer`, tells you what proportion of your real read and write traffic would be routable with each candidate — which is the targeting criterion measured instead of assumed. ## When the answer is not a shard key at all Sometimes the whale is a business fact rather than a modelling error, and the right architectural response is isolation rather than distribution. Giving the largest tenants their own collection, database or cluster removes their load from the shared cluster entirely, makes their capacity planning independent, and lets you offer them different backup, restore and upgrade schedules. It costs operational complexity — a routing layer that knows which tenants live where, and a migration path when a tenant graduates — but it turns a permanent skew problem into a provisioning decision. Zone-based placement is another lever for pinning particular tenants to particular shards. Naming isolation as an option, and being able to say when its complexity is justified, is the part of the answer that distinguishes a lead from a strong senior. ## What to say A complete answer moves through four beats: why `tenantId` alone fails on frequency and why hashing it does not help; the compound key that fixes it and what it costs the whale's queries; how you would validate the candidate against real data and real traffic before committing; and the point at which you stop trying to make one cluster hold every tenant and isolate the largest ones instead.
- Why doesn't hashing tenantId spread the dominant tenant's documents?Hashing is deterministic: one tenant id produces one hash value, so all of that tenant's documents still share a single shard-key value and sit in one indivisible range. Hashing distributes tenants relative to each other, which helps when tenant ids are clumpy or sequential, but it cannot divide a tenant that is larger than a shard.
- What does a tenant-only query cost once you shard on { tenantId: 1, _id: 1 }?For a small tenant, nothing meaningful — its documents fit in one range on one shard, so the query is still targeted. For the whale, its data now spans several shards, so a tenant-wide query touches all of them and the router merges results. That is inherent once one tenant exceeds a shard; the mitigation is to query at a narrower scope, such as adding a project or date dimension the key can exploit.
- How would you validate a candidate shard key against production data before sharding?Run `analyzeShardKey` (MongoDB 7.0+) on the candidate: it reports cardinality, the frequency distribution of key values, and whether the key changes monotonically, all measured against the real collection. Pair it with query sampling via `configureQueryAnalyzer` to learn what share of actual traffic would be routable. That converts three of the four criteria from argument into measurement.
- When would you give a tenant its own cluster instead of solving this with a shard key?When one tenant's working set or throughput approaches what the shared cluster can give it, when its traffic patterns degrade other tenants' latency, or when it needs its own backup, restore and upgrade cadence. The cost is a routing layer that maps tenants to deployments and a migration path for graduating tenants — worth paying when the alternative is sizing the whole cluster for one customer.
Filing customer paperwork by company name works until one client fills a whole cabinet on its own — you cannot split a single name across cabinets. Filing by company name plus invoice number keeps each client's papers together while letting the biggest one spill into as many cabinets as it needs.
saying these in an interview costs you the question
- Picks tenantId alone because every query filters on it
- Believes hashing tenantId spreads a single large tenant
- Assumes many distinct tenants guarantees even distribution
- Plans to fix a bad choice later without pricing the rewrite
- Never considers isolating the largest tenant from the shared cluster