Skip to main content
Create your own

MongoDB Replica Sets: Primary-Secondary Replication and Failover

Hello! Welcome back to our module on MongoDB.

In our previous lesson, we explored MongoDB's tunable consistency model, focusing on what guarantees you can achieve with read and write concerns. We discussed how w: "majority" is the cornerstone of durable writes that can survive node failures.

Today, we'll build upon that by examining the mechanism that makes this durability and high availability possible: the MongoDB replica set. This lesson moves from the logical controls of consistency to the physical architecture of replication and automatic failover. This is the practical foundation for building resilient, production-grade MongoDB deployments.

By the end of this ~60-minute lesson, you will be able to set up a MongoDB replica set, understand the replication and election processes, and witness automatic failover in action.

1. Replica Set Architecture

A replica set is a group of mongod processes that maintain an identical data set, providing redundancy and high availability. Within a replica set, nodes have distinct roles:

  • Primary: There is only one primary node at any time. It receives all write operations from clients and records them in its operations log (oplog).
  • Secondaries: These nodes asynchronously replicate the primary's oplog and apply the operations to their own data sets. They can handle read requests, which helps scale read-heavy workloads.
  • Arbiter (Optional): An arbiter is a mongod instance that participates in elections for a new primary but does not hold a copy of the data. Its sole purpose is to add a vote to break ties in replica sets with an even number of data-bearing nodes.

A standard, recommended architecture is a three-member replica set with one primary and two secondaries (P-S-S). This configuration can tolerate the loss of one node while still maintaining a majority ((3/2) + 1 = 2) to elect a new primary and satisfy w: "majority" writes.

A diagram showing a Primary-Secondary-Secondary (P-S-S) replica set architecture. A client application sends write operations to the Primary node. The Primary node then replicates these operations to two Secondary nodes. The client can be configured to send read operations to either the Primary or the Secondaries.
A standard three-member replica set. The Primary handles all writes, which are then replicated to the Secondaries. Reads can be distributed across the nodes.

As we discussed in the last lesson, using an arbiter to create a P-S-A architecture can be problematic. If the single secondary becomes unavailable, the replica set can no longer achieve a majority for writes, as the arbiter doesn't count towards the w: "majority" calculation. Therefore, using an odd number of data-bearing nodes is the standard best practice.

2. The Replication and Failover Process

The mechanics of replication and failover are driven by two key components: the oplog and the election protocol.

The Oplog: Heart of Replication

The oplog (operations log) is a special capped collection (oplog.rs) stored in the local database of each replica set member. The primary node writes all data-modifying operations (inserts, updates, deletes) to its oplog in a sequential, idempotent manner. Secondaries continuously tail this log, retrieve the operations, and apply them to their own data sets.

This design has important implications:

  • Asynchronicity: Because secondaries pull data from the primary, there is a natural delay known as replication lag. This is why reading from a secondary can return stale data.
  • Idempotency: Oplog entries are designed to have the same effect whether applied once or multiple times. This is critical for safely resuming replication after a network interruption or node restart.
  • Size: The oplog must be large enough to hold all transactions that occur during the longest potential downtime of a secondary. If a secondary is down for too long and the primary's oplog has "rolled over," the secondary will be too far behind to catch up and will require a full, resource-intensive resynchronization from another member.

Elections and Automatic Failover

High availability is achieved through a robust election process, which is an implementation of a consensus algorithm (similar to Raft).

  1. Heartbeats: All members of the replica set send heartbeats to each other every two seconds.
  2. Failure Detection: If a node is unreachable for a default period of 10 seconds, other members mark it as unavailable.
  3. Election Trigger: If the unavailable node was the primary, one of the eligible secondaries will call an election. An eligible secondary must be up-to-date, meaning its oplog is not too far behind the last known state of the primary.
  4. Voting: The candidate requests votes from all other members in the set. A node wins the election if it receives votes from a strict majority of all members in the replica set configuration (including those that are down).
  5. New Primary: The winner transitions to the primary state, and all other nodes begin replicating from it. When the old primary comes back online, it detects the new primary, rolls back any writes that were not replicated before it failed, and rejoins the set as a secondary.

This process ensures that a new primary is elected automatically, typically within seconds, minimizing downtime for the application.

3. Practical Guide: Setting Up a Replica Set

Let's walk through the steps to create and manage a replica set on a single machine for development purposes. In a production environment, each mongod instance would run on a separate server.

We will use mongosh, the MongoDB Shell.

Step 1: Launch mongod Instances

First, create three data directories for our three nodes:

mkdir -p /data/rs1 /data/rs2 /data/rs3

Now, open three separate terminal windows and launch a mongod process in each. Note the --replSet option, which must be the same for all members, and the different --port and --dbpath for each instance.

Terminal 1:

mongod --port 27017 --dbpath /data/rs1 --replSet myReplicaSet

Terminal 2:

mongod --port 27018 --dbpath /data/rs2 --replSet myReplicaSet

Terminal 3:

mongod --port 27019 --dbpath /data/rs3 --replSet myReplicaSet

At this point, you have three independent mongod instances running. They are aware they should belong to a replica set named myReplicaSet, but the set has not been initialized yet.

Step 2: Initiate the Replica Set

Connect to any one of the instances using mongosh. Let's use the first one.

mongosh --port 27017

Inside the shell, define a configuration object that lists the members of your set and then use rs.initiate() to create the set.

// Define the replica set configuration
config = {
    _id: "myReplicaSet",
    members: [
        { _id: 0, host: "localhost:27017" },
        { _id: 1, host: "localhost:27018" },
        { _id: 2, host: "localhost:27019" }
    ]
};

// Initiate the replica set
rs.initiate(config);

After you run this, the shell prompt will change to indicate which member of the replica set you are connected to (e.g., myReplicaSet:PRIMARY>). The instance where you ran initiate becomes the first primary.

Step 3: Verify the Replica Set Status

The rs.status() command is your primary tool for inspecting the health and status of a replica set.

myReplicaSet:PRIMARY> rs.status()

The output is verbose but contains critical information. Look for these key fields:

  • set: The name of the replica set.
  • members: An array of objects, one for each member.
    • name: The host and port of the member.
    • stateStr: The current role of the member ("PRIMARY", "SECONDARY", or "ARBITER").
    • health: 1 for healthy, 0 for unhealthy.
    • optime: The timestamp of the last operation applied from the oplog. Comparing optimes between the primary and secondaries helps you gauge replication lag.

Step 4: Test Automatic Failover

Let's see the failover in action.

  1. Write some data to the primary.

    myReplicaSet:PRIMARY> use testDB
    myReplicaSet:PRIMARY> db.myCollection.insertOne({ message: "Hello from the primary" }, { writeConcern: { w: "majority" } })
    

    Using w: "majority" ensures this write is propagated to at least one secondary before returning success.

  2. Verify replication. Connect to a secondary in a new terminal and check for the data. You must first run rs.secondaryOk() to allow reads on a secondary node.

    # New terminal
    mongosh --port 27018
    
    myReplicaSet:SECONDARY> rs.secondaryOk()
    myReplicaSet:SECONDARY> use testDB
    myReplicaSet:SECONDARY> db.myCollection.findOne()
    # You should see the document inserted on the primary.
    
  3. Simulate a primary failure. Go back to the primary's mongosh session and shut it down.

    myReplicaSet:PRIMARY> db.adminCommand({ shutdown: 1 })
    

    The mongod process in Terminal 1 will exit.

  4. Observe the election. Quickly switch to the mongosh session connected to the secondary (port 27018). Run rs.status() again. You may see the nodes change state as they detect the failure and hold an election. Within a few seconds, the prompt will change from SECONDARY to PRIMARY.

    myReplicaSet:SECONDARY> rs.status() // Check status during election
    # ... wait a few seconds ...
    myReplicaSet:PRIMARY> // The prompt has changed!
    

    You have successfully witnessed an automatic failover. The replica set remains available for writes, now served by the new primary.

  5. Bring the old primary back. Restart the mongod process in Terminal 1.

    # In Terminal 1
    mongod --port 27017 --dbpath /data/rs1 --replSet myReplicaSet
    

    It will detect that an election has occurred, see the new primary, and rejoin the set as a secondary. You can verify this with rs.status() from the new primary.

4. Read Preferences: Directing Read Traffic

By default, MongoDB drivers send all read and write queries to the primary. However, you can configure the driver to route read queries to secondaries using read preferences. This is a powerful tool for scaling read throughput.

The common read preference modes are:

  • primary: Default. All reads go to the primary. Guarantees the most recent data.
  • secondary: All reads go to available secondaries. Distributes read load but risks reading stale data due to replication lag.
  • primaryPreferred: Tries the primary first; if unavailable, reads from secondaries.
  • secondaryPreferred: Tries secondaries first; if none are available, reads from the primary. Useful for analytics workloads where recency is less critical than offloading the primary.
  • nearest: Reads from the member with the lowest network latency, regardless of its role.

Choosing a read preference is a critical architectural decision that balances load distribution against data consistency. For use cases requiring "read-your-own-writes," you must read from the primary.

Conclusion

In this lesson, we moved from theory to practice, establishing the operational backbone of a highly available MongoDB system. You now understand how replica sets function and how to configure one.

Key Takeaways:

  • Replica sets provide high availability and data redundancy using a primary-secondary architecture.
  • Replication is powered by the oplog, which secondaries tail asynchronously.
  • Automatic failover is handled by a consensus-based election process triggered by heartbeats and failure detection.
  • A three-member data-bearing replica set is the standard topology for balancing fault tolerance and cost.
  • The rs.initiate() and rs.status() commands are your fundamental tools for creating and monitoring replica sets.

In our next lesson, "Set up a MongoDB sharded cluster architecture," we will address the challenge of horizontal scalability. While replica sets solve for high availability, sharding is MongoDB's solution for distributing massive datasets and write-heavy workloads across multiple replica sets.

Can't find a good explanation? Sign up and we'll make it for you

Sign up