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
- A replica set is a primary plus secondaries (usually 3 members); elections pick a new primary on failure within seconds.
- A sharded cluster has shards (each a replica set), config servers (metadata), and mongos routers that clients connect to.
- The shard key determines data distribution. Choose one with high cardinality, even write distribution, and alignment with your most common queries.
- Ranged sharding keeps adjacent keys together (good for range queries, risks hotspots); hashed sharding spreads writes evenly (bad for range queries).
- Queries including the shard key are targeted to one shard; others are scatter-gather across all shards.
- Shard keys can be changed with resharding (5.0+), but it's expensive, so choose carefully up front.
Quick Example
Sharding an events collection for a multi-tenant SaaS:
Core Concepts
Replica Sets
- The primary receives all writes and records them in the oplog, a capped collection of operations.
- Secondaries tail the oplog and apply operations asynchronously.
- If the primary becomes unreachable, members hold an election. A new primary is typically chosen within about 10–12 seconds, and drivers automatically redirect writes to it.
- Use an odd number of voting members (3 or 5) across failure domains such as availability zones. Arbiters (vote-only members) are discouraged for production because they weaken majority write guarantees.
- Special members: hidden (for backups or analytics), delayed (a protection window against accidental deletes), and priority 0 (never becomes primary).
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:
- High cardinality: many distinct values, so data can be split finely.
- Low frequency: no single value dominates (one giant tenant makes a jumbo chunk).
- 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.
- 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
- MongoDB — The database overview
- Database Sharding — Sharding strategies in general
- MongoDB Transactions — Consistency across replicas and shards
- MongoDB Indexes — The shard key must be indexed
- High Availability — Replica set placement and failover
- Redis Cluster — A different approach to sharding with hash slots