Limited Offer

30% OFF Lifetime Access ($139) with code SYSTEM30

TOPIC #57Beginner 9 min read

Replication (Primary-Replica / Master-Slave)

💡
Core Architecture Summary

Scale read capacity and provide disaster recovery: Primary-Replica topology, WAL streaming, replication lag, replica promotion, and split-brain prevention.

Key Glossary Concepts in this TopicAll Glossary Terms

Primary-Replica (Leader-Follower) Database Architecture 👑👥

All write mutations target the Primary node and stream via WAL logs to Read Replicas, scaling read throughput linearly while providing standby failover redundancy.

Primary-Replica (Leader-Follower) Database Architecture 👑👥
100%
Rendering visual architecture flowchart...

01.1. What is Primary-Replica Replication and Why is it Used?

In modern web architectures, read queries typically outnumber write operations by 10:1 to 100:1 (Read-Heavy workloads). A single database server processing 100,000 queries per second will rapidly exhaust its CPU and memory bandwidth.

Primary-Replica (Leader-Follower) Replication distributes database load across multiple nodes:

  1. The Primary Node (Leader / Master): The single authoritative node that accepts all write mutations (INSERT, UPDATE, DELETE) and executes DDL schema migrations.
  2. Read Replicas (Followers / Slaves): Read-only copies of the database that continuously replay the Primary's log stream, serving read queries (SELECT) to scale read throughput horizontally.
  3. High Availability & Disaster Recovery: If the Primary node suffers hardware failure, a monitoring orchestrator (Patroni, Orchestrator, AWS RDS Multi-AZ) automatically promotes a healthy Read Replica to become the new Primary, restoring write availability in under 30 seconds.

02.2. Physical WAL Streaming Mechanics (PostgreSQL & MySQL)

How do changes flow from the Primary to Replicas?

  • Physical Streaming Replication (PostgreSQL): The Primary's walsender process streams raw binary Write-Ahead Log (WAL) byte records over a dedicated TCP connection to the replica's walreceiver process. The replica applies the WAL changes directly to its local buffer pool and disk pages.
  • Statement-Based Replication (Legacy MySQL): The primary logs the raw SQL query string (e.g., UPDATE users SET status='active'). Flaw: Non-deterministic functions (NOW(), UUID(), RAND()) cause replicas to diverge from the primary!
  • Row-Based Replication (Modern MySQL Binlog): The primary logs the exact byte diff of the modified row (Before-Image and After-Image), guaranteeing 100% deterministic replica fidelity.

03.3. Replication Lag & The Read-Your-Own-Writes Problem

Because replication is almost universally asynchronous to maximize write throughput, there is an unavoidable delay called Replication Lag (typically 5ms to 500ms over healthy networks, but stretching to seconds during heavy batch writes).

The "Read-Your-Own-Writes" Disaster:

  1. User Alice updates her profile name to "Alice Smith" (writes to Primary).
  2. The browser immediately reloads and queries a Read Replica for the updated profile.
  3. The replica has not yet processed the WAL stream (50ms lag).
  4. The replica returns the old name "Alice Jones".
  5. Alice panics and assumes her update failed!

Architectural Solutions:

  • Pin User to Primary After Writes: After a user performs a write mutation, route all subsequent read queries for that specific user to the Primary node for 5–10 seconds, falling back to replicas only after replication lag has cleared.
  • Track Monotonic Commit Timestamps: The client receives a Commit Log Sequence Number (LSN / GTID). Replicas only satisfy reads if their local LSN ≥ client LSN; otherwise, the request waits or routes to Primary.

⚖️Architectural Trade-offs & Production Realities

Architectural Advantages

  • Scales read throughput linearly: adding 5 read replicas increases read QPS capacity by 5x.
  • Enables high availability: standby replicas provide fast disaster recovery failover.
  • Allows offloading heavy reporting, analytical queries, and database backups to replicas without degrading production API performance.

Trade-offs & Constraints

  • Asynchronous replication lag causes temporary stale reads and inconsistent user views.
  • Does not scale write capacity: all writes are still bottlenecked by the single Primary node.
Production Implementation in Big Tech
GitHub & GitLab• Primary-Replica Scaling with Orchestrator

GitHub runs MySQL clusters with 1 Primary and up to 10 Read Replicas per cluster. Using GitHub Orchestrator and ProxySQL, read traffic is balanced across replicas, while write traffic routes to the Primary. Automated failovers promote standby replicas within 10 seconds during hardware faults with zero dropped transactions.

🎯 Staff+ Engineering Takeaways

  • Primary handles 100% of writes; Replicas handle read queries.
  • Replication scales read throughput horizontally but does not scale write capacity.
  • Asynchronous replication introduces Replication Lag (stale reads).
  • Solve Read-Your-Own-Writes by routing recent writers to the Primary for a short window.
  • Automated failover promotes the most up-to-date replica to become the new Primary.

Topic Knowledge Assessment 🧠

Step through 2 scenario questions to test your staff-level grasp.

Question 1 of 20 answered
#1

What is the "Read-Your-Own-Writes" consistency problem in a database cluster with asynchronous Read Replicas?

Rate This Architecture Chapter4.9 / 5.0 (38 ratings)

How clear and staff-actionable was this system breakdown?