Shard a Live Database Without Downtime, stage 8 of 9: change it
Queries that cross shards
Rows for different workspaces now live in different databases.
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
0 of 2 checks done
Sharding makes every operation inside the shard key cheaper and every operation across keys harder:
- Queries across shards become scatter-gathers, or need an index maintained on write.
- Uniqueness across shards needs one place that owns it.
- Transactions across shards become multi-step workflows.
Check
Users are sharded by workspace. Each shard has a unique index on email. Are emails globally unique?Think first
A user moves a page to a workspace on another host. Why can't that be one local transaction any more?