A single server can lose a disk, a process or a data centre. A single server also has a ceiling on RAM, disk and throughput. MongoDB answers the first problem with replica sets and the second with sharding, and every Atlas cluster is built from these pieces. After this lesson you will be able to explain how elections and the oplog keep data available, tune write and read concerns for the guarantees you need, and know what a shard key does and how to choose one.
A replica set is a group of mongod processes holding the same data. One primary accepts writes and records them in its oplog, an append-only log that the secondaries tail and apply. Members exchange heartbeats every two seconds; when the primary is unreachable for electionTimeoutMillis (10 seconds by default) the remaining members elect a new one, and drivers reconnect automatically.
Design rules:
priority: 0 (never primary), hidden (for analytics or backups) and secondaryDelaySecs (delayed, as protection against operator error).# Local three-member set for learning
mongod --replSet rs0 --port 27017 --dbpath ./r0 &
mongod --replSet rs0 --port 27018 --dbpath ./r1 &
mongod --replSet rs0 --port 27019 --dbpath ./r2 &
mongosh --eval 'rs.initiate({ _id: "rs0", members: [
{ _id: 0, host: "localhost:27017" }, { _id: 1, host: "localhost:27018" }, { _id: 2, host: "localhost:27019" } ] })'| Setting | Values | Meaning |
|---|---|---|
| Write concern w | 1, "majority" (default), a number | how many members must acknowledge a write; j: true adds journal durability; wtimeout bounds the wait |
| Read concern | local, majority, linearizable, snapshot | how durable the data you read must be; majority never returns data that could be rolled back |
| Read preference | primary, primaryPreferred, secondary, secondaryPreferred, nearest | which member serves reads |
db.orders.insertOne(order, { writeConcern: { w: "majority", wtimeout: 5000 } })
db.orders.find({ status: "shipped" }).readPref("secondaryPreferred")Secondary reads can be stale because replication is asynchronous; they suit analytics, while a user who just saved a form should read from the primary or use a causally consistent session.
When one replica set can no longer hold the data or the write load, a sharded cluster partitions each collection across several shards (each itself a replica set). Config servers store the metadata, and mongos routers present the cluster to applications as a single endpoint.
Data is divided into chunks by a shard key, and a balancer moves chunks between shards to keep them even. Queries that include the shard key are routed to one shard (targeted); queries without it are sent to all shards and merged (scatter-gather).
sh.enableSharding("shop")
sh.shardCollection("shop.orders", { customerId: "hashed" }) // hashed: even spread of writes
sh.shardCollection("shop.events", { region: 1, ts: 1 }) // ranged: range queries stay targeted
sh.status()A good key has:
ObjectId sends every insert to the last chunk on one shard. Hash it, or combine it with a prefix like region.reshardCollection (5.0+) can change the key later, but it is expensive, so choose deliberately, and shard only when a well-indexed replica set is truly exhausted.
| | Replica set | Sharded cluster |
|---|---|---|
| Purpose | availability, read scaling | data volume and write scaling |
| Writes go to | one primary | the shard that owns the key range |
| Application connects to | the set (replicaSet= URI) | mongos routers |
What happens when a replica set primary becomes unreachable?
"majority" makes writes durable across failover; read concern and read preference control freshness and which member serves reads.mongos routers; queries including the key are targeted.Next lesson: Capstone Project: Building an Order Management API — combine modeling, transactions, aggregation and change streams in one working service.