Replication, Partitioning & Sharding

Scale reads with replicas, split data safely with partitions and shards, and plan for consistency, hot keys and operational recovery.

These solve different bottlenecks

Primary writes replicate to read replicas; partitioning and sharding split large data by a stable keyPrimary writes replicate to read replicas; partitioning and sharding split large data by a stable key

TechniqueMain problem solvedStill hard
Read replicaPrimary read loadReplication lag and failover
Table partitionOne huge table/query maintenanceCross-partition queries and pruning
ShardingOne machine cannot hold write/data loadRouting, rebalancing and cross-shard work

Replication: more reads, not instant truth

Write to the primary; replicas apply changes later. After a profile update, reading a replica may show the old value. Use read-your-own-write routing, a short primary stickiness period, or a product experience that explicitly accepts delay. Monitor replica lag and promote only a safe, current replica during failover.

Partitioning: one logical table, smaller pieces

Partition by time for events/orders, range for geography, or list for known groups. Good partitioning enables pruning: a March query reads March partitions, not every year. Choose a partition key that matches retention and common filters. Too many tiny partitions add planning and operational overhead.

Sharding: independent database owners

A shard key routes each row to one shard: tenant_id is often safer than an arbitrary hash when tenancy is a hard isolation boundary. Avoid hot keys such as one celebrity account. Plan the router, shard map, ID generation, backups, rebalancing and degraded behavior before splitting.

Cross-shard joins and transactions are expensive. Prefer colocating related data, asynchronous aggregates and explicit workflows. Do not shard because it sounds scalable; exhaust indexing, caching, replicas and partitioning where appropriate first.