
How to Choose a Good Shard Key in MongoDB
You can add shards to a MongoDB cluster in an afternoon. What you can't easily fix is a bad shard key. Pick one that sends every new write to the same place, and you've built a cluster where one shard works at 100% while the others idle. Pick one your queries don't use, and every read becomes a broadcast to every shard. Either way, you've paid for horizontal scale and gotten none of it.
The shard key is the field (or fields) MongoDB uses to decide which shard each document lives on. It determines how writes are distributed, which queries can be routed to a single shard, and whether data can be split into small enough pieces to balance. Recent versions let you refine or completely reshard a collection, so the choice is no longer permanent, but resharding a multi-terabyte collection is slow and expensive. It's worth getting right up front.
This guide covers the properties that make a shard key good, the common patterns and anti-patterns, worked examples for typical workloads, how to measure candidate keys with analyzeShardKey, and the constraints the shard key places on your schema. If you're new to how sharded clusters are put together, start with MongoDB Sharding Explained.
The Four Properties of a Good Shard Key
A shard key is judged on four things. You'll rarely maximize all of them, so the goal is to understand the trade-offs for your workload.
1. High Cardinality
Cardinality is the number of distinct shard key values. MongoDB groups documents into chunks by key range, and a chunk can never be split between two documents with the same key value. So the number of distinct values caps how finely the data can be divided.
Shard on country, and you have maybe 200 values. The chunk for "US" might hold 40% of your data, and it can never be split. It becomes a jumbo chunk: too big to move, stuck on one shard forever.
Shard on userId, and you have millions of values. The balancer can split and move ranges freely.
2. Low Frequency
Frequency is how often a single value appears. Even with high cardinality, if one value is extremely common, all documents with that value land in one unsplittable chunk. A customerId key sounds great until you learn one enterprise customer accounts for 30% of all orders.
3. Non-Monotonic Values
A key is monotonic if new documents always have a larger (or smaller) value than existing ones. Timestamps, auto-incrementing counters, and ObjectIds are all monotonic. With a ranged shard key, every new document belongs to the chunk covering the highest values, which lives on one shard. That shard takes 100% of inserts until the balancer eventually moves things around, and then the next chunk at the top becomes the new hot spot.
4. Query Targeting
A query that includes the shard key (or a prefix of a compound key) can be routed to exactly the shards that own matching data. A query without it must be sent to every shard. The best shard key is one that your most frequent, latency-sensitive queries already filter on.
| Property | Question to ask | Bad sign |
|---|---|---|
| Cardinality | How many distinct values will exist? | Hundreds or fewer |
| Frequency | Does any single value hold a big share of documents? | One value is more than a few percent |
| Monotonicity | Do new values always increase? | Timestamps, counters, ObjectIds (ranged) |
| Targeting | Do my hottest queries filter on this field? | Most reads don't include it |
Ranged, Hashed, and Compound Keys
The type of key you use changes which properties you get.
Ranged keys keep documents with nearby values together. Queries for a range of key values touch few shards. They're a good fit when the leading field has high cardinality and isn't monotonic:
sh.shardCollection("app.users", { email: 1 });
Hashed keys store a hash of the value, so consecutive values scatter across chunks. This fixes monotonic keys completely: inserts with increasing ObjectIds spread evenly across all shards. The cost is that range queries on the key must go to every shard, because neighbouring values are no longer stored together.
sh.shardCollection("logs.requests", { _id: "hashed" });
Compound keys combine fields to get the benefits of both. The most useful pattern is a leading field that groups data your queries need together, followed by a field that adds cardinality or distribution:
// All of a tenant's documents are grouped, but a big tenant is still splittable
sh.shardCollection("saas.records", { tenantId: 1, _id: 1 });
// Group by device, spread writes within a device by hashing the time
sh.shardCollection("iot.readings", { deviceId: 1, ts: "hashed" });
A compound key can include at most one hashed field. Queries are targeted when they include a prefix of the key, so { tenantId: "t_42" } is targeted with the first example, while { _id: someId } alone is not.
Worked Examples
Abstract rules are easier to apply with real workloads in front of you.
Multi-Tenant SaaS
Every query in a typical SaaS app includes the tenant, and tenant sizes vary wildly.
{ tenantId: 1 }alone: excellent targeting, but a large tenant becomes a jumbo chunk.{ tenantId: "hashed" }: even distribution of tenants, but still one chunk per tenant value.{ tenantId: 1, _id: 1 }: targeting by tenant, and large tenants can be split into many chunks across shards.
The compound key is the usual winner. Small tenants stay on one shard (so their queries touch one shard), while large tenants spread across several. For more on tenancy patterns, see Multi-Tenant Application Design Patterns in MongoDB.
E-Commerce Orders
Access patterns: "show this customer's orders" (very frequent), "look up an order by ID" (frequent), "orders placed today across all customers" (reporting, occasional).
sh.shardCollection("shop.orders", { customerId: 1, orderDate: 1 });
Customer order history is targeted and sorted by date within the shard. Lookup by order ID alone would be scatter-gather, so include customerId in order URLs and API calls, or keep a small unsharded lookup collection mapping order IDs to customers. The daily report becomes a scatter-gather query, which is fine for an occasional analytics job.
Event Logs and Telemetry
Access pattern: heavy insert volume, reads mostly for recent events of a single device or source.
// Bad: every insert goes to the "latest" chunk on one shard
sh.shardCollection("iot.events", { ts: 1 });
// Better: writes spread by device; a device's events stay together and ordered
sh.shardCollection("iot.events", { deviceId: 1, ts: 1 });
With { deviceId: 1, ts: 1 }, each device's inserts land in that device's range. As long as you have many devices, writes spread across shards, and "last hour for device X" is a targeted range scan. If you have only a handful of very chatty devices, switch the second field to hashed or add a bucket field. For raw time-series ingestion, also consider time series collections, which can be sharded too.
User Profiles
Access pattern: lookup by user ID on every request, rare range queries.
sh.shardCollection("app.profiles", { _id: "hashed" });
If user IDs are ObjectIds (monotonic), hashing is essential. Point lookups by _id are targeted because MongoDB hashes the query value to find the right chunk.
Chat Messages
Access pattern: load the latest messages for a conversation, paginate backwards.
sh.shardCollection("chat.messages", { conversationId: 1, sentAt: -1 });
Every read is scoped to a conversation and wants messages in time order, which this key provides on a single shard. Extremely large group conversations can still be split by sentAt ranges.
Measuring Candidates with analyzeShardKey
Since MongoDB 7.0, you don't have to guess. The analyzeShardKey command evaluates a candidate key against your real data and, optionally, your real query traffic. It works on both unsharded and sharded collections, so you can run it before you shard.
First, if you want read and write distribution metrics, turn on query sampling for the collection so MongoDB records a sample of real operations:
db.adminCommand({
configureQueryAnalyzer: "shop.orders",
mode: "full",
samplesPerSecond: 10,
});
Let it run through a representative period (a full business day is ideal), then analyze a candidate. The collection needs an index that starts with the candidate key for the key characteristics part:
db.adminCommand({
analyzeShardKey: "shop.orders",
key: { customerId: 1, orderDate: 1 },
keyCharacteristics: true,
readWriteDistribution: true,
});
The output (trimmed) tells you most of what you need:
{
keyCharacteristics: {
numDocsTotal: 150012877,
numDistinctShardKeyValues: 149870112,
mostCommonValues: [
{ value: { customerId: "c_ent_001", orderDate: ISODate("2026-09-01") }, frequency: 42 },
// ...
],
monotonicity: { type: "not monotonic" },
avgDocSizeBytes: 812
},
readDistribution: {
sampleSize: { total: 861203, find: 790112, aggregate: 71091 },
percentageOfSingleShardReads: 91.4,
percentageOfMultiShardReads: 2.1,
percentageOfScatterGatherReads: 6.5
},
writeDistribution: {
sampleSize: { total: 211480, update: 150902, delete: 3105, findAndModify: 57473 },
percentageOfSingleShardWrites: 98.7,
percentageOfScatterGatherWrites: 1.3,
percentageOfShardKeyUpdates: 0.2
}
}
How to read it:
numDistinctShardKeyValuesclose tonumDocsTotalmeans high cardinality.mostCommonValuesreveals frequency problems. If the top value covers a big fraction of documents, you'll get jumbo chunks.monotonicity.typeflags keys that would create an insert hot spot under ranged sharding.percentageOfSingleShardReadstells you how well the key targets your actual queries. Over 90% is a strong result.percentageOfShardKeyUpdatesshows how often writes would change a document's shard key value, which is expensive because the document has to move between shards.
Run it for two or three candidates and compare. Turn sampling off when you're done:
db.adminCommand({ configureQueryAnalyzer: "shop.orders", mode: "off" });
Constraints the Shard Key Puts on Your Schema
Once a collection is sharded, the key affects more than distribution.
The key must be indexed. Every shard key needs an index that starts with the key fields. On an empty collection, shardCollection creates it. On a populated collection, create it first.
Unique indexes must include the shard key. MongoDB can only enforce uniqueness per shard, so a unique index on a sharded collection must be prefixed by the shard key. If you shard users on { _id: "hashed" }, you can't have a unique index on email alone. The usual workaround is a separate, unsharded (or differently sharded) collection that holds the unique values, written together with the main document.
Updating a key value moves the document. Since MongoDB 4.2, you can change a document's shard key value (unless the key includes _id), but the update must run as a retryable write or in a transaction, and it physically moves the document to another shard. If your key changes often, it's the wrong key.
Documents can lack key fields. Since 4.4, documents missing a shard key field are treated as having null for it. That's convenient for migrations but dangerous if many documents end up with null, since they all fall into one chunk.
Include the key in operations for efficiency. Single-document updates and deletes that include the shard key in the filter are routed to one shard. Recent versions can run updateOne without it, but the router must do extra work to find the document first.
Fixing a Key That Went Wrong
If you're already stuck with a poor key, you have two options.
Refine it by adding suffix fields. This is cheap (it's mostly a metadata change) and fixes low-cardinality keys:
db.orders.createIndex({ customerId: 1, orderId: 1 });
db.adminCommand({
refineCollectionShardKey: "shop.orders",
key: { customerId: 1, orderId: 1 },
});
Reshard it when the leading field itself is wrong. reshardCollection copies the data into a new layout while the collection stays online, then cuts over. It needs roughly enough free space to hold another copy of the collection and generates heavy I/O, so plan it for a quiet period. In recent versions, resharding has become significantly faster, but it's still a major operation for large collections.
Common Mistakes
Sharding on _id with ranged sharding. Default ObjectIds are monotonic, so a ranged _id key sends every insert to one shard. Use { _id: "hashed" } or a different key.
Picking the key by data distribution alone. A hashed _id distributes perfectly, but if every query filters by customerId, every query becomes scatter-gather. Balance distribution with targeting.
Using a low-cardinality field as the whole key. status, type, country, and booleans are fine as a prefix in a compound key, but never on their own.
Ignoring the biggest tenant or customer. Averages hide outliers. Look at mostCommonValues, or run a $group by the candidate field sorted by count, and plan for the largest value.
Not sampling real queries. The access patterns you imagine and the ones production runs are often different. Run configureQueryAnalyzer for a representative period before deciding.
Forgetting about unique constraints. Discovering after sharding that you can no longer enforce unique emails is painful. List every unique index before choosing a key.
Conclusion
A good shard key has high cardinality, no dominant values, doesn't increase monotonically (or is hashed if it does), and appears in the queries that matter most. Compound keys like { tenantId: 1, _id: 1 } or { deviceId: 1, ts: 1 } often give you targeting and distribution together. Measure candidates against real data and traffic with analyzeShardKey, account for unique indexes before you commit, and use refineCollectionShardKey or reshardCollection if the workload changes under you.
Your next step: pick the collection you're most likely to shard, turn on configureQueryAnalyzer for a day, and run analyzeShardKey on your top two candidate keys. Put the single-shard read percentages side by side, and the decision usually makes itself.


