Prompt 17 · Lesson 01 · Evidence before distribution
Choose a Shard Key: Cardinality, Frequency, Monotonicity, Query Isolation, and Write Distribution
A shard key is both the data-placement key and the router contract. Measure its value distribution and request patterns before you make it permanent.
Learning objectives
Evaluate candidate shard keys from cardinality, per-value frequency, monotonicity, query isolation, and write-distribution evidence rather than one scalar score.
Explain how a high-cardinality key can still be bad because of skew, monotonic growth, or poor query targeting.
Use a deterministic AtlasMart workload to compare prospective keys before changing topology.
Demonstrate the hot-MaxKey behavior of a monotonically increasing ranged shard key in a disposable cluster.
Identify balancing/headroom and future-zone constraints that should be documented before choosing a production key.
This lesson pins MongoDB Community Server
8.3.8 with
mongodb/mongodb-community-server:8.3.8-ubuntu2204-slim, mongosh 2.10.0, and PyMongo
4.17.0 where the driver is used. The local topology
is a disposable single-host sharded cluster with one
single-member config server replica set, two
single-member shard replica sets, and one
mongos, exposed only on loopback diagnostic ports
27135–27138. This is sufficient for routing/distribution
mechanics but is not production high availability.
Authentication and TLS are disabled only for the isolated lab.
Feature Compatibility Version (FCV) is inspected in the lab. FCV
is observed and never changed. Application/DDL operations go
through mongos; direct shard connections are
diagnostics only. Lesson 1 includes a deterministic Python
analysis as the low-resource mandatory fallback. The real
topology then demonstrates one monotonically hot ranged key; the
small fixture is not a benchmark. Atlas, Search, Vector Search,
KMS, and Enterprise Advanced are not mandatory. Product commands
were not executed in this generation environment because Docker,
mongod, mongos, mongosh, and PyMongo are unavailable here, so
all timing/load outputs must be measured on the learner machine
rather than copied as invented results.
1. A shard key is simultaneously a placement key and a routing contract
AtlasMart has millions of orders and four plausible fields that
engineers could shard on: region,
tenantId, createdSeq, and the compound
pair {tenantId, orderId}. A
shard key is the indexed field pattern MongoDB
uses to divide a collection into non-overlapping ranges and
route operations. Choosing it is therefore not merely a storage
decision: it decides which writes become hot, which queries can
be isolated to a shard subset, which zone boundaries are
expressible, and how much future redistribution work is likely.
| Property | Question | Failure if ignored |
|---|---|---|
| Cardinality | How many distinct key values exist? | Too few values cap how finely data can be split and can create indivisible/large ranges. |
| Frequency | How much data/work sits behind the hottest values? | A high-frequency value can remain hot even when total cardinality is large. |
| Monotonicity | Do new keys trend toward MinKey or MaxKey? | Ranged inserts concentrate on the current edge range/shard. |
| Query isolation | Do hot request predicates contain the key or a compound-key prefix? | Requests scatter to more shards and merge work grows. |
| Write distribution | Which ranges receive inserts/updates under real traffic? | One shard can bottleneck cluster-wide write capacity. |
| Future zones | Can placement/compliance boundaries be expressed with shard-key fields/prefixes? | Later zone design may require resharding instead of a metadata-only change. |
Cardinality, frequency, monotonicity, query targeting, growth, and placement constraints trade off. A high-cardinality key is not automatically good; a hashed key is not automatically good; and a key that balances bytes can still scatter the latency-sensitive workload.
2. Quantify the candidates on a deterministic AtlasMart sample
The following simulator does not mimic MongoDB's internal
analyzeShardKey implementation. It computes
transparent workload statistics so you can see why each
candidate fails or succeeds before the server-specific analyzer
is introduced in Lesson 3.
from collections import Counter# 1,200 synthetic orders. tenant "mega" is intentionally heavy.docs = []for i in range(1, 1201): tenant = "mega" if i <= 420 else f"t{((i - 421) % 39) + 1:02d}" region = ["EU", "US", "APAC"][i % 3] docs.append({ "region": region, "tenantId": tenant, "createdSeq": i, "orderId": f"o-{i:05d}", })queries = ( [{"tenantId": "mega"}] * 50 + [{"tenantId": "t07"}] * 25 + [{"region": "EU"}] * 10 + [{"orderId": "o-00900"}] * 10 + [{"createdSeq": {"$gte": 1100}}] * 5)candidates = { "region": lambda d: (d["region"],), "tenantId": lambda d: (d["tenantId"],), "createdSeq": lambda d: (d["createdSeq"],), "tenantId+orderId": lambda d: (d["tenantId"], d["orderId"]),}for name, keyfn in candidates.items(): values = [keyfn(d) for d in docs] counts = Counter(values) top = counts.most_common(1)[0][1] print({ "candidate": name, "cardinality": len(counts), "maxFrequency": top, "maxFrequencyPct": round(top / len(docs) * 100, 1), })# createdSeq is exactly increasing by construction; this is a workload fact,# not a generic rule that all integer keys are monotonic.print({"createdSeqMonotonic": all(docs[i]["createdSeq"] < docs[i+1]["createdSeq"] for i in range(len(docs)-1))})print({"querySampleCount": len(queries)})
Expected deterministic evidence: region has
cardinality 3; createdSeq and
{tenantId,orderId} have cardinality 1,200; and
tenantId is highly skewed because
mega owns 420 documents. That still does not answer
routing: the query sample shows that tenant-scoped traffic
dominates, so a key that omits tenantId may
distribute bytes yet scatter common requests.
3. Boundary case: high cardinality plus monotonic growth can still hotspot
docker rm -f atlasmart-ch17-l1-cfg atlasmart-ch17-l1-s1 atlasmart-ch17-l1-s2 atlasmart-ch17-l1-mongos 2>/dev/null || truedocker network rm atlasmart-ch17-l1-net 2>/dev/null || truedocker volume rm atlasmart-ch17-l1-cfg-data atlasmart-ch17-l1-s1-data atlasmart-ch17-l1-s2-data 2>/dev/null || truedocker network create atlasmart-ch17-l1-netdocker run -d --name atlasmart-ch17-l1-cfg --network atlasmart-ch17-l1-net -p 127.0.0.1:27135:27017 -v atlasmart-ch17-l1-cfg-data:/data/db mongodb/mongodb-community-server:8.3.8-ubuntu2204-slim --configsvr --replSet atlasmart-cfg17-l1 --bind_ip_alldocker run -d --name atlasmart-ch17-l1-s1 --network atlasmart-ch17-l1-net -p 127.0.0.1:27136:27017 -v atlasmart-ch17-l1-s1-data:/data/db mongodb/mongodb-community-server:8.3.8-ubuntu2204-slim --shardsvr --replSet atlasmart-shard17-l1-a --bind_ip_alldocker run -d --name atlasmart-ch17-l1-s2 --network atlasmart-ch17-l1-net -p 127.0.0.1:27137:27017 -v atlasmart-ch17-l1-s2-data:/data/db mongodb/mongodb-community-server:8.3.8-ubuntu2204-slim --shardsvr --replSet atlasmart-shard17-l1-b --bind_ip_allfor PORT in 27135 27136 27137; do until mongosh "mongodb://127.0.0.1:$PORT/admin?directConnection=true" --quiet --eval 'quit(db.runCommand({ping:1}).ok===1?0:1)'; do sleep 1; donedonemongosh "mongodb://127.0.0.1:27135/admin?directConnection=true" --quiet --eval 'rs.initiate({_id:"atlasmart-cfg17-l1",configsvr:true,members:[{_id:0,host:"atlasmart-ch17-l1-cfg:27017"}]})'mongosh "mongodb://127.0.0.1:27136/admin?directConnection=true" --quiet --eval 'rs.initiate({_id:"atlasmart-shard17-l1-a",members:[{_id:0,host:"atlasmart-ch17-l1-s1:27017"}]})'mongosh "mongodb://127.0.0.1:27137/admin?directConnection=true" --quiet --eval 'rs.initiate({_id:"atlasmart-shard17-l1-b",members:[{_id:0,host:"atlasmart-ch17-l1-s2:27017"}]})'for PORT in 27135 27136 27137; do until mongosh "mongodb://127.0.0.1:$PORT/admin?directConnection=true" --quiet --eval 'quit(db.hello().isWritablePrimary?0:1)'; do sleep 1; donedonedocker run -d --name atlasmart-ch17-l1-mongos --network atlasmart-ch17-l1-net -p 127.0.0.1:27138:27017 --entrypoint mongos mongodb/mongodb-community-server:8.3.8-ubuntu2204-slim --configdb atlasmart-cfg17-l1/atlasmart-ch17-l1-cfg:27017 --bind_ip_all --port 27017until mongosh "mongodb://127.0.0.1:27138/admin" --quiet --eval 'quit(db.runCommand({ping:1}).ok===1?0:1)'; do sleep 1; donemongosh "mongodb://127.0.0.1:27138/admin" --quiet --eval 'printjson(sh.addShard("atlasmart-shard17-l1-a/atlasmart-ch17-l1-s1:27017"));printjson(sh.addShard("atlasmart-shard17-l1-b/atlasmart-ch17-l1-s2:27017"));printjson({hello:db.hello(),fcv:db.runCommand({getParameter:1,featureCompatibilityVersion:1}).featureCompatibilityVersion,shards:db.adminCommand({listShards:1}).shards});'
const app = db.getSiblingDB("atlasmart");sh.enableSharding("atlasmart", "atlasmart-shard17-l1-a");app.orders_hot.drop();app.orders_hot.createIndex({createdSeq:1});printjson(sh.shardCollection("atlasmart.orders_hot", {createdSeq:1}));for (let start=1; start<=1000; start+=100) { const batch=[]; for (let i=start;i<start+100;i++) batch.push({createdSeq:i,tenantId:`t${i%20}`,status:"paid"}); app.orders_hot.insertMany(batch);}printjson(sh.splitAt("atlasmart.orders_hot", {createdSeq:500}));printjson(sh.moveChunk("atlasmart.orders_hot", {createdSeq:900}, "atlasmart-shard17-l1-b"));const cfg=db.getSiblingDB("config");const c=cfg.collections.findOne({_id:"atlasmart.orders_hot"});printjson(cfg.chunks.find({uuid:c.uuid},{min:1,max:1,shard:1}).sort({min:1}).toArray());// Every appended key is above the current historical maximum.const tail=[];for (let i=1001;i<=1100;i++) tail.push({createdSeq:i,tenantId:`t${i%20}`,status:"paid"});app.orders_hot.insertMany(tail);printjson(app.orders_hot.find({createdSeq:{$gte:1001}}).explain("executionStats"));
The added 100 writes all fall in the range whose upper bound is
MaxKey. On a ranged monotonic key, that edge range
becomes the current write hotspot. The balancer may split/move
edge ranges over time, but it competes with the continuing
inserts; this is not the same as evenly distributing each
incoming write.
This fixture proves routing concentration for the declared range map. It does not prove a specific CPU percentage, throughput ceiling, migration rate, or p99 latency. Record those on the learner machine with the same dataset and concurrency before making capacity claims.
4. Frequency and “jumbo” reasoning require distribution evidence
A frequent shard-key value constrains split flexibility because identical values cannot be split across arbitrary ranged-key boundaries. Modern MongoDB can move large ranges in more situations than older folklore suggests, so do not diagnose every large range as “jumbo.” Instead inspect range size, key frequency, balancer status, migration errors, and current version behavior. If one tenant dominates, adding a suffix or changing the key may create more split points; hashing can distribute a monotonic field, but it may damage range targeting.
Hashing does not manufacture cardinality. Three region values remain three logical values before hashing, and a single very frequent value still represents a large amount of data. Hashed distribution is especially useful for high-cardinality monotonic keys, but range queries on the original value become poor targeting candidates.
5. Production judgment: write the decision record before the command
A production shard-key review should attach evidence:
distinct-value counts and top frequencies, write distribution
over time, the percentage of hot queries that contain the
key/prefix, expected zone boundaries, per-shard disk/RAM/CPU
headroom, balancer windows, backup/restore implications, and a
rollback/evolution option. Include tenant-isolation/security
separately: a shard key containing tenantId is a
routing mechanism, not authorization.
Bridge. Lesson 2 compares ranged, hashed, and compound shard keys directly so you can see how distribution and prefix-aware targeting pull in different directions.
docker rm -f atlasmart-ch17-l1-cfg atlasmart-ch17-l1-s1 atlasmart-ch17-l1-s2 atlasmart-ch17-l1-mongos 2>/dev/null || truedocker volume rm atlasmart-ch17-l1-cfg-data atlasmart-ch17-l1-s1-data atlasmart-ch17-l1-s2-data 2>/dev/null || truedocker network rm atlasmart-ch17-l1-net 2>/dev/null || true
Check your understanding
- Why is high cardinality not sufficient?
- Why does an increasing ranged key hotspot?
- Does hashing solve a three-value region key?
- Is tenantId in the shard key an authorization boundary?
- What should accompany a shard-key choice?
Review the answers
1. Because frequency skew, monotonic growth, query targeting, and placement requirements can still create hotspots or scatter/gather.
2. New values fall into the current MaxKey range, so incoming writes concentrate on whichever shard owns that edge range.
3. No. Hashing does not create new logical distinct values or remove a very frequent value.
4. No. Authorization must be enforced independently through application and MongoDB security controls.
5. Measured key/workload distributions, targeting evidence, capacity/headroom, zone constraints, and an evolution/rollback plan.
Authoritative references
- Choose a Shard Key — Cardinality, frequency, monotonicity, query patterns, and shard-key tradeoffs.
- Troubleshoot Shard Keys — Hot ranges, uneven load, jumbo/indivisible ranges, and scatter/gather symptoms.
- analyzeShardKey — Key-characteristic and sampled read/write-distribution metrics.
- configureQueryAnalyzer — Query sampling for prospective shard-key analysis.
- Hashed Indexes for Sharding — Write distribution and range-query limitations.
- Zones — Zone membership, non-overlapping zone ranges, prefix requirements, and balancing behavior.
- addShardToZone — Associate shards with zone labels.
- Change a Shard Key — Refine versus reshard decision boundary.
- Refine a Shard Key — Append suffix fields without changing existing shard-key field types.
- reshardCollection — Full shard-key change, redistribution phases, demo mode, and abort/commit boundaries.
- unshardCollection — MongoDB 8.0+ consolidation of a sharded collection onto one shard.
- moveCollection — MongoDB 8.0+ relocation of an unsharded collection between shards.
- MongoDB Sharding — Routing, ranges, balancer behavior, and mongos-only client access.
- MongoDB 8.3 Release Notes — Current 8.3 release and sharding/DDL changes.
- mongosh Release Notes — mongosh 2.10.0 baseline.
- PyMongo Release Notes — PyMongo 4.17 baseline.