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.

Intermediate–Advanced130–210 minutesShard-key/distribution engineering labMongoDB 8.3.8 · mongosh 2.10.0 · PyMongo 4.17.0Last reviewed: September 2026

Learning objectives

01

Evaluate candidate shard keys from cardinality, per-value frequency, monotonicity, query isolation, and write-distribution evidence rather than one scalar score.

02

Explain how a high-cardinality key can still be bad because of skew, monotonic growth, or poor query targeting.

03

Use a deterministic AtlasMart workload to compare prospective keys before changing topology.

04

Demonstrate the hot-MaxKey behavior of a monotonically increasing ranged shard key in a disposable cluster.

05

Identify balancing/headroom and future-zone constraints that should be documented before choosing a production key.

Reproducible lab baseline

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.
No single “best shard-key score”

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.

candidate-key evidence from deterministic data and query samples
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

start compact sharded topology (l1)
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});' 
shard a disposable collection on a monotonically increasing key
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.

Evidence boundary

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.

Wrong approach: hash everything

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.

cleanup / full reset
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

  1. Why is high cardinality not sufficient?
  2. Why does an increasing ranged key hotspot?
  3. Does hashing solve a three-value region key?
  4. Is tenantId in the shard key an authorization boundary?
  5. 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

Keep knowledge open

Help the academy stay free and grow.

If these tutorials save you time, a small donation supports new lessons, technical review, diagrams, examples, and long-term maintenance.

ETHEthereum / ERC-20 only
0x716c4Ab160C4B66F31a28AE2448BfF68fc3a2ef0

Send only Ethereum or ERC-20 compatible assets to this address.