MongoDB Replica Sets & Sharding

MongoDB scales in two directions. A replica set keeps copies of the same data on several servers: one primary accepts writes, secondaries replicate from it, and an election promotes a new primary automatically if the current one fails. That gives high availability and read scaling. Sharding splits a collection's data across multiple replica sets (shards), so storage and write throughput can grow beyond a single machine.

Replica sets are the baseline for any production deployment. Sharding is a bigger commitment: it adds routing components and makes the shard key one of the most consequential decisions in your data model, because it determines whether queries hit one shard or all of them and whether load spreads evenly.

TL;DR

Quick Example

Sharding an events collection for a multi-tenant SaaS:

Core Concepts

Replica Sets

Replication behavior interacts with write and read concern; see MongoDB transactions.

Sharded Cluster Architecture

Data in a sharded collection is split into chunks, contiguous ranges of shard key values. The balancer moves chunks between shards to keep data evenly distributed. Unsharded collections live on each database's primary shard.

Choosing a Shard Key

A good shard key has:

  1. High cardinality: many distinct values, so data can be split finely.
  2. Low frequency: no single value dominates (one giant tenant makes a jumbo chunk).
  3. Non-monotonic writes: an ever-increasing key such as a timestamp or ObjectId alone sends every insert to the same "last" chunk, a hot shard.
  4. Query alignment: most queries include the key (or its prefix) so they can be targeted.

Compound keys often satisfy all four: { tenantId: 1, createdAt: 1 } targets per-tenant queries while spreading each tenant's data over time.

Ranged vs Hashed Sharding

Zone Sharding

Zones pin ranges of the shard key to specific shards, for example keeping EU customers' data on shards in EU regions for data residency (GDPR), or keeping recent hot data on faster hardware.

Best Practices

Don't Shard Prematurely

A well-indexed replica set on appropriately sized hardware handles a great deal. Shard when you're approaching limits on storage, write throughput, or working-set memory on a single replica set, not "just in case". Sharding adds operational and query-design complexity.

Test the Shard Key With Real Query Patterns

Before sharding, list your top queries and check that the key targets them. Simulate insert patterns to confirm writes spread across shards. MongoDB's analyzeShardKey command (7.0+) reports cardinality, frequency, and monotonicity for candidate keys.

Include the Shard Key in Queries and Updates

Targeted queries touch one shard and scale linearly. Scatter-gather queries touch every shard and get slower as you add shards. Where possible, design APIs so requests carry the shard key (tenant ID, user ID).

Plan Capacity Per Shard

Watch per-shard disk, CPU, and working set. The balancer evens out data, not load: a shard holding a few very hot tenants can be overloaded while evenly sized.

Common Mistakes

Monotonic Shard Keys

Low-Cardinality Keys

Sharding on country or status means at most a few dozen distinct values, so chunks can't be split further and become "jumbo", and data can't balance. Combine with a high-cardinality field.

Running Two-Member Replica Sets

With two members, losing one leaves no majority, so no primary can be elected and writes stop. Use three voting members minimum, placed in separate availability zones.

FAQ

When should I shard MongoDB?

When one replica set can no longer handle your storage, write throughput, or working-set size even after indexing and vertical scaling, or when you need geographic data placement. Many large applications run happily on unsharded replica sets for years.

Can I change the shard key later?

Yes. Since MongoDB 5.0, reshardCollection rewrites the collection with a new key online, and 4.4+ allows refining a key by adding suffix fields. Resharding is resource-intensive and can take a long time on big collections, so it's a recovery path, not a routine operation.

Do transactions work across shards?

Yes, since MongoDB 4.2. Cross-shard transactions use a two-phase commit coordinated by the cluster and cost more than single-shard ones. Designing so related writes share a shard key value keeps most transactions on one shard.

Is Atlas different?

MongoDB Atlas runs the same replica set and sharding architecture as a managed service. It automates provisioning, upgrades, backups, scaling, and even global cluster zone setup, but shard key design is still your responsibility.

Related Topics

References