Part 9 · 11 chapters · ~20 min

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.

89

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.

FormatLogsSizeProblem
STATEMENTThe SQL textTinyUnsafe. Non-deterministic functions (NOW(), UUID(), RAND()), LIMIT without ORDER BY, and triggers can all produce different results on the replica.
ROWBefore and after images of each changed rowLargeThe default since 5.7, and correct. A single UPDATE touching a million rows writes a million row images.
MIXEDStatement, switching to row when unsafeMediumA reasonable compromise, but the switching rules are subtle. Prefer ROW and be explicit.
binlog_row_image reduces the cost
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.
run it
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
90

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:

code
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 PREPARED in redo and present in the binlog → it committed. Roll it forward.
  • Transaction is PREPARED in 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.

say this and the interview shifts
“MySQL has two write-ahead logs — InnoDB's redo and the server's binlog — so every commit is an internal XA two-phase commit between them. That is a real throughput cost, it is why the double-one configuration exists, and it is a structural consequence of the storage engine API being a hard boundary. Postgres has one WAL and does not pay this.”
91

The replication pipeline

Four moving parts, three of them threads. Knowing which one is behind is the whole of replication debugging.

run it
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
async replication
binlog → relay log → apply
swipe the figure sideways, or tap expand for full screen
1/6
commit
A transaction commits on the primary. Everything that follows is asynchronous — the client has already been told it succeeded.
92

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.

code
-- 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.

errant transactions
If anyone writes directly to a replica, that replica generates a GTID with its own UUID that the primary has never seen. Promote a different server later and the topology can never converge: the old replica has a transaction nobody else has, and it will be sent to everyone on the next failover, potentially applying a stray change months later.

Prevention: super_read_only=ON on every replica, always. Not read_only — that still permits users with SUPER, which includes most admin accounts.
run it
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;
93

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 pointPrimary waits untilOn primary crash
AFTER_SYNC (default)The replica acknowledges, before the primary commits in InnoDBNo phantom reads. The transaction was never visible to clients on the primary, and the replica has it. Safe to fail over.
AFTER_COMMITThe replica acknowledges, after the primary commitsA 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.
semi-sync degrades to async silently
If no replica acknowledges within 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.
94

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:

PolicyParallelism ruleEffectiveness
DATABASETransactions on different schemasUseless for the common single-schema application.
LOGICAL_CLOCKTransactions that committed together on the primary — same group commit — cannot have conflicted, so they can be applied in parallelGood. Depends on the primary actually forming groups.
WRITESETTransactions whose modified rows do not overlap, computed from row hashesMuch better. Finds parallelism even in transactions that committed at different times on the primary.
the configuration that fixes most lag
Set 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.
run it
-- 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.
95

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.
measure it with a heartbeat instead
Write a row with the current timestamp on the primary once a second; on the replica, compare that row's value to the replica's clock. 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.
run it
# 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 != '';
96

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.

ModeWritesUse for
Single-primary (default)One member accepts writes; automatic election on failureMost deployments. You get automatic failover without external tooling.
Multi-primaryAny member accepts writesRarely 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.

97

Router, ProxySQL, Vitess

Applications should not hold a list of database hostnames. Three layers sit in front, at increasing levels of ambition.

LayerDoesDoes not
MySQL RouterRoutes to the current primary of an InnoDB Cluster; reads to secondaries. Lightweight, deploy beside the app.No query rewriting, no caching, no sharding.
ProxySQLQuery-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.
VitessFull 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.
the progression
Router when you have an InnoDB Cluster and want failover handled. ProxySQL when connection count or read/write splitting is the problem — it is the pragmatic choice for most teams. Vitess when one machine genuinely cannot hold the write volume, which is later than most people think. The Scaling module covers the decision in full.
98

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.

DefenceMechanism
Fencing / STONITHForcibly power off or network-isolate the old primary before promoting. The only truly reliable answer.
QuorumOnly a majority partition may elect. This is what Group Replication does natively.
VIP / proxy cutoverMove the address atomically so clients cannot reach the old primary even if it lives.
super_read_onlySet 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.

automatic failover is a decision, not a default
Automatic failover trades a fast recovery from real failures against the risk of a spurious failover during a network blip. Whether that trade is right depends on how costly divergence is for your data. For a ledger, many teams deliberately choose manual promotion with good alerting. Decide deliberately, and write down which you chose and why.
99

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

  1. Read from the primary for N seconds after a write. Crude, common, and it works. Costs primary capacity for a window.
  2. Sticky routing per session after a write. Same idea, scoped to the session. Better, still coarse.
  3. 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.
run it
-- 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.