Replication & HA
Everything so far has been one machine. Replication adds a second, and with it a set of problems that have no single-server equivalent: two logs that must agree, a replica that is always slightly in the past, and the question of what "committed" means when the machine holding your commit might be about to die.
Binlog formats
The binary log is a server-layer log of changes, separate from InnoDB's redo log. It drives replication and point-in-time recovery. Three formats, and the choice has real consequences.
| Format | Logs | Size | Problem |
|---|---|---|---|
STATEMENT | The SQL text | Tiny | Unsafe. Non-deterministic functions (NOW(), UUID(), RAND()), LIMIT without ORDER BY, and triggers can all produce different results on the replica. |
ROW | Before and after images of each changed row | Large | The default since 5.7, and correct. A single UPDATE touching a million rows writes a million row images. |
MIXED | Statement, switching to row when unsafe | Medium | A reasonable compromise, but the switching rules are subtle. Prefer ROW and be explicit. |
FULL (default) logs every column of both images. MINIMAL
logs only the primary key plus changed columns — often a large saving on wide
tables. NOBLOB omits unchanged BLOB/TEXT columns. Set
MINIMAL unless a downstream consumer (a CDC pipeline, for
instance) needs the full before-image.SELECT @@binlog_format, @@binlog_row_image, @@log_bin; SHOW BINARY LOGS; SHOW BINLOG EVENTS IN 'binlog.000003' LIMIT 20; # ROW events are opaque in SHOW BINLOG EVENTS. # In a shell, decode them properly: # mysqlbinlog --verbose --base64-output=DECODE-ROWS binlog.000003 # # Row events then appear as readable pseudo-SQL: # ### UPDATE `shop`.`orders` # ### WHERE # ### @1=4471 /* the BEFORE image */ # ### @4='pending' # ### SET # ### @1=4471 /* the AFTER image */ # ### @4='paid' # Extract a window of time — the basis of PITR (ch 103): # mysqlbinlog --start-datetime="2026-09-21 14:00:00" \ # --stop-datetime="2026-09-21 14:05:00" binlog.000003
Two logs, and the XA commit between them
This is the consequence of the pluggable engine architecture from chapter 6, and it is the answer to a genuinely good interview question.
InnoDB has its redo log. The server has its binary log. A transaction must appear in both or neither — if it were in redo but not binlog, the replica would silently diverge from the primary forever; if in binlog but not redo, the replica would have data the primary lost.
So every commit runs an internal XA two-phase commit between the two logs:
1. InnoDB PREPARE -- write redo, mark the trx PREPARED, fsync 2. Binlog WRITE -- write the events 3. Binlog SYNC -- fsync the binlog ← THE commit point 4. InnoDB COMMIT -- mark committed (no fsync needed)
Step 3 is the moment the transaction becomes real. Recovery uses that fact:
- Transaction is
PREPAREDin redo and present in the binlog → it committed. Roll it forward. - Transaction is
PREPAREDin redo but absent from the binlog → it did not commit. Roll it back.
This is also why full durability needs two settings at 1 —
innodb_flush_log_at_trx_commit=1 and
sync_binlog=1. One alone leaves a window where the two logs can
disagree after a crash.
The replication pipeline
Four moving parts, three of them threads. Knowing which one is behind is the whole of replication debugging.
SHOW REPLICA STATUS\G -- (SHOW SLAVE STATUS on 8.0.21 and earlier) -- The fields that matter, and what they mean: -- -- Replica_IO_Running must be Yes -- Replica_SQL_Running must be Yes -- Seconds_Behind_Source ← misleading. see chapter 95. -- Last_IO_Error network / auth problems -- Last_SQL_Error a statement failed to apply -- -- Compare these two pairs to locate the bottleneck: -- Source_Log_File + Read_Source_Log_Pos ← IO thread got this far -- Relay_Source_Log_File + Exec_Source_Log_Pos ← applier got this far -- -- read ≈ exec, both behind source → NETWORK is the bottleneck -- read far ahead of exec → APPLIER is the bottleneck -- The modern equivalent, with more detail: SELECT * FROM performance_schema.replication_connection_status\G SELECT * FROM performance_schema.replication_applier_status_by_worker\G
GTIDs
Classic replication positions a replica by file and offset:
binlog.000042, position 19384756. Those coordinates
are meaningless on any other server, so promoting a new primary meant
hand-calculating where every other replica should resume — error-prone and slow
at exactly the wrong moment.
A Global Transaction Identifier is
source_uuid:transaction_number, assigned once and carried
everywhere. It means the same thing on every server in the topology.
-- a GTID set: one uuid, ranges of transaction numbers 3E11FA47-71CA-11E1-9E33-C80AA9429562:1-104721, 8A2C1B33-91DD-11E5-A0B1-B8CA3A6A1B24:1-38
With SOURCE_AUTO_POSITION=1, a replica simply tells its new
primary which GTIDs it has already executed, and the primary sends everything
else. Failover becomes a single command.
Prevention:
super_read_only=ON on every replica, always. Not
read_only — that still permits users with SUPER, which includes
most admin accounts.SELECT @@gtid_mode, @@enforce_gtid_consistency, @@super_read_only; SELECT @@global.gtid_executed; -- everything applied here SELECT @@global.gtid_purged; -- applied but binlogs since deleted -- Find errant transactions: GTIDs on the replica that did -- NOT come from the primary's UUID. SELECT GTID_SUBTRACT( @@global.gtid_executed, 'PRIMARY-UUID-HERE:1-999999999' ) AS errant_gtids; -- non-empty = someone wrote to this replica. investigate. -- Repoint a replica at a new primary — the whole failover: -- STOP REPLICA; -- CHANGE REPLICATION SOURCE TO -- SOURCE_HOST='new-primary', SOURCE_AUTO_POSITION=1; -- START REPLICA;
Semi-synchronous replication
Default replication is asynchronous: the primary commits and returns to the client without waiting for any replica. Fast, and it means a primary crash can lose transactions the client was told had succeeded.
Semi-sync makes the primary wait for at least
rpl_semi_sync_source_wait_for_replica_count replicas to
acknowledge receipt — receipt into the relay log, not application.
| Wait point | Primary waits until | On primary crash |
|---|---|---|
AFTER_SYNC (default) | The replica acknowledges, before the primary commits in InnoDB | No phantom reads. The transaction was never visible to clients on the primary, and the replica has it. Safe to fail over. |
AFTER_COMMIT | The replica acknowledges, after the primary commits | A window exists where other clients on the primary can read a transaction that the replica may not have. Those reads become phantom reads after failover. |
rpl_semi_sync_source_timeout (default 10 seconds), the primary
gives up and switches to asynchronous so writes are not
blocked forever. That is usually the right behaviour — but it means your
durability guarantee can disappear without anyone noticing.
Monitor
Rpl_semi_sync_source_status. If it is 0, you are running
async and believing you are not.Multi-threaded appliers
The classic cause of replica lag: a primary with 32 cores accepting writes from hundreds of connections, and a replica applying them with one thread. The replica cannot keep up by construction.
Parallel application needs to know which transactions may safely run concurrently. Two schemes:
| Policy | Parallelism rule | Effectiveness |
|---|---|---|
DATABASE | Transactions on different schemas | Useless for the common single-schema application. |
LOGICAL_CLOCK | Transactions that committed together on the primary — same group commit — cannot have conflicted, so they can be applied in parallel | Good. Depends on the primary actually forming groups. |
WRITESET | Transactions whose modified rows do not overlap, computed from row hashes | Much better. Finds parallelism even in transactions that committed at different times on the primary. |
binlog_transaction_dependency_tracking=WRITESET on the
primary (it computes the dependency information that goes into the
binlog) and replica_parallel_type=LOGICAL_CLOCK with a
replica_parallel_workers of 8–16 on the replica. Keep
replica_preserve_commit_order=ON so replicas remain consistent
snapshots.
This combination routinely takes a replica from hours behind to seconds behind, and it is the first thing to try before considering sharding for write volume.
-- ON THE PRIMARY:
SET GLOBAL binlog_transaction_dependency_tracking = 'WRITESET';
SET GLOBAL transaction_write_set_extraction = 'XXHASH64'; -- pre-8.0.26
-- ON THE REPLICA:
SET GLOBAL replica_parallel_type = 'LOGICAL_CLOCK';
SET GLOBAL replica_parallel_workers = 16;
SET GLOBAL replica_preserve_commit_order = ON;
-- requires STOP REPLICA / START REPLICA to take effect
-- Are the workers actually being used?
SELECT worker_id, service_state,
last_applied_transaction,
applying_transaction
FROM performance_schema.replication_applier_status_by_worker;
-- If only worker 1 ever has work, parallelism is not happening —
-- check that the PRIMARY has WRITESET enabled.Measuring lag truthfully
Seconds_Behind_Source is the number everyone watches and it is
unreliable. It is computed as the difference between the replica's clock and
the timestamp of the event currently being applied — which breaks in several
ordinary situations:
- It reads 0 when the IO thread is broken. No events are arriving, so there is nothing to be behind on. Total failure reports as perfect health.
- A long-running transaction on the primary carries the timestamp of when it started, so lag appears to jump by the transaction's duration.
- In a chained topology it measures against the intermediate server, not the true origin.
- Clock skew between the machines corrupts it directly.
pt-heartbeat does exactly this. It measures the property you
actually care about — how stale is the data I would read here — and it
correctly reports a broken IO thread as ever-growing lag rather than zero.# The honest measurement.
# ON PRIMARY:
# pt-heartbeat --update --daemonize -D percona --table=heartbeat
# ON REPLICA:
# pt-heartbeat --monitor -D percona --table=heartbeat
-- Or the same idea with no extra tooling:
SELECT TIMESTAMPDIFF(SECOND, ts, NOW()) AS real_lag_seconds
FROM heartbeat WHERE id = 1;
-- 8.0 also exposes precise timestamps per transaction:
SELECT
TIMESTAMPDIFF(MICROSECOND,
last_applied_transaction_original_commit_timestamp,
last_applied_transaction_end_apply_timestamp) / 1e6 AS apply_lag_s
FROM performance_schema.replication_applier_status_by_coordinator;
-- and the queue depth — how much is waiting to be applied:
SELECT COUNT(*) FROM performance_schema.replication_applier_status_by_worker
WHERE applying_transaction != '';Group Replication and InnoDB Cluster
A different model. Instead of a designated primary streaming to passive replicas, a group of servers agrees on the order of transactions using a consensus protocol derived from Paxos.
The mechanism is certification. A transaction executes locally, then its write set is broadcast to the group. Each member independently checks whether it conflicts with any concurrent transaction already certified. Same input, same deterministic rule, so every member reaches the same verdict without further coordination. Conflicts cause a rollback on the originating member.
| Mode | Writes | Use for |
|---|---|---|
| Single-primary (default) | One member accepts writes; automatic election on failure | Most deployments. You get automatic failover without external tooling. |
| Multi-primary | Any member accepts writes | Rarely worth it. Certification conflicts under contention cause rollbacks that applications must handle, and hot rows perform badly. |
InnoDB Cluster is the packaged product: Group Replication plus MySQL Router for connection routing plus MySQL Shell for administration. Flow control throttles fast members so slow ones can keep up, which means one struggling member slows the whole group — a trade-off worth knowing before you adopt it.
Router, ProxySQL, Vitess
Applications should not hold a list of database hostnames. Three layers sit in front, at increasing levels of ambition.
| Layer | Does | Does not |
|---|---|---|
| MySQL Router | Routes to the current primary of an InnoDB Cluster; reads to secondaries. Lightweight, deploy beside the app. | No query rewriting, no caching, no sharding. |
| ProxySQL | Query-rule-based routing, read/write split, connection multiplexing (thousands of app connections onto few backend ones), query caching, rewriting, mirroring. | No sharding of a single logical table. |
| Vitess | Full horizontal sharding. vtgate parses queries and routes by a vindex; vttablet manages each MySQL instance; online resharding with cutover. | A large operational commitment. It is a distributed system you now run. |
Failover and split brain
The hard problem in failover is not promoting a replica. It is being certain the old primary is really dead — because a network partition looks exactly like a dead server from the other side.
If you promote a new primary while the old one is alive and still accepting writes, you have split brain: two servers taking writes, diverging, with no automatic way to merge them afterwards.
| Defence | Mechanism |
|---|---|
| Fencing / STONITH | Forcibly power off or network-isolate the old primary before promoting. The only truly reliable answer. |
| Quorum | Only a majority partition may elect. This is what Group Replication does natively. |
| VIP / proxy cutover | Move the address atomically so clients cannot reach the old primary even if it lives. |
super_read_only | Set on demotion so a returning old primary cannot accept writes. |
Orchestrator is the standard tool for classic topologies: it discovers the replication graph, detects failure by cross-checking from multiple vantage points (rather than trusting one health check), and can promote with configurable hooks for fencing.
Read-your-writes on replicas
The bug every team with read replicas ships at least once. A user updates their profile, the write goes to the primary, the redirect issues a read that lands on a replica 200ms behind, and the user sees their old data. They update again. Support gets a ticket saying "it doesn't save".
The options, worst to best
- Read from the primary for N seconds after a write. Crude, common, and it works. Costs primary capacity for a window.
- Sticky routing per session after a write. Same idea, scoped to the session. Better, still coarse.
- Wait for the replica to catch up to a known position. The correct answer. Capture the GTID of your write, then have the replica wait for exactly that transaction before reading.
-- ON THE PRIMARY, right after your write's COMMIT: SELECT @@session.gtid_executed INTO @write_gtid; -- carry @write_gtid in the session / a cookie / the request context -- ON THE REPLICA, before the read that must see it: SELECT WAIT_FOR_EXECUTED_GTID_SET(@write_gtid, 1); -- returns 0 → the replica has it, safe to read -- returns 1 → timed out after 1s, fall back to the primary -- THEN do the read. SELECT * FROM profiles WHERE user_id = ?; -- This gives real read-your-writes with a bounded wait, and -- degrades safely to a primary read instead of showing stale data. -- ProxySQL can do this transparently via its GTID-aware routing.
That is the distributed picture. Part 10 is everything involved in keeping this running day to day: schema changes, backups, observability and the config that matters.