Shard a Live Database Without Downtime, stage 5 of 9: break it
Shards that are slightly in the past
Here is the backfill and catch-up code. Find every line that contributes to these three problems.
System so far· 7 parts
Select a component to see what it is responsible for and which state it owns.
- 1Users → Application servers: Requests
- 2Application servers → Monolith Postgres: Reads and writes until cutover
- 3Application servers → Write audit log: Log every write
- 4Backfill and catch-up → Monolith Postgres: Copy existing rows
- 5Backfill and catch-up → Write audit log: Replay logged writes
- 6Backfill and catch-up → Shard hosts: Write rows by shard
- 7Application servers → Connection poolers: Queries by shard
- 8Connection poolers → Shard hosts: Schema on host
- Request / response
- Asynchronous
What you need to know
Backfill and catch-up run at the same time, against the same rows. The backfill copies a row as it was when scanned; catch-up applies newer writes from the log. Whichever writes last wins, unless writes are conditional.
The rule that works is newest version wins: write only if the row is absent or the stored version is older.
Think first
Catch-up writes version 7 of a block to its shard. A moment later the backfill reaches the same block, which it scanned at version 5, and upserts it unconditionally. What's on the shard?Check
Catch-up always starts reading the audit log from position 0. A deploy restarts it on day 4. What happens?The backfill also competes with production for the primary's CPU and disk. Throttle it, back off when production latency rises, or read from a replica or snapshot. See Backpressure and capacity.