Explain Sharding and Horizontal Scaling
Sharding & Horizontal Scaling
When datasets grow beyond the memory or I/O capacity of a single server, MongoDB employs Sharding to partition data across a distributed cluster of machines.
Sharded Cluster Components
| Component | Architectural Role | High Availability Setup |
|---|---|---|
| Shard | Stores a partitioned subset of data (chunks) | Deployed as an independent Replica Set |
| mongos (Router) | Stateless query router; routes client operations to target shards | Run multiple instances behind a load balancer or on app servers |
| Config Servers | Stores cluster metadata, routing tables, and chunk distributions | Dedicated 3-member Replica Set (CSRS) |
Shard Keys & Partitioning Strategies
- Ranged Sharding: Divides data into ranges based on shard key values. Efficient for range queries, but can lead to write hotspots if values are monotonically increasing (e.g. timestamps).
- Hashed Sharding: Hashes the shard key field to evenly distribute writes across shards. Eliminates write bottlenecks at the cost of broadcast range scans.
- Good Shard Key Qualities: High cardinality, balanced write distribution, and query isolation.