Database Replication and Its Impact on Read Performance
Database replication distributes copies of data across multiple servers, enabling read scaling, geographic distribution, and high availability. The performance implications of replication architecture extend beyond simple redundancy. Properly configured replication can multiply read throughput by distributing queries across replicas, while poorly configured replication introduces latency inconsistencies that confuse applications and frustrate users.
Understanding replication mechanics, lag characteristics, and consistency trade-offs is essential for any database deployment that serves more than trivial traffic levels. The choice between synchronous and asynchronous replication, the number and placement of replicas, and the read routing strategy all directly impact both performance and data correctness.
Replication Architectures
PostgreSQL and MySQL both support streaming replication where the primary server continuously ships write-ahead log (WAL) records or binary log events to replica servers. Each replica applies these changes to maintain a near-identical copy of the primary's data. The fundamental architectural choice is whether replication commits synchronously or asynchronously.
Asynchronous Replication
Asynchronous replication commits transactions on the primary without waiting for replicas to acknowledge receipt. The primary writes the WAL record and returns success to the client immediately. Replicas receive and apply changes with a delay that ranges from milliseconds under normal conditions to seconds or minutes during high write throughput or network congestion.
This architecture provides the best write performance because the primary never waits for replicas. However, it creates a window where the primary and replicas have different data. If the primary fails during this window, committed transactions that have not yet reached the replicas are lost. For most web applications, the sub-second replication lag and the possibility of losing the most recent fraction of a second of transactions is an acceptable trade-off for write performance.
Synchronous Replication
Synchronous replication waits for at least one replica to acknowledge receipt of the WAL data before committing the transaction on the primary. This eliminates the data loss window but adds network round-trip latency to every write operation. For replicas in the same data center with sub-millisecond network latency, the overhead is typically 1 to 3 milliseconds. For cross-region replicas, the overhead equals the network round-trip time between regions, often 50 to 200 milliseconds.
# PostgreSQL synchronous replication
# postgresql.conf on primary
synchronous_standby_names = 'FIRST 1 (replica1, replica2)'
synchronous_commit = on
# For quorum-based sync (any 2 of 3 must confirm)
synchronous_standby_names = 'ANY 2 (replica1, replica2, replica3)'
Semi-synchronous replication, available in MySQL, provides a middle ground. The primary waits for at least one replica to acknowledge receipt of the binary log event before returning success to the client, but does not wait for the replica to apply the change. This guarantees that the data exists on at least two servers but does not guarantee that the replica can serve reads from the committed transaction immediately.
Measuring Replication Lag
Replication lag is the delay between a write committed on the primary and that write becoming visible on a replica. Monitoring replication lag is critical because applications routing reads to replicas may see stale data if the lag exceeds the application's staleness tolerance.
PostgreSQL exposes replication lag through the pg_stat_replication view on the primary and pg_last_wal_replay_lsn() on replicas. The difference between the primary's current WAL position and the replica's replayed position, converted to bytes and then to time using the write rate, gives the lag duration.
-- On primary: check lag for each replica
SELECT
client_addr,
state,
pg_wal_lsn_diff(
pg_current_wal_lsn(),
replay_lsn
) AS replay_lag_bytes,
replay_lag
FROM pg_stat_replication;
Lag spikes commonly occur during large bulk operations (a batch import that generates megabytes of WAL in seconds), DDL operations (adding a column to a large table), and maintenance operations (VACUUM FULL, REINDEX). Plan these operations during low-traffic periods and monitor replica lag throughout execution. For alerting thresholds, a lag exceeding 5 seconds warrants investigation, and lag exceeding 30 seconds should trigger automatic failover of read traffic back to the primary.
Read Routing Strategies
Distributing read queries across replicas requires a routing layer that directs each query to an appropriate server. The routing strategy must account for replication lag, query type, and consistency requirements.
| Strategy | How It Works | Consistency | Best For |
|---|---|---|---|
| Round-robin | Distribute reads evenly across replicas | Eventually consistent | Analytics, reporting |
| Least-connections | Route to the replica with fewest active queries | Eventually consistent | Mixed workloads |
| Lag-aware | Skip replicas with lag above threshold | Bounded staleness | User-facing reads |
| Session affinity | Route a user's reads to the same replica | Read-your-own-writes within session | Post-write reads |
| Primary fallback | Read from primary when replicas lag | Strong when needed | Critical operations |
The read-your-own-writes problem is the most common consistency issue in replicated databases. A user submits a form, the write goes to the primary, and the subsequent page load reads from a replica that has not yet received the write. The user sees their old data and believes the update failed. Solutions include routing post-write reads to the primary for a brief window, using session affinity to pin users to a specific replica, or checking replication lag before routing to replicas.
External connection poolers like ProxySQL provide built-in read/write splitting that routes SELECT queries to replicas and all other statements to the primary. This transparent routing requires no application code changes but does not handle the read-your-own-writes problem without additional configuration.
Replica Scaling Patterns
Adding replicas increases read throughput linearly up to the point where replication overhead on the primary becomes a bottleneck. Each replica requires the primary to ship WAL data, and very high replica counts can saturate the primary's network bandwidth or CPU with replication processing.
Cascading replication addresses this scaling limit by having replicas replicate from other replicas rather than directly from the primary. A tree structure where the primary feeds 2 to 3 replicas, and each of those feeds additional replicas, distributes the replication load across multiple servers. The trade-off is increased replication lag for downstream replicas because each cascade hop adds its own processing delay.
For workloads requiring extreme read throughput, consider complementing database replicas with an application caching layer. A Redis cache serving the most frequently accessed queries can handle millions of reads per second, reserving database replicas for queries that require fresh data or are too complex to cache effectively.
Failover and Promotion
When the primary server fails, a replica must be promoted to primary to restore write capability. The promotion process involves stopping replication, replaying any remaining WAL records, and reconfiguring the server to accept write connections. Automated failover tools like Patroni (PostgreSQL) and MySQL Orchestrator manage this process, detecting primary failure and coordinating promotion across the cluster.
The critical performance metric during failover is the recovery time objective (RTO): how long write operations are unavailable. With automated failover and health checks running at 1-second intervals, promotion typically completes in 10 to 30 seconds. This includes failure detection, leader election, replica promotion, and DNS or virtual IP updates to redirect client connections.
After promotion, the new primary must handle both read and write traffic until the old primary is recovered and rejoined as a replica. Monitor the new primary's CPU and I/O metrics closely during this period, and consider temporarily reducing read traffic by falling back to cached data until a replacement replica is available.
Monitoring Replication Health
Comprehensive replication monitoring tracks lag, throughput, and replica availability. Integrate these metrics into your APM dashboard for correlation with application performance.
- Replication lag — Track p50, p95, and maximum lag across all replicas. Alert on sustained lag above your staleness threshold.
- WAL generation rate — Bytes of WAL produced per second on the primary. Sudden increases indicate bulk operations or write traffic spikes.
- Replay rate — WAL bytes applied per second on each replica. Replay rate consistently below generation rate means lag will grow unboundedly.
- Replica connection status — Monitor the streaming state of each replica connection. Disconnected replicas stop receiving updates and fall behind.
- Replication slots — PostgreSQL replication slots prevent WAL recycling until the replica has consumed it. An offline replica with a replication slot causes WAL accumulation on the primary, eventually filling the disk.
Frequently Asked Questions
How many read replicas should I create?
Start with 2 replicas for redundancy and read scaling. Each replica roughly doubles read throughput for queries that can tolerate eventual consistency. Add replicas based on measured read query load and latency requirements. Most applications see diminishing returns beyond 5 to 8 replicas because the primary's WAL shipping overhead and network bandwidth become bottlenecks. For extreme read scaling, complement replicas with application-level caching.
How do I handle the read-your-own-writes problem?
Route reads to the primary for a brief window after a user performs a write, typically 2 to 5 seconds depending on normal replication lag. Implement this using a session cookie or in-memory flag that marks the user's session as requiring primary reads. After the window expires, resume routing reads to replicas. Alternatively, use lag-aware routing that checks the replica's current lag before directing a query to it.
What causes sudden replication lag spikes?
The most common causes are large batch operations that generate high WAL volume, DDL operations like adding columns or creating indexes on large tables, long-running transactions that prevent WAL replay on the replica, replica server resource contention from heavy read queries, and network bandwidth limitations between the primary and replica. Monitor WAL generation rate on the primary and replay rate on replicas to identify the bottleneck.
Can I use replicas for backup instead of a separate backup process?
Replicas provide high availability but are not a substitute for proper backups. Replicas replicate all changes including accidental deletions and data corruption. A dropped table on the primary is immediately replicated to all replicas. Maintain separate point-in-time recovery backups using pg_basebackup with WAL archiving or equivalent MySQL backup tools that allow restoring to a specific timestamp before the data loss event.
Does replication work across different PostgreSQL versions?
PostgreSQL supports streaming replication between consecutive major versions during upgrades, where the replica runs the newer version. However, this is intended as a temporary state during rolling upgrades, not a permanent configuration. For ongoing replication, all servers should run the same major version to avoid compatibility issues with WAL format changes and feature differences that could cause replication failures.