Sharding Strategies: Partition Keys, Distribution, and Rebalancing
Sharding moves the hard problems out of the database and into your schema. We walk through the mechanics: choosing a partition key, hashing vs. range distribution, and how data physically lands on shards. Then we look at what breaks — cross-shard joins, distributed transactions, hot keys, and rebalancing when you add a shard. We compare the real strategies: hash sharding for even distribution, range sharding for locality, and directory-based routing — and what each costs at the application layer. The honest take: sharding is a last resort with permanent costs, and the key choice is made before the first insert.
Topics covered:
- Partition keys and distribution functions
- Hash vs. range vs. directory strategies
- Cross-shard queries and distributed transactions
- Hot keys, rebalancing, and operational costs
Related articles
Horizontal Partitioning and Shard Keys
Range vs hash partitioning, shard key selection, and the hidden complexities of distributing data across nodes.
Replication Lag and Consistency Guarantees
Async vs sync replication, read-after-write consistency, and the physical limits of replicating data across nodes.
Why Your Query Is Slow (And It's Not the Index)
Buffer pool misses, WAL contention, lock waits, and connection pooling — the non-index causes of database slowness.
More in Databases
B-Trees: The Shape of Databases
Why every major database is a tree shaped like a disk page — and how to read your index's health from its shape.
WatchPostgres Internals Tour: Processes, Buffer Pool, WAL, and MVCC
A guided tour of PostgreSQL internals — process model, buffer manager, WAL, and MVCC — the mechanisms that make Postgres behave the way it does.
DetailsSQL Joins and Execution Plans: Nested Loop, Hash, and Merge
How the database executes joins — nested loop, hash join, and merge join — and how to read execution plans to see which strategy your query gets.
DetailsConnection Pools, Database-Side: What a Connection Really Costs
What actually happens to your database when connections pile up — the connection lifecycle, pool sizing math, and why max_connections is not a tuning knob.
DetailsStorage Engines: LSM-Trees vs. B-Trees
LSM-trees vs. B-trees — how each storage engine writes, compacts, and reads, and what that means for write and read amplification in your workload.
DetailsReplication Explained: WAL Shipping, Lag, and Failover
How database replication actually works — the transaction log, the lag, and the failure modes of synchronous and asynchronous replication in production.
DetailsTransactions and Isolation Levels: ACID, MVCC, and Anomalies
What transactions actually guarantee — ACID mechanics, MVCC, and the real behavior behind each isolation level, demonstrated with concrete anomalies.
DetailsDatabase Indexes Visualized: B-Trees, Covering Indexes, and Planner Decisions
How database indexes actually work — B-trees, hash indexes, covering indexes, and when the planner will or won't use the index you made.
DetailsQuery Optimizer Internals: From Parse Tree to Execution Plan
What happens inside a query optimizer — parse, rewrite, join ordering, cost models, and how the planner decides the plan your query gets.
DetailsDepth, delivered weekly
One technical dispatch a week — articles and episode notes before they go public.
One technical dispatch per week. No noise.