Scaling a relational database
How far one Postgres or MySQL database goes and how to stretch it: query and index tuning, connection pooling, read replicas and lag, caching, partitioning tables, vacuum and bloat, vertical scaling, splitting by function, and when to shard or move to a distributed SQL database.
Reading is half of it. See this used in a real interview: walk through Design a URL Shortener →
A single well-run relational database handles far more than many interview answers assume: tens of thousands of transactions per second and terabytes of data on modern hardware. Jumping straight to sharding adds enormous complexity. Interviewers like to see the progression: what you do first, what each step buys, and the signals that tell you it is time for the next one.
Step 1: make the queries cheap
- Indexes for every frequent query; composite indexes in the right column order; covering indexes for hot reads. See database indexing.
- Read query plans (
EXPLAIN ANALYZE); look for sequential scans on big tables and sorts that spill to disk. - Remove N+1 query patterns from application code; batch writes.
- Use keyset pagination, not large offsets. See pagination.
- Keep transactions short; long ones hold locks and block cleanup.
This step alone often buys an order of magnitude.
Step 2: connection pooling
Each Postgres connection is a process with its own memory; thousands of connections from autoscaled app servers exhaust the server. A pooler (PgBouncer, or a managed proxy) multiplexes many client connections over a small number of database connections, typically in transaction mode. Size the database side with queueing theory: a few times the number of cores is usually plenty.
Step 3: scale up
Bigger machines with more memory (so the working set fits in cache) and faster NVMe storage go a long way. Vertical scaling is boring and effective until the largest instance is not enough, or a single machine is too risky.
Step 4: read replicas and caching
- Read replicas take read traffic; the primary handles writes. Replication is usually asynchronous, so replicas lag by milliseconds to seconds.
- Handle lag: read your own writes from the primary (for a short period after writing, or for that user), and send only lag-tolerant reads to replicas. See replication and consistency.
- Cache hot, read-mostly data in Redis or in-process to remove most reads entirely. See caching.
Step 5: partition big tables
Table partitioning within one database (by time or by key range) keeps indexes small and makes retention trivial: dropping an old month is instant, while deleting millions of rows is slow and bloats the table. Queries that include the partition key touch only relevant partitions. See time-series data.
Keep it healthy
- Vacuum and bloat (Postgres): updates and deletes leave dead row versions that autovacuum must clean; long transactions and badly tuned autovacuum cause bloat and slow queries. Monitor dead tuples and transaction age.
- Replication slots and lagging replicas can fill the disk with retained WAL.
- Migrations must not lock hot tables. See online schema migrations.
- Backups and tested restores, with point-in-time recovery. See backups and disaster recovery.
Step 6: split by function
Move separate domains into separate databases: users, orders, analytics, logs. Each grows independently. Move analytics and reporting off the primary to a warehouse fed by change data capture, since big scans compete with transactional traffic.
Step 7: shard, or go distributed
When write volume or data size exceeds one primary even after all of the above:
- Application-level sharding: split rows by a shard key (tenant id, user id) across many databases, with a routing layer. Cross-shard queries and transactions become hard. See sharding and partitioning.
- Sharding extensions and proxies: Citus for Postgres, Vitess for MySQL, which handle routing and resharding.
- Distributed SQL: CockroachDB, Spanner, YugabyteDB and similar give a SQL interface with automatic sharding and consensus replication, at the cost of higher write latency and different performance characteristics.
Signals for each step
| Signal | Likely step |
|---|---|
| Slow queries, high CPU from scans | indexes and query fixes |
| Connection errors, memory pressure from connections | pooling |
| CPU or I/O saturated, reads dominate | replicas and caching |
| Huge tables, slow deletes, retention needs | partitioning |
| Write throughput or storage beyond one primary | sharding or distributed SQL |
In the interview
Start with one primary plus replicas and say what it can handle, using your estimates. For a URL shortener with tens of thousands of writes per second at peak, a single tuned primary may be enough with caching for reads; say so, then explain the next steps if growth continues. See Design a URL Shortener. For money and inventory, the relational database's transactions are a feature worth keeping as long as possible. See Design a Payment System and Design Ticketmaster.
Checklist
- Indexes, query plans, short transactions and keyset pagination first.
- Connection pooling sized from concurrency, not client count.
- Vertical scaling before distribution.
- Replicas with lag-aware routing; caching for hot reads.
- Time or range partitioning for large tables and retention.
- Vacuum, WAL and migration hygiene.
- Functional splits, then sharding or distributed SQL when the numbers require it.