Write Scaling and Sharding
Functional decomposition before sharding, hash, range, directory and geo strategies with rebalancing, choosing a shard key, resharding live, cross-shard queries and transactions, Vitess and Citus internals, app versus proxy versus database-layer sharding, and global secondary indexes.
Decompose first, then choose a strategy and a key
Before sharding one table across machines, split by function: move the ledger, notifications, analytics and sessions into separate databases owned by separate services. Each gets its own capacity and its own failure domain, with no cross-shard queries inside a function.
- Most queries (and every transaction) should include it, so they touch one shard.
- It should spread load evenly: high cardinality, no single giant value (one huge merchant).
- Related data should share it (co-location): accounts, transfers and ledger entries of one tenant together.
- It should rarely change: changing a row's shard key means moving the row.
- For a multi-tenant product, the tenant id is usually right; for consumer apps, the user id.
Resharding live, cross-shard queries and transactions
live resharding, the five steps 1 copy backfill rows for the moving key range into the target shard (consistent snapshot) 2 stream apply ongoing changes from the source's binlog / WAL to the target until lag ≈ 0 3 verify compare row counts and checksums per range; diff samples 4 cut over briefly stop writes for the range (seconds), apply the tail, switch routing, resume 5 clean up keep the source read-only for a while, then delete the moved rows
| cross-shard need | options |
|---|---|
| a query without the shard key | scatter-gather to every shard (slow, fan-out tax), or a global secondary index (a separate table keyed by the other column, maintained asynchronously) |
| a transaction across shards | avoid by co-location; else 2PC (Vitess has it, with costs) or a saga with compensations (Distributed Systems part 8) |
| reporting across everything | stream all shards into a warehouse (part 7) instead of querying shards |
Vitess, Citus and where sharding lives
| Vitess (MySQL) | Citus (Postgres) | |
|---|---|---|
| router | vtgate: parses SQL, uses VSchema and vindexes to route | coordinator node plans distributed queries |
| per shard | vttablet beside each mysqld (health, query rewriting, pooling) | worker nodes holding shards (ordinary Postgres tables) |
| mapping | vindexes: hash, lookup (a global index), custom functions | distribution column; reference tables copied to every node |
| resharding | VReplication workflows (Reshard, MoveTables) with verification | shard rebalancer moves shards between nodes online |
| used by | YouTube origin, Slack, GitHub, PlanetScale | multi-tenant SaaS, real-time analytics, Azure Cosmos DB for PostgreSQL |
Where to shard: in the app (full control, every service must know), in a proxy (Vitess, ProxySQL rules: transparent to apps), or in the database (Citus, NewSQL systems: SQL looks unchanged). The further down the stack, the less application change and the more you depend on one product's limits.