Database Optimization Tricks in Postgres
Published 29 September 2026
A running collection of PostgreSQL optimization tricks, one section per trick.
Covering index¶
A covering index is an index that contains all the columns a query needs, so PostgreSQL can answer the query from the index alone without looking up the full table row.
For example, if you often run:
An index on email can find the matching row, but PostgreSQL may still need to visit the
table to fetch name.
You could make the index cover both columns:
Now email helps PostgreSQL find the row, and name is stored in the index for this query.
PostgreSQL may then use an index-only scan. It still sometimes checks the table for visibility information, so a covering index doesn't guarantee that every query avoids table access.
Use covering indexes for queries you run often and keep them focused: extra index columns use disk space and make inserts and updates a bit more expensive.
Keyset pagination¶
Offset-based pagination gets slower the deeper you go:
PostgreSQL generally has to find and step past the first 100,000 rows before returning 20, so every later page does more wasted work. Keyset pagination avoids the skip entirely: instead of saying "skip N rows", it asks for rows after the last sort key from the previous page.
SELECT id, created_at, total
FROM orders
WHERE (created_at, id) < ('2026-09-01 10:00:00', 12345)
ORDER BY created_at DESC, id DESC
LIMIT 20;
With a matching index, PostgreSQL jumps straight to that position and reads the next 20 rows, so page 5,000 costs about the same as page 1:
The id tie-breaker matters when multiple rows share the same timestamp — without it, rows
can be skipped or repeated across pages. The trade-off is that you can't jump to an
arbitrary page number; you can only move forward (or backward) from a known row.
Partitioning¶
Scenario: the table grows to 500 million rows, and 90% of queries only touch the last 30
days of data. Someone suggests partitioning. Would you? On which column, and what kind of
partitioning? And what happens to a query like
SELECT * FROM users WHERE email = 'x@y.com' after you partition?
Decision¶
I'd consider partitioning if the table is large enough that managing and querying it is becoming difficult. The fact that 90% of queries touch only the last 30 days is a strong reason to consider it, but partitioning won't automatically make every query faster. I'd check actual query plans and maintenance needs first.
But the first move should be to question the premise. This is a users table. Users don't expire, so "retention" would mean deleting customers. And do 90% of queries really only touch users who signed up in the last 30 days? That pattern fits an events, orders or logs table much better. A stronger answer:
Time-based partitioning makes sense if this is really an append-heavy table like events or orders. For a users table, I'd push back: lookups are by
user_idoruser_idif the table size itself becomes a problem.
The rest of this section assumes the time-oriented case.
How¶
Partition by the date or timestamp column used to restrict queries — for example,
created_at — using range partitioning, with one partition per month or week. A query
such as:
can then use partition pruning to avoid scanning older partitions.
What "range" means here¶
With range partitioning, each partition owns a non-overlapping interval of created_at
values:
CREATE TABLE users (
id bigint,
email text,
created_at timestamptz
) PARTITION BY RANGE (created_at);
CREATE TABLE users_2026_09 PARTITION OF users
FOR VALUES FROM ('2026-09-01') TO ('2026-10-01');
users_2026_09 holds rows where created_at is at least September 1 and less than
October 1 (the upper bound is exclusive).
When a query filters on created_at, PostgreSQL compares that filter with the known
partition boundaries and skips partitions that cannot contain matching rows. It doesn't
"extract the month" from the timestamp — it reasons about the actual condition against the
bounds. A one-month condition can leave just one partition to scan; a rolling 30-day window
usually overlaps two monthly partitions, so both get scanned.
Types of partitioning¶
PostgreSQL supports three kinds:
- Range — each partition holds an interval of values (dates, IDs). The natural fit when
queries mostly filter by recent dates, and it makes retention easy: dropping an old month
is
DROP TABLE, not a hugeDELETE. - List — each partition holds an explicit set of values, such as one per country or status.
- Hash — PostgreSQL hashes a key and spreads rows across a fixed number of partitions. This evens out data and work, but a date-range query can't skip any partitions.
A real-world hash example¶
Take a high-volume events table where most queries look up one account's events:
If there's no natural date-based retention requirement, hash-partition by account_id
into, say, 16 partitions:
CREATE TABLE events (
account_id bigint,
event_id bigint,
payload jsonb
) PARTITION BY HASH (account_id);
CREATE TABLE events_p0 PARTITION OF events
FOR VALUES WITH (MODULUS 16, REMAINDER 0);
CREATE TABLE events_p1 PARTITION OF events
FOR VALUES WITH (MODULUS 16, REMAINDER 1);
-- Repeat for remainders 2 through 15
PostgreSQL hashes account_id to find its remainder, so a lookup for account_id = 42
prunes to a single partition. Rows are spread evenly, which keeps each table and its
indexes smaller and allows partition-level maintenance.
Hash partitioning fits when you want an even spread by a key that queries commonly specify,
but there's no useful range or list grouping. If queries instead ask for recent events
across all accounts, hashing by account_id won't let PostgreSQL skip anything for a date
range.
What you actually gain¶
Pruning is only one of the wins, and often not the biggest:
- Cheap retention. Dropping (or detaching) an old partition is a near-instant metadata
operation. The unpartitioned alternative is a huge
DELETEthat bloats the table and leavesVACUUMa mess to clean up. - Smaller indexes. Each partition has its own, smaller indexes. The ones for hot, recent partitions are more likely to fit entirely in memory.
- Per-partition maintenance.
VACUUM,ANALYZEand reindexing run on one partition at a time, so each run is faster, and old partitions that no longer change rarely need it.
Trade-offs¶
You'll need to manage partition creation and retention, and queries that don't filter on the partition key may still need to consider many partitions.
The email lookup gets worse, not just "slightly slower". Take this query:
There's no created_at condition, so PostgreSQL can't prune anything. Say there are 40
monthly partitions. Postgres has no global index, so each partition gets its own index
on email, and the query probes all 40 of them — 40 index lookups instead of 1. That's fine
occasionally, but on a hot path like login it adds up.
The bigger trap: unique constraints must include the partition key. In PostgreSQL, every
UNIQUE (and PRIMARY KEY) constraint on a partitioned table has to contain the
partitioning column. Once users is partitioned by created_at, you can't enforce
UNIQUE (email) across the table — only UNIQUE (email, created_at), which allows the same
email in two different months. For a users table, that's a serious problem: email
uniqueness now has to be enforced somewhere else (a separate lookup table, or the
application), and both options are weaker than a real constraint.
The same rule hits the primary key: it has to become PRIMARY KEY (id, created_at)
rather than just id. That ripples out to foreign keys — any table referencing users
must now reference the full (id, created_at) key, so every child table has to carry
created_at as well.
So partitioning users by created_at is likely the wrong call. It fits better on
append-heavy, time-oriented tables — events, logs, orders — where queries filter on the
date and nothing needs to be unique across all of time.
Partitioning vs sharding¶
These two get used interchangeably, but they solve different problems.
Scenario: the app is now global with 2 billion users, and a single Postgres instance can't
handle the write load anymore, even with partitioning. You need to shard. What's the actual
difference, would you shard by user_id or by region, and what happens when one shard gets
too hot or you need more shards?
The difference¶
Partitioning breaks one big table into smaller tables on the same host. Indexes get
smaller and maintenance gets cheaper, but the total rows, CPU, memory and disk all still
belong to one machine. It's also transparent to the application: you still query users,
and Postgres routes to the right partition.
Sharding splits the data across multiple physical hosts, each its own database. That's what actually adds write capacity. It's usually not transparent: routing logic lives in the application or in a proxy/extension (Citus, Vitess) that knows which shard owns which key.
| Partitioning | Sharding | |
|---|---|---|
| Where the data lives | One host | Many hosts |
| Fixes | Big-table pain: index size, vacuum, retention | Single-host limits: write load, total size |
| Application sees | One table | Usually needs routing logic |
| Unique constraints | Must include partition key | Only enforceable per shard |
| Transactions | Normal | Cross-shard ones are expensive |
Shard by user_id or by region?¶
By user_id, as the default. Sharding by region suffers from skew: if 90% of users are
in the US, the US shard is nearly the whole dataset and you've gained little.
There's one legitimate reason to shard by region: data residency. GDPR may require EU
users' data to stay in the EU. The real-world answer is often region first, then hash by
user_id within each region.
Within user_id, use hash sharding rather than range. user_id is usually a numeric,
auto-incrementing value, so with range sharding every new user gets the highest ID and
every insert lands on the last shard. Hashing spreads new users evenly.
What sharding costs you¶
- Scatter-gather queries. Any query without a
user_idfilter has to be broadcast to every shard and the results merged. - No global unique constraints. Uniqueness is only enforced per shard, so constraints across shards need other tricks.
- Login by email needs a lookup table. After sharding by
user_id, a lookup by email would be a scatter-gather. The fix is a separateemail → user_idtable, itself sharded byemail, which routes the query to the right user shard. This is a global secondary index built by hand — and it's also where global email uniqueness gets enforced. - Cross-shard transactions are possible but expensive. Two-phase commit, or systems like Spanner, CockroachDB, Citus and Vitess, can do them, but they're slow, complex and hurt availability. Design the schema so a user's data lives on one shard and most transactions stay single-shard.
- Moving data between shards is hard and takes real engineering time.
Adding shards¶
With naive hash(user_id) % N, going from 4 to 5 shards remaps about 80% of keys: a key
stays put only if hash % 4 == hash % 5, which holds for 4 out of every 20 hash values. That
means a massive data migration. Two standard fixes:
- Consistent hashing. Keys and nodes sit on a hash ring; adding a node moves only about 1/N of the keys.
- Many logical shards. Create, say, 4,096 virtual shards up front and map them onto 8 physical hosts. To scale out, move whole virtual shards to new hosts — the key-to-virtual-shard mapping never changes, only the virtual-to-physical map does. Vitess, Citus and Instagram's setup work roughly like this.
When one shard gets too hot¶
First diagnose why. On a users table the usual cause is uneven key distribution, or several heavy tenants landing on the same shard — not a single row being hammered. With logical shards the fix is the same as scaling out: move some of the hot host's virtual shards somewhere else.
Key salting solves a different problem: one key getting hammered, like a celebrity's
likes counter or a viral post. You split the key into post-1, post-2, … spread across
shards, and reads have to gather and combine all the salts. It works, but it adds separate
read logic for hot keys, and a single user row is rarely that hot.
Replication¶
Scenario: your shards handle writes fine, but reads are 20x the writes, so you've added read replicas. A user updates their profile picture, refreshes, and sees the old picture. They refresh again and see the new one. What's happening? How do you fix it without sending all reads to the primary? Should replication be sync or async, and what do you give up with each?
What's happening¶
Replication lag. The write went to the primary; the first refresh was served by a replica that hadn't applied it yet; the second refresh hit a replica (or the same one, later) that had caught up. The guarantee being violated is called read-your-writes consistency.
Vocabulary that's easy to get wrong¶
Most of the mistakes I made on this question came from blending terms that belong to different systems.
Leader-based vs leaderless replication. These are different architectures, and the reasoning for one doesn't carry over to the other.
| Leader-based (single-leader) | Leaderless (Dynamo-style) | |
|---|---|---|
| Who accepts writes | One primary (leader) | Any node |
| How writes spread | Leader ships its log to followers | Client/coordinator writes to W of N nodes |
| Read guarantee comes from | Reading the leader, or a replica known to be caught up | Quorum overlap: R + W > N |
| Examples | Postgres, MySQL, Raft-based systems | Cassandra, DynamoDB, Riak |
- Quorums and the pigeonhole argument (
R + W > Nmeans every read set overlaps every write set) belong to leaderless systems. They're correct there. - Postgres and MySQL are single-leader. With async replication the primary waits for nobody — not a majority, not one replica. So reading from 3 of 4 replicas guarantees nothing: all 4 could be behind.
- Raft is leader-based too, even though it uses majorities. The leader commits an entry once a majority has it in their logs, but reads still go through the leader (via a lease or a ReadIndex check) to be linearizable — you don't get fresh reads by polling a majority of followers. Postgres streaming replication isn't Raft either way.
Other terms worth having straight:
- WAL (write-ahead log): the log Postgres writes changes to and ships to replicas.
- LSN (log sequence number): a position in the WAL.
pg_current_wal_lsn()on the primary,pg_last_wal_replay_lsn()on a replica for how far it has applied. - Written vs flushed vs applied: a replica can have received the WAL, flushed it to disk, and still not have replayed it — so a query on that replica won't see it yet.
- Failover: promoting a replica to primary when the primary dies.
Fixing read-your-writes without routing everything to the primary¶
From simplest to most involved:
- Show the change on the client. The frontend already has the new picture, so it can display it immediately. The problem disappears without touching the database.
- Read from the primary for a short window after a write. Route that user's reads to the primary for ~5–10 seconds, tracked in their session or a cookie. Simple and very common.
- LSN-based routing. The write returns its LSN (
pg_current_wal_lsn()after commit), stored in the user's session. On the next read, pick any replica and check that itspg_last_wal_replay_lsn()is at least the stored LSN. If not, wait briefly or fall back to the primary. No broadcast and no extra primary query needed.
It's still worth asking whether the complexity is worth it — for most reads, brief staleness is fine.
Sync vs async¶
| Async | Sync | Semi-sync / quorum commit | |
|---|---|---|---|
| Primary waits for | Nobody | Every replica (rare in practice) | A configured number k of standbys |
| Write latency | Lowest | Highest | Middle |
| On primary crash | Can lose recent committed writes | No loss | No loss if an acked replica survives |
| If a replica dies | Nothing happens | Writes block | Fine while k replicas are up |
- The real cost of async is data loss on failover, not just stale reads. If the primary dies before shipping its latest writes, those writes are gone — even though the user was told they succeeded.
- The real cost of sync is availability: one slow or dead replica stalls every write.
Sync replication doesn't give you read-your-writes by default. Postgres's
synchronous_commit = on only guarantees the replica has written the WAL to disk, not
applied it. You need remote_apply for replica reads to see the write.
How many replicas does "sync" actually wait for?¶
Usually only a subset, and you configure how many. Waiting for every replica exists in theory, but almost nobody runs it, because one slow or dead replica would stall every write.
The typical production setup: one primary, k synchronous standbys (usually k = 1) in a different availability zone, and every other replica async. At least two nodes always have every committed write, you only pay for one extra round trip per commit, and the standby survives a zone outage. Designing Data-Intensive Applications calls this semi-synchronous replication.
Postgres¶
synchronous_standby_names chooses which standbys count and how many to wait for:
# Wait for the first 1 available standby, in priority order
synchronous_standby_names = 'FIRST 1 (replica_a, replica_b, replica_c)'
# Quorum: wait for ANY 2 of these 3 (whichever ack first)
synchronous_standby_names = 'ANY 2 (replica_a, replica_b, replica_c)'
FIRST kis priority-based: the primary waits for the k highest-priority standbys that are connected. If one drops, the next in the list takes its place.ANY kis quorum-based: the primary waits for any k of the listed standbys.- Standbys that aren't listed are async.
synchronous_commit then decides what an ack means:
| Value | The replica has… | Can a read on the replica see the write? |
|---|---|---|
remote_write |
received the WAL and handed it to the OS (not yet flushed) | Not guaranteed |
on (default) |
flushed the WAL to disk | Not guaranteed |
remote_apply |
replayed the WAL | Yes |
If not enough sync standbys are available, Postgres blocks commits until they come back. It chooses durability over availability.
MySQL semi-sync¶
rpl_semi_sync_source_wait_for_replica_countdefaults to 1. The primary waits for one replica to acknowledge it received the event, not that it applied it.- The catch: if no replica acks within
rpl_semi_sync_source_timeout(10 seconds by default), MySQL quietly falls back to async. A crash right after that fallback can still lose data — a good detail to bring up in an interview.
Consensus-based databases¶
CockroachDB, Spanner, etcd and TiDB work differently. There, "synchronous" means waiting for a majority (Raft or Paxos), so a minority of slow nodes never blocks writes.
A strong answer:
Async for replicas that serve reads, because the latency matters and brief staleness is fine for most reads. I'd handle read-your-writes at the application layer. For durability, I'd keep one semi-sync standby in another availability zone, so a failover doesn't lose acknowledged writes, without every write waiting on every replica.
Disaster recovery¶
What happens when a node in a sync setup fails, and what's the recovery workflow? Take this config:
replica_a is the sync standby; replica_b and replica_c are async (Postgres labels
them potential sync standbys). There are three separate failure cases, and they play out
very differently.
RPO and RTO¶
Two terms frame every DR discussion:
- RPO — Recovery Point Objective. The most data you can afford to lose, measured in time. "RPO = 5 minutes" means losing up to the last 5 minutes of writes is acceptable. It looks backward from the moment of failure.
- RTO — Recovery Time Objective. The longest you can afford to be down. "RTO = 30 minutes" means service has to be back within 30 minutes. It looks forward from the moment of failure.
Case 1: the sync standby fails¶
Postgres handles this automatically:
- The primary loses its connection to
replica_a. - It picks the next connected standby in the priority list,
replica_b, and makes it the sync standby. No config change or reload is needed. - Commits that were waiting on
replica_aare released oncereplica_bconfirms it has flushed WAL up to their LSN.
You can watch this in pg_stat_replication.sync_state:
Things that can go wrong:
- The next standby is lagging. If
replica_bis in statecatchuprather thanstreaming, commits stall until it catches up. Failover is automatic, but not always instant. - All standbys are down. Commits hang indefinitely — they don't error, clients just wait. Postgres chooses durability over availability.
- The escape hatch. Set
synchronous_standby_names = ''and runSELECT pg_reload_conf();. Writes resume, but you're now async, and a primary crash can lose data. This is a deliberate business decision, not something to automate. - "Cancelled" doesn't mean rolled back. If a client cancels a commit that's waiting on
replication, Postgres returns
WARNING: canceled wait for synchronous replication ... transaction has already committed locally. The transaction is committed on the primary but may not be replicated.
Bringing replica_a back:
- Restart it. If the primary still has the WAL it missed (kept by a replication slot or
wal_keep_size), it streams and catches up. - If that WAL is gone, rebuild it with
pg_basebackup. - Once caught up, it becomes sync again automatically, since it's first in the list.
- Watch out: a replication slot for a dead replica retains WAL forever and can fill the
primary's disk. Set
max_slot_wal_keep_sizeas a safety limit.
Case 2: the primary fails¶
This is the real DR scenario. Postgres has no built-in automatic failover. You need tooling: Patroni (most common, backed by etcd or Consul), repmgr or pg_auto_failover — or a managed service like RDS Multi-AZ.
The workflow:
- Detect. The primary's leader lease (in Patroni, a key in etcd) expires after missed heartbeats. Tune the timeout to avoid false positives from network blips.
- Fence the old primary. Make sure it can't accept writes anymore — STONITH, cutting its network access, or Patroni demoting it when it loses the lease. Skip this and you risk split brain: two primaries both accepting writes.
- Choose the right standby. Promote the one that was sync at the time of the crash.
It's the only node guaranteed to have every acknowledged commit, which gives RPO = 0.
Promoting
replica_c(async) would silently lose data. This is why Patroni'ssynchronous_moderecords the current sync standby in etcd — so it knows exactly which node is safe to promote. - Promote it with
pg_promote()orpg_ctl promote. The new primary starts a new timeline. - Repoint the other replicas at the new primary. They follow the timeline switch
(
recovery_target_timeline = 'latest'). - Update
synchronous_standby_nameson the new primary to list the remaining standbys, e.g.FIRST 1 (replica_c, old_primary). Patroni does this automatically. - Redirect clients by moving a virtual IP, updating DNS, or using HAProxy/PgBouncer with health checks against Patroni's REST API.
- Rejoin the old primary as a standby. It may have written WAL that never reached the
sync standby (transactions that never got an ack), so its history has diverged.
pg_rewindrolls it back to where the timelines forked (it needswal_log_hints = onor data checksums). If that fails, rebuild withpg_basebackup.
In-flight transactions. A client whose commit never got an ack before the crash didn't get a success, so losing that write is acceptable. But the client can't tell whether it committed — which is why idempotency keys matter: the retry has to be safe.
Case 3: the whole region goes down¶
The sync standby is often in the same region as the primary (for latency), so a region failure takes out both.
- A cross-region async replica gives fast recovery, but RPO > 0: you lose whatever was lagging, typically seconds.
- Continuous WAL archiving to S3 (pgBackRest or WAL-G) plus base backups gives
point-in-time recovery (PITR). It's slower to restore — hours for a large database — but
it's the last line of defense. It also protects against a
DROP TABLE, which replication would faithfully copy to every standby.
Picking a setup¶
| Setup | RPO | RTO |
|---|---|---|
| Sync standby + Patroni failover | 0 (no acknowledged writes lost) | ~30 seconds to a few minutes |
| Async replica failover | A few seconds (whatever was lagging) | ~1 minute |
| Restore from S3 backups + WAL (PITR) | Seconds to minutes (depends on archive frequency) | Hours for a large database |
Lower RPO and RTO cost more money and add write latency, so the business sets the targets and engineering picks the setup that meets them.
The interview one-liner¶
If the sync standby dies, Postgres promotes the next one in the
FIRSTlist automatically. If all standbys die, commits block, and dropping to async is a business decision. If the primary dies, Postgres doesn't fail over on its own. Patroni detects it, fences the old primary to prevent split brain, and promotes the sync standby specifically, since that's what gives RPO = 0. Then it repoints the other replicas and the clients, andpg_rewinds the old primary back in as a standby. For a region-level disaster, you rely on a cross-region async replica or WAL archives in S3, and accept some RPO.
When should you shard?¶
Scenario: back to the start. The table is 1.2 million rows, one Postgres instance, and the query is slow. A junior engineer says: "We should shard it now so we don't have problems later." How do you respond? And at what point would you actually consider partitioning, replicas or sharding? Concrete signals, not just "when it gets big".
Diagnose before you architect¶
1.2 million rows is tiny — it likely fits entirely in RAM on a laptop. A slow query at that size means a missing index or a bad query, not a scaling problem:
Sharding won't fix a missing index. It'll give us the same slow query on 8 machines, plus cross-shard problems. Let's run
EXPLAIN ANALYZEfirst. At 1.2M rows, this is a tuning problem, not a scaling problem.
The ladder¶
Sharding goes last because its engineering cost is very high. Climb the cheaper rungs first, and don't skip them unless you're genuinely sure the load will outgrow them within months:
- Tune queries and indexes.
EXPLAIN ANALYZE, add or fix indexes, rewrite bad queries. - Connection pooling (PgBouncer).
- Caching (Redis) for hot, read-heavy data.
- Read replicas — for reads only. Replicas don't help a write bottleneck at all, and writes are usually the harder problem.
- Vertical scaling. Buy a bigger machine; large instances go a long way, and most companies never need to shard.
- Shrink the hot table. Move rarely used or static columns into a separate table so the hot rows are narrower.
- Partitioning + archiving — for time-series or archival data. It isn't free (see above: it changes your primary key, breaks global unique constraints and complicates foreign keys), so don't do it "early just because".
- Sharding.
Signals that you actually need to shard¶
Row count alone is a weak signal — a 1B-row table with good indexes can be perfectly fine. Look at load and operational signals instead:
| Signal | What it tells you |
|---|---|
| Primary CPU consistently above ~70% after tuning | You've run out of compute on one box |
| Buffer cache hit ratio falling below ~99% | The hot data no longer fits in RAM |
| Write IOPS or WAL throughput near the instance's limit | The write ceiling that replicas can't fix |
| Replication lag growing under normal load | Replicas can't keep up with the write volume |
VACUUM can't keep up, bloat keeps growing |
The table is too large to maintain |
| Backup or restore takes longer than your RTO | Often the real forcing function |
| Already on the largest instance type | No vertical headroom left |
The backup/restore signal is the catch in "just get a host with a 16 TB disk". A 16 TB database is possible, but restoring it after a failure can take many hours — sometimes a day.
The NoSQL off-ramp¶
If your access patterns fit, DynamoDB, Cassandra or MongoDB can scale elastically. But they get there by dropping joins, ad hoc queries and often strong consistency. That only works when you know your access patterns up front.