- Sources: PlanetScale blog, HN discussion
- Summary: A PlanetScale engineering post dated 2026-07-15 explains horizontal sharding as the way to scale a relational database past a single primary. The worked example spreads one petabyte across 256 shards of three servers each (one primary and two replicas, 768 servers total, about four terabytes per shard) serving millions of queries per second. A proxy router layer parses each SQL query, routes it to the correct shard by a hash of the sharding key, plans and aggregates cross-shard queries, and sits behind a network load balancer, so an application connects through a single connection string and never sees the 768 servers. The post points to Vitess for MySQL and PlanetScale's Neki for Postgres as the systems that implement this pattern.
- Why it matters: It lays out the routing and cross-shard-query mechanics teams hit once a single-primary database runs out of write throughput and storage headroom.
send feedback on this story