Replication and Sharding Concepts

Advanced
14 min

Replication and Sharding Concepts

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.

Replica Sets: High Availability

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:

  • Use an odd number of voting members, at least three, so elections always reach a majority.
  • Prefer a third data-bearing member over an arbiter, which votes but holds no data and cannot help with majority write concern.
  • Special members: priority: 0 (never primary), hidden (for analytics or backups) and secondaryDelaySecs (delayed, as protection against operator error).
bash
# 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" } ] })'

Write Concern, Read Concern and Read Preference

| 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 |

javascript
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.

Sharding: Horizontal Scale

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).

javascript
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()

Choosing a Shard Key

A good key has:

  • High cardinality — many distinct values, so chunks can be split.
  • Low frequency — no single value that holds a large share of the documents.
  • Non-monotonic change (for ranged keys) — an always-increasing key such as a timestamp or a plain ObjectId sends every insert to the last chunk on one shard. Hash it, or combine it with a prefix like region.
  • Presence in common queries, so they are targeted rather than broadcast.

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 |

Common Mistakes

  • Even member counts (two or four) that cannot form a majority after one failure.
  • Reading from secondaries and expecting your own last write to be visible.
  • A monotonically increasing ranged shard key, creating a hot shard for inserts.
Quick Quiz
Question 1 of 3

What happens when a replica set primary becomes unreachable?

Key Takeaways

  • A replica set replicates the primary's oplog to secondaries and elects a new primary automatically; use an odd number of voting members.
  • Write concern "majority" makes writes durable across failover; read concern and read preference control freshness and which member serves reads.
  • Secondary reads can be stale; use the primary or causal sessions for read-your-own-writes.
  • Sharding partitions collections by a shard key across shards behind mongos routers; queries including the key are targeted.
  • Choose shard keys with high cardinality, low frequency and non-monotonic values, and shard only when a replica set is truly exhausted.

Next lesson: Capstone Project: Building an Order Management API — combine modeling, transactions, aggregation and change streams in one working service.

Replication and Sharding Concepts - MongoDB | CodeYourCraft | CodeYourCraft