Skip to content

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
12345678CLIENTUsersSERVICEApplicationserversDATABASEMonolithPostgresLOG / STREAMWrite audit logWORKERBackfilland catch-upSERVICEConnectionpoolersDATABASEShard hosts

Select a component to see what it is responsible for and which state it owns.

  1. 1Users → Application servers: Requests
  2. 2Application servers → Monolith Postgres: Reads and writes until cutover
  3. 3Application servers → Write audit log: Log every write
  4. 4Backfill and catch-up → Monolith Postgres: Copy existing rows
  5. 5Backfill and catch-up → Write audit log: Replay logged writes
  6. 6Backfill and catch-up → Shard hosts: Write rows by shard
  7. 7Application servers → Connection poolers: Queries by shard
  8. 8Connection poolers → Shard hosts: Schema on host
  • Request / response
  • Asynchronous

What you need to know

0 of 2 checks done
  1. 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.

  2. 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?