Learn Labs
Replication & Scaling

Scaling & Pooling

Read replicas, PgBouncer, sharding, high availability

Replication gives you copies of your data. Scaling is about deciding what to do with them — and what to do once copies alone aren't enough.

Vertical scaling, and its ceiling

The simplest lever is a bigger machine: more RAM (more room for shared_buffers and OS page cache), faster disks, more CPU cores for parallel query execution. It requires no architecture changes and it's usually the right first move. It also has a hard ceiling — at some point there isn't a bigger single machine to buy, and a single primary can only accept writes as fast as one WAL stream can be fsynced to disk.

Read replicas

If your workload is read-heavy (most web apps are), streaming replicas let you scale reads roughly linearly by adding more of them. Writes still go through the primary — replicas don't help write throughput at all, and application code has to be aware of replication lag: a read immediately after a write, on a replica, can return stale data if it hasn't caught up yet.

Primary
Replica 1serves reads
Replica 2serves reads
Replica 3serves reads

Connection pooling

Postgres backends are not free — each connection is a full OS process with its own memory (including work_mem headroom), and the default max_connections (100) is reached faster than most people expect once an app has a few dozen instances each holding their own pool.

PgBouncer sits between clients and Postgres and multiplexes many client connections onto far fewer real backend connections:

App 1one of hundreds
App 2one of hundreds
App 3one of hundreds
PgBouncer
[databases]
learning = host=postgres port=5432 dbname=learning

[pgbouncer]
pool_mode = transaction
max_client_conn = 1000
default_pool_size = 20

Pooling modes trade compatibility for reuse:

  • Session — one client connection maps to one backend for the whole session. Fully compatible (session-level state like SET and advisory locks works normally), least reuse.
  • Transaction — a backend is only held for the duration of one transaction, then returned to the pool. Much better reuse, but session state (prepared statements outside the pooler's support, SET outside a transaction) doesn't survive across transactions.
  • Statement — a backend is returned after every single statement. Best reuse, breaks multi-statement transactions entirely; rarely used.

Sharding — the last resort

Sometimes writes themselves need to scale past what one primary can do. Sharding splits data across multiple independent Postgres instances by some key (e.g. customer_id), so each shard only holds and serves a slice of the data. Postgres has no built-in sharding; it's done either at the application layer (your code decides which shard to query) or with an extension like Citus, which distributes tables and rewrites queries across a cluster of Postgres nodes transparently.

Query Router
Shard 1customer_id % 3 = 0
Shard 2customer_id % 3 = 1
Shard 3customer_id % 3 = 2

Sharding is powerful and also the most operationally expensive option on this list — cross-shard joins, transactions, and rebalancing all become hard problems. Reach for it only after vertical scaling, read replicas, and pooling are genuinely exhausted, not as a first move.

High availability

None of the above automatically handles the primary failing. Tools like Patroni watch a primary's health, hold leader election via a distributed store (etcd/Consul/ZooKeeper), and promote a replica to primary automatically on failure — the same "one primary, several standbys, automatic failover" shape as most other database HA systems.

This repo's postgres/compose.yaml runs a single node for learning purposes; none of these tools are wired up locally, but the vertical → replicas → pooling → sharding progression is the order production systems actually climb it in. Next: Users, Roles & Security.

On this page