Learn Labs
6. Replication

6.6 Technology deep dives


6.1 PostgreSQL streaming replication (single-leader, WAL shipping + logical)

Problem it solves. Keep a byte-identical hot standby for failover and read scaling, with no application changes.

Why not statement-based? Nondeterminism (now(), random(), nextval) and ordering dependencies make it unsafe under concurrency.

Why add logical replication later? WAL shipping is tightly coupled to the storage format, so leader and follower must run the same major version — which makes zero-downtime major upgrades impossible. Logical decoding solves both that and CDC.

How it works internally. The leader streams WAL records over a replication connection; each standby has a replication slot recording the LSN it has consumed, so the leader retains WAL until it's consumed. synchronous_commit controls the durability level (off / local / remote_write / on / remote_apply); synchronous_standby_names picks which standbys are synchronous. Logical replication runs the WAL through an output plugin (pgoutput) that turns physical records into row-level insert/update/delete events.

Deployment. Primary + N standbys; a failover manager (Patroni + etcd/Consul, repmgr, or a cloud provider's). Patroni's real job is fencing — using a distributed lock so only one node believes it's primary.

Monitoring.

  • pg_stat_replication: sent_lsn − replay_lsn per standby, in bytes and in seconds (replay_lag) — both matter; a standby can be byte-close but seconds-behind if replay is blocked
  • Replication slot retained WAL — the #1 way to fill a primary's disk: an inactive slot pins WAL forever
  • pg_stat_database_conflicts — queries on the standby cancelled by replay
  • Sync-standby availability; failover events; timeline ID after promotion

Scaling. Read replicas for reads; cascading replication to reduce load on the primary; but all writes still go through one node.

Backup. Base backup + WAL archiving for PITR (pgBackRest, WAL-G). Replication is not backup — §0.

What actually breaks.

  • An abandoned replication slot fills pg_wal and takes the primary down. This is the single most common Postgres replication outage.
  • max_standby_streaming_delay vs long-running read queries on the standby — either queries get cancelled ("conflict with recovery") or replay stalls and lag grows. You must choose which.
  • Split brain after a network partition if the failover manager lacks proper fencing — the old primary keeps accepting writes.
  • Lost writes on promotion with synchronous_commit = off or an async standby promoted.
  • Timeline divergence — the old primary rejoins on a diverged timeline and needs pg_rewind or a full rebuild.
  • Logical replication doesn't replicate DDL — a schema change on the publisher breaks the subscriber.

6.2 MySQL binlog replication + GTIDs

Problem it solves. Same as above, but with a logical log by default, which makes cross-version replication and CDC first-class.

Why row-based over statement-based? MySQL shipped statement-based first and switched to row-based by default whenever a statement is nondeterministic — the practical concession that statement-based replication cannot be made safe in general.

How it works internally. The primary writes the binlog (row events by default, plus a commit marker per transaction). Replicas have an I/O thread copying the binlog to a local relay log, and one or more SQL applier threads replaying it. GTIDs give each transaction a globally unique source_uuid:txn_id, which makes failover safe: a replica knows exactly which transactions it has, rather than a file+offset that becomes meaningless after a primary change. Semi-sync (rpl_semi_sync_source_wait_for_slave_count) waits for N replicas to receive (not apply) before acking.

Monitoring. Seconds_Behind_Source (crude — it measures the applier, and reads 0 when the I/O thread is broken, so always pair it with SHOW REPLICA STATUS errors), applier thread errors, GTID gaps (gtid_executed vs the primary's), relay log disk usage, semi-sync ack timeouts (which silently downgrade to async).

What actually breaks.

  • The GitHub incident from §1.4 — a lagging replica promoted, its autoincrement counter behind, reusing primary keys that were already referenced by Redis, leaking private data to the wrong users. The lesson generalizes: cross-system consistency breaks catastrophically when the promoted node's counters go backwards.
  • Semi-sync silently degrading to async after a timeout — you think you have durability and you don't.
  • Replication drift with statement-based mode still enabled for some statements.
  • A single-threaded applier falling behind a multi-threaded writer (fixed by parallel replication, which introduces its own ordering subtleties).
  • sync_binlog=0 + innodb_flush_log_at_trx_commit=2 — fast, and loses committed transactions on a power failure.

6.3 Cassandra / ScyllaDB (leaderless, Dynamo-style)

Problem it solves. Multi-region, always-writable storage that tolerates node and region failure without failover, with tunable staleness.

Why not single-leader? Failover time, leader sensitivity to gray failures, and cross-region write latency. Leaderless doesn't distinguish the normal case from the failure case (§4.6), so a slow node degrades nothing.

Why not multi-leader? Reads can be arbitrarily stale; quorums give a middle ground.

How it works internally. Consistent hashing ring; RF replicas per key; per-request consistency level (ONE, QUORUM, LOCAL_QUORUM, EACH_QUORUM, ALL, ANY) implementing the w/r knobs; read repair, hinted handoff, and anti-entropy repair (Merkle-tree based) as the three catch-up mechanisms; LWW at the cell level using wall-clock timestamps; storage is an LSM tree (Ch 4).

Consistency levelw / rEffect
LOCAL_QUORUMw = 2 local, r = 2 localFast, no cross-region wait, may be stale across regions
QUORUMw = 3 of 6, r = 3 of 6Cross-region wait on every operation
EACH_QUORUMquorum in every regionStrongest, slowest
ONE / ANYw = 1 / sloppyHighest availability, weakest guarantee

Deployment. Multi-DC keyspace with NetworkTopologyStrategy; clients pinned to a local coordinator; LOCAL_QUORUM as the default; scheduled repairs (weekly, within gc_grace_seconds).

Monitoring. Per-CL latency and timeout rate; hint count and hint delivery rate (§4.5 — the only real staleness proxy); repair completion time vs gc_grace_seconds (miss this and deleted data resurrects); pending compactions; dropped mutations (a write accepted but silently discarded under load); tombstone-scanned-per-read histograms; clock skew across nodes (§4.4 ⑤ makes this a data-correctness metric, not just an ops metric).

Scaling. Add nodes to the ring; quorums seldom exceed 4-of-7 or 5-of-9 because bigger quorums raise the chance of hitting a slow replica.

Backup. Snapshots (hard links over immutable SSTables) + incremental backups shipped off-node.

What actually breaks.

  • Clock skew silently dropping writes. LWW with wall-clock timestamps means a node with a fast clock can make later writes lose. NTP failure becomes data loss with no error anywhere.
  • Zombie data. If a node misses a delete and repair doesn't run within gc_grace_seconds, the tombstone is purged elsewhere and the deleted row comes back on the next repair.
  • Failed writes that partly succeeded aren't rolled back (§4.4 ④) — "the write returned an error" tells you nothing about whether the data is there.
  • Sloppy quorum / CL=ANY accepting a write that no correct replica will ever serve.
  • Read repair storms after a node returns from a long outage.
  • Using it for a workload needing uniqueness or invariants — LWW cannot express "this username is taken."

6.4 Sync engines / CRDT libraries (Automerge, Yjs, PouchDB/CouchDB, Firestore)

Problem it solves. Sub-frame-latency local reads/writes, offline operation, and real-time multi-user collaboration on the same document.

Why not request/response to a server? Every interaction pays a round trip; every call needs error handling in the UI; offline requires a separate code path.

Why CRDTs rather than LWW? LWW on a document means one user's paragraph silently vanishes. Users notice.

How it works internally. Each character/element gets an immutable unique ID; operations reference the ID of the element they follow, not an index — so no transformation is needed and replicas converge from any delivery order (§3.4). Automerge/Yjs implement RGA/YATA-family sequence CRDTs plus maps, lists, counters, and text, all composable into a JSON-shaped document. Changes are exchanged as compact binary deltas; the full change history is retained to allow merging from any peer.

Deployment. A relay/persistence server (y-websocket, Automerge sync server, CouchDB) that is just a dumb replica, not an authority — which is what makes local-first possible.

Monitoring. Document size growth (the metric — CRDT metadata and tombstones accumulate), sync round-trip latency, peer count per document, merge conflict/sibling counts, client memory.

Scaling. Per-document, not per-dataset. Sync engines assume all data the user may need is downloaded in advance — fine for a user's own files, not for an ecommerce catalog.

What actually breaks.

  • Document bloat. A long-lived collaborative document accumulates tombstones and per-character metadata until it's tens of MB and slow to load. Requires compaction/snapshotting and, sometimes, history truncation.
  • Invariants you cannot express. "No more than 5 items" — concurrent adds will exceed it and your only option is to drop some (§3.4).
  • Merge results that are correct but semantically wrong. Both users' edits preserved, producing a sentence neither wrote.
  • Access control. If every replica has the whole document history, you cannot hide part of it.
  • Schema migration across offline clients — a client that's been offline for six months syncs a document written by three schema versions later.
  • Deleted-data compliance — the change history is, by design, permanent.

On this page