Lesson 0009 · Horizontal scale
When one replica set isn't enough — data volume or write throughput exceeds single-machine capacity. MongoDB's answer is built in; Postgres's isn't. The shard key is the most consequential schema decision you'll make, and a bad one can't be undone without downtime.
j: true and w: "majority" — how are they different?
With default w: 1, when can an acknowledged write be silently rolled back?
readConcern: "majority" guarantees the data you read:
A replica set gives you redundancy and read scale — you can spread reads across secondaries. But every member holds a full copy of all data, and all writes go to one primary. When the data outgrows one machine's storage, or when write throughput exceeds one machine's capacity, the replica set model hits a ceiling. Sharding partitions the data horizontally across multiple replica sets, each holding a slice of the total.
Postgres has no built-in equivalent. CitusDB (now a Microsoft extension) adds it externally; Vitess does the same for MySQL. In MongoDB it's a first-class cluster topology — no third-party tool required. — MongoDB Manual: "Sharding is a method for distributing data across multiple machines. MongoDB uses sharding to support deployments with very large data sets and high throughput operations."
The shard key is a field (or compound of fields) you designate when enabling sharding on a collection. Every document in the collection must have the shard key field(s). MongoDB uses the shard key to assign each document to a chunk — a contiguous range of key values. Each chunk lives on exactly one shard. — MongoDB Manual: "The shard key determines the distribution of the collection's documents among the cluster's shards. The shard key is either an indexed field or indexed compound fields that exists in every document in the collection."
Three properties determine whether a shard key is good or bad:
| Property | What it means | Bad choice | Good choice |
|---|---|---|---|
| Cardinality | How many distinct values exist. Each distinct value can produce at most one chunk — low cardinality caps how many chunks (and therefore shards) you can distribute across. | status with 3 values → max 3 chunks → max 3 shards regardless of data size |
userId or a hashed field — millions of distinct values → many chunks |
| Write distribution | Are inserts spread across shards or piling up on one? A monotonically increasing key (ObjectId, timestamp) sends all new writes to the last chunk on the last shard — a write hotspot. | createdAt timestamp with range sharding → all new docs go to one shard |
{ userId: 1, createdAt: 1 } compound, or hashed sharding on a high-cardinality field |
| Query isolation | Can most queries include the shard key, hitting only one shard? Without the shard key, mongos must ask every shard (scatter-gather). | Shard on a field you never filter on — every query scatters | Shard on the field your most common query uses in its filter |
All three matter. A shard key that's perfect on cardinality but terrible on write distribution creates a hotspot that defeats horizontal scale. The MongoDB Manual calls this a "non-monotonically increasing" requirement for range-sharded keys.
The cluster divides each collection's shard key space into chunks — contiguous key ranges. As data grows, chunks split when they exceed the configured size (default 128 MB). The balancer is a background process that migrates chunks between shards when the distribution becomes uneven. — MongoDB Manual: "The MongoDB balancer is a background process that monitors the number of chunks on each shard. When the number of chunks on a given shard reaches certain migration thresholds, the balancer automatically migrates chunks between shards."
Chunk migration is not free: while a chunk is being moved, the source shard still serves reads and writes for it, then the chunk is committed to the destination. This can cause brief write latency spikes during heavy migration. The balancer runs in a maintenance window you can configure.
| Strategy | How chunks are assigned | Good for | Bad for |
|---|---|---|---|
| Range sharding | Adjacent key values go to the same chunk (and shard). Contiguous ranges are co-located. | Range queries — createdAt between A and B hits one or few shards. Good read locality. |
Monotonically increasing keys (ObjectId, timestamp) — all new writes pile into the last chunk = write hotspot. |
| Hash sharding | MongoDB hashes the shard key value; the hash determines the chunk. Adjacent key values scatter across shards. | Write distribution — even for monotonic keys; each new ObjectId hashes to a random chunk. | Range queries — a date range scatters across all shards because the hash is non-monotonic. |
In practice: use hash sharding when write throughput is the problem and queries are mostly point lookups; use range sharding when range queries dominate and your key is not monotonic (or you use a compound key to break the monotonicity).
This is the operational consequence of the shard key choice — the equivalent of using an index vs doing a COLLSCAN, but at cluster scale.
The parallel to indexing (0003): a query without an index causes a COLLSCAN — one machine reads every document. A sharded query without the shard key causes a scatter-gather — every machine reads every document. Both are proportional to total data; neither is what you designed sharding to achieve.
{ status: 1 } with three values means three possible chunks — the cluster can never distribute across more than three shards. Adding a fourth shard buys nothing.
{ _id: 1 } (ObjectId) or { createdAt: 1 } — all inserts go to the rightmost chunk. One shard takes 100% of writes; the rest idle. This is the B-tree right-page contention problem from MySQL at cluster scale.
{ country: 1, userId: 1 } where 90% of users are in one country — 90% of data lands on a small number of chunks for that country. Uneven chunk distribution defeats the balancer's ability to help.
{ userId: 1, createdAt: 1 } for a user-activity collection — userId provides high cardinality and even distribution; most queries filter by userId, hitting one shard; createdAt provides sort ordering within a shard. This is the canonical pattern.
A query that does not include the shard key field causes mongos to:
You shard an orders collection on { createdAt: 1 } with range sharding. All new orders have increasing timestamps. The result is:
mongos holds:
Why does low shard key cardinality (e.g., a boolean field) limit sharding effectiveness?
Hash sharding vs range sharding — the key trade-off:
sh.enableSharding("mydb"). Shard a collection on a field with sh.shardCollection("mydb.col", { userId: 1 }). Run sh.status() to see chunk distribution across shards. Then run a query with and without the shard key and compare explain("executionStats") — look for "SHARD_MERGE" (scatter) vs a single-shard plan (targeted).