Read Replicas (Primary-Replica)

⭐ Interview Importance: HIGH
⏱️ Revision Time: 4 min

Concept

Most modern web applications are heavily skewed toward reading data. In a system like Twitter, for every 1 person who writes an INSERT statement (posting a tweet), 1,000 people write SELECT statements (reading the timeline).

The Primary-Replica architecture (formerly Master-Slave) exploits this ratio to achieve Horizontal Scaling.

  1. You create one Primary Database. It is the absolute source of truth. It handles 100% of the INSERT, UPDATE, and DELETE (Write) queries.
  2. You create 3 Read Replica Databases. They are identical clones of the Primary. They handle 100% of the SELECT (Read) queries.

You have effectively multiplied your read capacity by 3x, taking a massive CPU load off the Primary server.

How the Synchronization Works

When you execute an INSERT on the Primary, how do the 3 Replicas get the new data?

They use Asynchronous Replication (Streaming WAL).

  1. The Primary executes the INSERT and instantly returns “Success” to the user.
  2. In the background, the Primary streams its Write-Ahead Log (WAL) binary file over the TCP network to the Replicas.
  3. The Replicas receive the binary instructions and “replay” the transaction against their own hard drives.

The Big Catch: Eventual Consistency (Replication Lag)

Because the replication happens asynchronously over a network, it takes time (usually a few milliseconds). This introduces Replication Lag.

The Classic Bug:

  1. User changes their profile picture. The Node.js server sends the UPDATE to the Primary Database.
  2. The Node.js server redirects the user back to the Homepage.
  3. The Homepage sends a SELECT query to fetch the profile. The Load Balancer routes this read request to Replica #2.
  4. The Bug: Replica #2 hasn’t received the WAL file from the Primary yet. It returns the old profile picture.
  5. The user thinks the website is broken and frantically clicks “Update” again.

How to Fix Replication Lag

You solve this in the Node.js application layer.

  • Read-After-Write Consistency: If the Node.js server knows it just wrote data for a specific user, it sets a temporary cookie or flag in Redis (recent_writer: true). For the next 5 seconds, any read queries for that specific user bypass the Replicas and are routed directly to the Primary database to guarantee fresh data. Everyone else continues to read from the Replicas.

Interview Questions

Q: If the Primary database physically crashes and the motherboard catches fire, what happens to the system?
A: This is called a Failover.
The system detects the Primary is dead. An automated script immediately promotes one of the Read Replicas (e.g., Replica #1) to become the new Primary. The Node.js application’s DNS router automatically updates to point write traffic to the new Primary.
However, because replication is Asynchronous, any transactions that occurred in the 5 milliseconds right before the crash—which hadn’t been streamed over the network to Replica #1 yet—are permanently lost.

Q: A developer wants to guarantee zero data loss during a crash, so they configure “Synchronous Replication”. How does this affect the application?
A: Synchronous Replication forces the Primary to wait.
When the user executes an INSERT, the Primary writes it, sends it over the network to the Replica, and physically freezes the user’s API request until the Replica acknowledges it also successfully wrote the data to its hard drive. Only then does it return “Success”.
This guarantees absolute zero data loss. However, it effectively cuts your database write performance in half and doubles your API latency, because every single write is now gated by network speed and a secondary hard drive.