Part 5 · 3 chapters · ~18 min

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.

13

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.

choosing a shard key
  1. Most queries (and every transaction) should include it, so they touch one shard.
  2. It should spread load evenly: high cardinality, no single giant value (one huge merchant).
  3. Related data should share it (co-location): accounts, transfers and ledger entries of one tenant together.
  4. It should rarely change: changing a row's shard key means moving the row.
  5. For a multi-tenant product, the tenant id is usually right; for consumer apps, the user id.
SHARDING STRATEGIES
how a key finds its shard, and what rebalancing costs in each
shard keymerchant_idrouterhash · range · directory · geoshard 1shard 2shard 3shard 4
swipe the figure sideways, or tap expand for full screen
1/5
hash
shard = hash(key) mod N (or a consistent hash ring, or fixed slots). Even spread, no hot ranges; range queries across keys hit every shard. Adding shards moves data unless you use many fixed slots (Redis: 16,384; Vitess: keyspace ranges).
even spread; range scans hit every sharduse many virtual slots to rebalance cheaply
14

Resharding live, cross-shard queries and transactions

code
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 needoptions
a query without the shard keyscatter-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 shardsavoid by co-location; else 2PC (Vitess has it, with costs) or a saga with compensations (Distributed Systems part 8)
reporting across everythingstream all shards into a warehouse (part 7) instead of querying shards
15

Vitess, Citus and where sharding lives

Vitess (MySQL)Citus (Postgres)
routervtgate: parses SQL, uses VSchema and vindexes to routecoordinator node plans distributed queries
per shardvttablet beside each mysqld (health, query rewriting, pooling)worker nodes holding shards (ordinary Postgres tables)
mappingvindexes: hash, lookup (a global index), custom functionsdistribution column; reference tables copied to every node
reshardingVReplication workflows (Reshard, MoveTables) with verificationshard rebalancer moves shards between nodes online
used byYouTube origin, Slack, GitHub, PlanetScalemulti-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.