Type something to search...
Understanding MongoDB Replica Sets and High Availability

Understanding MongoDB Replica Sets and High Availability

A single MongoDB server is a single point of failure. When the machine reboots for a kernel patch, the disk fills up, or the cloud provider has a bad afternoon in one availability zone, your application goes down with it. Backups help you recover data after the fact, but they don't keep the application running while you restore.

A replica set is MongoDB's built-in answer to this. It's a group of mongod processes that hold copies of the same data. One member, the primary, accepts writes. The others, the secondaries, continuously copy those writes. If the primary disappears, the remaining members hold an election and promote a secondary within seconds, and modern drivers reconnect to the new primary on their own. For production deployments, a replica set isn't an optional upgrade; it's the baseline.

This guide covers how replica sets are structured, how elections and failover work, how to set one up locally and inspect it, how to design member layouts for real-world failures, and what your application needs to do to ride out a failover cleanly.

The Anatomy of a Replica Set

A typical replica set has three data-bearing members:

  • Primary: receives all writes and, by default, all reads. There's at most one primary at a time.
  • Secondaries: replicate the primary's changes by reading its oplog (the operations log) and applying the same operations to their own data. They can serve reads if your application asks for it.

Every write on the primary is recorded in the oplog, a special capped collection in the local database. Secondaries tail that oplog, much like tail -f on a log file, and apply each entry in order. If you want the details of that mechanism, MongoDB Oplog Explained goes deep on it.

Members send heartbeats to each other every two seconds. If the primary stops responding for longer than the election timeout (10 seconds by default), the secondaries conclude it's gone and one of them calls an election.

Member Types and Options

Beyond plain secondaries, you can configure members for specific jobs:

Member typeVotesCan become primaryTypical use
SecondaryYesYesNormal redundancy
Priority 0YesNoA member in a distant region that shouldn't lead
HiddenYesNoDedicated to backups or analytics, invisible to apps
DelayedYesNoRuns an hour behind to recover from human error
ArbiterYesNo (holds no data)Breaking ties when you can't afford a third data node

A replica set can have up to 50 members, but only 7 of them can vote. In practice, most deployments use three or five voting members.

How Elections Work

MongoDB's election protocol is based on Raft. The short version:

  1. A secondary notices the primary has been unreachable for electionTimeoutMillis (10 seconds by default).
  2. It increments the election term and asks the other voting members for their votes.
  3. A member grants its vote only if the candidate's oplog is at least as up to date as its own and it hasn't already voted in this term.
  4. The candidate that collects votes from a majority of voting members becomes primary.

The majority requirement is the heart of the design. With three voting members, a majority is two. If a network partition splits the set into a group of two and a group of one, only the group of two can elect a primary. The isolated member steps down (if it was primary) and becomes read-only. That guarantees you never have two primaries accepting conflicting writes, which is the dreaded split brain scenario.

It also explains why you want an odd number of voting members. Four members need three votes for a majority, so they tolerate only one failure, the same as three members, while costing more. Five members tolerate two failures.

Voting membersMajority neededFailures tolerated
110
321
431
532
743

Priority influences which member wins. Members with a higher priority value are preferred, and if a higher-priority member catches up after an election, it will call a new election to take over. That's useful for keeping the primary in your main data center under normal conditions.

A typical failover, from primary failure to a new primary accepting writes, takes about 10 to 12 seconds with default settings: most of that is the election timeout, and the election itself is fast.

Setting Up a Replica Set Locally

The quickest way to experiment is Docker Compose. This file starts three members of a replica set called rs0:

services:
  mongo1:
    image: mongo:8.0
    command: ["mongod", "--replSet", "rs0", "--bind_ip_all", "--port", "27017"]
    ports: ["27017:27017"]
    volumes: ["mongo1:/data/db"]
  mongo2:
    image: mongo:8.0
    command: ["mongod", "--replSet", "rs0", "--bind_ip_all", "--port", "27018"]
    ports: ["27018:27018"]
    volumes: ["mongo2:/data/db"]
  mongo3:
    image: mongo:8.0
    command: ["mongod", "--replSet", "rs0", "--bind_ip_all", "--port", "27019"]
    ports: ["27019:27019"]
    volumes: ["mongo3:/data/db"]

volumes:
  mongo1:
  mongo2:
  mongo3:

After docker compose up -d, connect to the first member and initiate the set:

// docker compose exec mongo1 mongosh
rs.initiate({
  _id: "rs0",
  members: [
    { _id: 0, host: "mongo1:27017", priority: 2 },
    { _id: 1, host: "mongo2:27018", priority: 1 },
    { _id: 2, host: "mongo3:27019", priority: 1 },
  ],
});

The hostnames in the config are what drivers use to reach each member, so they must be resolvable from wherever your application runs. Inside the Compose network, mongo1 works. From your laptop, you'll need entries in /etc/hosts pointing those names at 127.0.0.1. This mismatch is the most common reason a local replica set "works in the shell but not in my app." For more Docker-specific setup, see Running MongoDB in Docker and Docker Compose.

This setup has no authentication, which is fine for a disposable local cluster. In any shared environment, enable access control and use a keyfile or x.509 certificates for internal member authentication.

Inspecting the Replica Set

rs.status() is the command you'll use most. It's verbose, so here's a trimmed view:

rs.status().members.map((m) => ({
  name: m.name,
  state: m.stateStr,
  health: m.health,
  optime: m.optimeDate,
}));
[
  {
    name: "mongo1:27017",
    state: "PRIMARY",
    health: 1,
    optime: ISODate("2026-09-14T06:12:04Z"),
  },
  {
    name: "mongo2:27018",
    state: "SECONDARY",
    health: 1,
    optime: ISODate("2026-09-14T06:12:04Z"),
  },
  {
    name: "mongo3:27019",
    state: "SECONDARY",
    health: 1,
    optime: ISODate("2026-09-14T06:12:03Z"),
  },
];

Other useful commands:

rs.conf(); // current configuration
rs.hello(); // this member's view: isWritablePrimary, hosts, setName
rs.printSecondaryReplicationInfo(); // how far each secondary lags behind
rs.printReplicationInfo(); // oplog size and time window

The lag output tells you how many seconds behind the primary each secondary is:

source: mongo2:27018
{
  syncedTo: 'Mon Sep 14 2026 06:12:04 GMT+0000',
  replLag: '0 secs (0 hrs) behind the primary '
}

Healthy secondaries are usually within a second or two. Sustained lag means a secondary can't keep up (slower disk, heavy reads, network issues), and a lagging member takes longer to become a viable primary.

Testing a Failover

Don't wait for production to find out how your application behaves during an election. You can trigger one deliberately:

// On the primary: step down for 60 seconds
rs.stepDown(60);

The primary becomes a secondary, and one of the others is elected. With the Compose setup, you can also simulate a crash:

docker compose stop mongo1
# watch the election happen
docker compose exec mongo2 mongosh --port 27018 --eval 'rs.status().members.map(m => m.name + " " + m.stateStr)'
docker compose start mongo1

When mongo1 comes back, it rejoins as a secondary, catches up from the oplog, and, because it has the highest priority, takes over as primary again shortly after.

Run your application's test suite or a load generator during these experiments. You should see a short pause in writes, a few retried operations, and no errors surfacing to users.

What Your Application Needs to Do

A replica set only delivers high availability if the application is configured to take advantage of it.

Connect to the set, not a single host. List multiple members and the set name, or use an SRV record:

import { MongoClient } from "mongodb";

const client = new MongoClient(
  "mongodb://mongo1:27017,mongo2:27018,mongo3:27019/?replicaSet=rs0",
);
// Atlas and many hosted setups use SRV instead:
// mongodb+srv://user:pass@cluster0.example.mongodb.net/

The driver discovers the full topology from any member it can reach, monitors all members, and routes writes to whoever is currently primary. If you connect to a single host without replicaSet, the driver may treat it as a direct connection and won't follow a failover.

Leave retryable writes and reads on. Current drivers enable them by default. When a write fails because the primary stepped down, the driver waits for a new primary and retries the operation once. This covers the vast majority of failover blips without any code on your part. MongoDB Retryable Writes and Reads explains exactly what is and isn't retried.

Use majority write concern for data you can't lose. With w: 1, the primary acknowledges a write as soon as it has it. If the primary crashes before any secondary copies that write, the write gets rolled back when the old primary rejoins. With w: "majority" (the default in recent versions for most configurations), the write is acknowledged only after a majority has it, so it survives any single failure.

Set sensible timeouts. serverSelectionTimeoutMS (default 30 seconds) controls how long the driver waits for a suitable server, which covers the election window. Don't lower it below 15 seconds or so, or operations will fail during normal elections.

Designing for Real Failures

Where you place members matters as much as how many you have.

Spread members across failure domains. Three members on the same physical host protect you from a mongod crash but not from a host failure. In the cloud, put each member in a different availability zone. That's what Atlas does by default for its dedicated clusters.

Plan for losing a whole data center. With two data centers, you can't survive the loss of either one with automatic failover, because whichever site has the majority of votes becomes a single point of failure. The robust layout uses three sites:

LayoutSite ASite BSite CSurvives loss of
2 sites, 3 members21noneSite B only
3 sites, 3 members111Any one site
3 sites, 5 members221Any one site

The five-member layout keeps two members in each primary site (so a single node failure doesn't force cross-site traffic) and puts a single tie-breaker member in the third site.

Be careful with arbiters. An arbiter votes but holds no data, which makes a primary-secondary-arbiter (PSA) set look like a cheap alternative to three data nodes. The problem: if the secondary goes down, the set still has a primary, but only one data-bearing member, so majority writes can't be acknowledged and majority read concern can stall. Cache pressure on the primary rises while it waits. Use three data-bearing members whenever you can, and treat arbiters as a last resort.

Use a delayed member as an undo button. A hidden, priority-0 member with secondaryDelaySecs: 3600 stays an hour behind. If someone runs a bad deleteMany in production, you have an hour to stop replication on that member and copy the data back. It's not a replacement for backups, but it's a much faster recovery for "oops" moments.

rs.add({
  host: "mongo-delayed:27017",
  priority: 0,
  hidden: true,
  secondaryDelaySecs: 3600,
  votes: 0,
});

Giving the delayed member zero votes keeps it from affecting elections or the write majority.

Common Pitfalls

Running an even number of voting members. Four voters tolerate one failure, the same as three, and are more likely to deadlock in a symmetric partition. Add a fifth member or remove the vote from one.

Treating secondaries as free read capacity. Reading from secondaries can return stale data, and if the set is sized so that secondaries handle significant read traffic, losing one member overloads the rest. Size the set so the primary alone can handle the workload, and use secondary reads deliberately.

Ignoring replication lag. A secondary that's minutes behind can't take over quickly, and reads from it are minutes stale. Alert on lag, not just on member health.

Undersizing the oplog. If a secondary goes offline for longer than the oplog window, it can't catch up and needs a full resync. Check rs.printReplicationInfo() and make sure the window covers your longest realistic maintenance.

Hardcoding the primary's hostname. Connection strings that name only one host defeat the purpose. Always include multiple seeds or use SRV.

Skipping failover drills. Many teams discover during their first real failover that some service uses an old driver, a custom retry wrapper, or a connection string without the set name. Drill regularly.

Conclusion

A replica set keeps copies of your data on multiple servers, elects a new primary automatically when the current one fails, and uses majority voting to make split brain impossible. To get real high availability, run an odd number of voting members across separate failure domains, prefer data-bearing members over arbiters, use majority write concern for important data, and connect with a replica-set-aware connection string so drivers can follow failovers with retryable writes.

This week, run the Docker Compose setup above, point a small script that inserts a document every 100 milliseconds at it, and call rs.stepDown() on the primary. Watching your script pause for a few seconds and then carry on without a single error is the best way to build confidence in your failover story.

Tags :
Share :

Related Posts

A Complete Guide to MongoDB Query Operators

A Complete Guide to MongoDB Query Operators

Your first MongoDB queries are usually simple equality filters: find the user with this email, find orders with this status. That covers a surprising

Continue Reading
Async MongoDB in Python with Motor and FastAPI

Async MongoDB in Python with Motor and FastAPI

FastAPI runs your endpoints on an event loop. That's what lets a single worker juggle hundreds of concurrent requests: while one request waits on the

Continue Reading
Atlas Online Archive: Tiering Cold Data to Cut Costs

Atlas Online Archive: Tiering Cold Data to Cut Costs

Look at almost any production database and you'll find the same shape. A small slice of recent data gets nearly all the reads and writes: this week's

Continue Reading