Sharding for Horizontal Scaling
Understand how MongoDB distributes data across multiple servers using sharding, and how to choose a good shard key.
What Is Sharding?
Sharding distributes a collection's data horizontally across multiple servers (shards), so no single server needs to store or process the entire dataset — MongoDB's primary mechanism for scaling beyond what one machine (or replica set) can handle.
Replication (the previous lesson) creates redundant copies of the same data for reliability. Sharding splits different data across servers for scale. Production MongoDB deployments commonly use both together — each shard is itself a replica set.
The Shard Key
A shard key is a field (or combination of fields) chosen when sharding a collection, determining how MongoDB distributes documents across the available shards.
sh.shardCollection("myapp.orders", { customerId: 1 });Choosing a Good Shard Key
A good shard key has high cardinality (many distinct values) and distributes writes evenly across shards — a poor choice can create a "hot shard" that receives a disproportionate share of traffic, defeating the purpose of sharding entirely.
Good Shard Key Characteristics
- High cardinality — many distinct values
- Evenly distributes both reads and writes
- Matches your most common query patterns
Poor Shard Key Choices
- Low cardinality (e.g. a boolean or a country with few options)
- Monotonically increasing values (e.g. a timestamp) — creates a hot shard for all recent writes
- A field rarely used in actual queries
The Architecture of a Sharded Cluster
A sharded cluster has three components: the shards themselves (each typically a replica set), config servers (storing metadata about which data lives on which shard), and mongos routers (the entry point applications actually connect to, which route each query to the correct shard(s)).
FAQs
Almost certainly not — sharding adds real operational complexity and is meant for datasets or workloads that have genuinely outgrown what a single replica set can handle.
Historically difficult and disruptive; more recent MongoDB versions have added support for reshardCollection to change it, though it remains a significant operation to plan carefully.
Summary
Sharding scales MongoDB horizontally by distributing data across servers based on a carefully chosen shard key, typically reserved for large-scale production deployments. Next, you'll cover security and authentication.