Lesson 0009 · Horizontal scale

Sharding

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.

~13 minWarm-up + retrieval quiz + no-Postgres-equivalent contrastCheat sheet: here

Warm-up — apply lesson 0008 first (cold, from memory)

j: true and w: "majority" — how are they different?

j = single-node disk durability: "did this node's WAL flush before acking?" About surviving a process crash. w = replication acknowledgment: "how many nodes confirmed the write?" About surviving a node failure. Setting one says nothing about the other — they are orthogonal.

With default w: 1, when can an acknowledged write be silently rolled back?

w: 1 acks as soon as the primary writes to its oplog buffer. If the primary then crashes before secondaries replicate that entry, the newly elected primary won't have it — the write is rolled back silently. w: "majority" closes this window: a majority have the entry before ack, so any elected primary already has it.

readConcern: "majority" guarantees the data you read:

readConcern: "majority" returns only data that a majority of voting members have durably acknowledged. That's the same bar as w: "majority" on writes. Any elected primary already has this data — it cannot disappear due to a failover or rollback. The default "local" reads the node's most recent data, which may include writes not yet replicated.

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 three components

CLIENT │ ▼ mongos ←──────────────────── Config servers (replica set) (query router) (stores cluster metadata) │ ├──► Shard 1 (replica set) ├──► Shard 2 (replica set) └──► Shard 3 (replica set)

The shard key — the most consequential decision

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:

PropertyWhat it meansBad choiceGood 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.

Chunks and the balancer

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.

Range vs hash sharding

StrategyHow chunks are assignedGood forBad 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).

Targeted query vs scatter-gather

This is the operational consequence of the shard key choice — the equivalent of using an index vs doing a COLLSCAN, but at cluster scale.

Targeted query: the filter includes the shard key. mongos reads the chunk map, identifies exactly which shard holds the matching documents, and routes to that one shard. Latency = one shard's latency. This is the goal.
Scatter-gather: the filter does not include the shard key. mongos doesn't know which shard has the data, so it sends the query to every shard, waits for all of them, merges the results. Latency = the slowest shard's latency, and all shards bear the load. Gets worse as you add shards.

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.

Shard key anti-patterns

Low cardinality. { 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.
Monotonically increasing key with range sharding. { _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.
Under-selective compound key. { 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.
High-cardinality compound key with good query isolation. { 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.
One irreversible decision: once you shard a collection, you cannot change the shard key without resharding — which requires a full data copy on MongoDB 5.0+ (or dropping and rebuilding on older versions). Choose carefully before enabling sharding.

A query that does not include the shard key field causes mongos to:

Without the shard key, mongos can't identify which chunk (and therefore which shard) holds the matching documents. It broadcasts to every shard, waits for all results, and merges them. Latency = slowest shard, and all shards bear the load — the sharding equivalent of a COLLSCAN. The fix: include the shard key in your most common query filters.

You shard an orders collection on { createdAt: 1 } with range sharding. All new orders have increasing timestamps. The result is:

Range sharding on a monotonically increasing key is the canonical anti-pattern. All new documents have the largest timestamp and therefore land in the rightmost chunk, which lives on one shard. That shard takes 100% of write traffic; the others sit idle. The fix: use hashed sharding (even distribution but no range-query locality) or a compound key like { userId, createdAt } that breaks the monotonicity across shards.

mongos holds:

mongos is stateless with respect to your data — it reads the chunk map from the config servers and routes queries to the correct shard(s). It holds no documents. This is why you can run multiple mongos instances for redundancy without data synchronization: there's no data to synchronize. Shards (each a full replica set) hold the actual documents.

Why does low shard key cardinality (e.g., a boolean field) limit sharding effectiveness?

Chunks are ranges of shard key values. A field with only 2 distinct values (true/false) can produce at most 2 chunks. Adding a third or fourth shard gains nothing — there are no more chunks to spread. High cardinality (millions of distinct values) is required to produce enough chunks to distribute meaningfully across many shards.

Hash sharding vs range sharding — the key trade-off:

Hash sharding hashes the key value before assigning a chunk, so even sequential ObjectIds scatter evenly across shards — great for write distribution. But a date-range query can't be targeted: the hash destroys locality, so it scatters to all shards. Range sharding keeps adjacent values together, enabling targeted range queries, but suffers write hotspots on monotonically increasing keys.
Optional — when you have a mongos instance: enable sharding on a test database with 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).