Explorer
MongoDB

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.

Finished this lesson?

Mark this chapter complete to update your learning streak and unlock the next lesson.