SysDesignPrep.com
Study guide 46 of 183

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

SignalLikely step
Slow queries, high CPU from scansindexes and query fixes
Connection errors, memory pressure from connectionspooling
CPU or I/O saturated, reads dominatereplicas and caching
Huge tables, slow deletes, retention needspartitioning
Write throughput or storage beyond one primarysharding 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.

Open in your browser to sign in

Google does not allow sign-in inside this app's built-in browser. Open this page in Safari and sign in there. The link opens this same page.

Tap the ⋯ or share button at the top or bottom of the screen, then Open in browser. Or copy the link and paste it into Safari.