DB Sharding & Partitioning
Studies in this cluster, in series order. Each one keeps its own URL.
Databases
Indexes, isolation, storage engines, shard and partition keys, and zero-downtime migrations you can ship without a maintenance window.
DB Sharding & Partitioning
6 studies- 1.Database Sharding & Partitioning — Keys, Hotspots & RebalancingVertical scaling and a single primary eventually hit CPU, IOPS, storage, or write-throughput walls. Sharding splits data across nodes for independent write capacity. Table partitioning splits one table inside an engine for prune and maintenance. This hub is the map for when each is enough, how a key avoids hotspots, why scatter-gather explodes cost, and how online resharding works.
- 2.Partition Strategies — Range, Hash, List & CompositeBefore you invent an application shard layer, master in-engine partitioning: range, hash, list, and composite. The same vocabulary maps onto Vitess, Citus, and TiDB. Pick the strategy from the access pattern. A wrong one creates hotspots, immovable partitions, or queries that never prune.
- 3.Shard Keys — Hotspots, Fan-out Queries & CardinalityThe shard key is the product decision you cannot easily undo. Evaluate cardinality, skew, stability, and whether the hot transaction stays on one shard. Low-cardinality columns, monotonic time, and celebrity keys become hotspots. Queries that omit the key become fan-out.
- 4.Cross-Shard Queries — Scatter-Gather, Aggregations & AvoidanceThe day you shard, any query without the shard key becomes a scatter-gather. Latency is the slowest shard plus the merge, and every such query multiplies load by N. Merge sums and counts, never averages of averages. Keep scatter off the OLTP hot path.
- 5.Resharding & Rebalancing — Dual-Write, Online Moves & Consistent Hash RingsIf you cannot reshard, you cannot shard safely. Online moves are dual-write plus backfill plus cutover, engine range splits, or a consistent-hash vnode remap. Modulo-N moves most keys when N changes. A ring moves a slice. This is a data-plane move, not a schema migration.
- 6.Sharding Failure Modes — Orphans, Partial Commits, Dual-Writes & IdempotencySharding turns single-node ACID into a distributed problem. Orphan rows, partial multi-shard commits, dual-write divergence, lost updates at cutover, and missing idempotency keys are the failures seniors name. Prefer one shard per transaction. Use a saga when you cannot. Two-phase commit is the rare exception.