Database Performance Tuning and Scaling Questions
System-level performance work beyond a single query: configuration and resource tuning, capacity planning, handling large data volumes, and scaling read and write throughput. Covers identifying bottlenecks, growth management, the vertical-versus-horizontal scaling decision, materialized views for expensive queries, and maintenance work such as vacuuming, index rebuilds, bulk loads, and safe schema changes on large tables. Tests whether a candidate can keep a database healthy as load grows.
A read-heavy service is seeing increased database read latency during traffic spikes. Walk through how you would reduce it: what would you look at first, what would you try next, and how would you decide between the different levers available to you, spanning data-access changes, caching, and precomputation? For each option you pick, explain the trade-offs and how you would measure whether it worked.
Sample Answer
Direct answer
Work from the cheapest, most targeted fix toward the most expensive, most structural one: first confirm what is actually slow, a specific query, a specific table, or genuinely just more concurrent load than any single fix will absorb, then reach for data-access changes before caching before precomputation, because each step down that list costs more to build and adds more that can go wrong.
How to think about it
Confirm the shape of the spike first. Is latency up because query count rose (more concurrent requests, same per-query cost) or because per-query cost itself rose (a plan regression, a lock, a missing index that only bites at higher row counts)? EXPLAIN ANALYZE (runs the query for real, instead of only estimating, and reports the actual plan, timing, and row counts) on the actual slow query, and wait events in pg_stat_activity, answer this in minutes and determine which lever below is even relevant.
Data-access changes first, usually the cheapest and lowest-risk: is there a missing or unused index on the filtered or joined columns; is the query doing more work than the request needs, selecting unused columns, or issuing many small queries where one would do; is the connection pool itself the bottleneck (a connection pool is a fixed-size set of already-open database connections that requests borrow and return, instead of opening a brand new connection every time), requests queueing for a pool slot look identical to "the database is slow" from the outside, but the fix is pool sizing, not the query.
Caching next, once data-access tuning is exhausted and the same read repeats often enough to be worth it: an application-side cache in front of the hot read, sized and invalidated for that specific access pattern, trades a new consistency and invalidation surface for latency.
Precomputation last, for read shapes that are expensive no matter how well-indexed the base tables are, a cross-row aggregate, a multi-table join for a dashboard: materialize the result ahead of time so the spike-time read becomes a plain table scan, not a live computation.
For each lever, measure before and after using the same metric the spike broke, typically the 95th or 99th percentile (p95 or p99, the latency value slower than 95% or 99% of requests) latency on the affected endpoint, not a proxy metric, and confirm the fix actually moved that number under comparable load, not just in isolation.
Worked example
A real, pinned demonstration of step two done correctly, executed against a live PostgreSQL 16 instance: a 50,000-row orders table (id, customer_id, order_date, amount, status) filtered by customer_id, first with no supporting index.
Seq Scan on orders
Filter: (customer_id = 408)
Rows Removed by Filter: 49990
Buffers: shared hit=417
Execution Time: 1.022 ms
After adding a single index on that column:
CREATE INDEX idx_orders_customer_id ON orders (customer_id);
Bitmap Heap Scan on orders
Recheck Cond: (customer_id = 408)
Heap Blocks: exact=10
Buffers: shared hit=10 read=2
Execution Time: 0.035 ms
About a 35x reduction in buffer touches (417 down to 12) and a 29x reduction in execution time (1.022 ms down to 0.035 ms) on this demo's data. Those two ratios do not have to match, and here they do not: buffer touches count pages visited, execution time is wall-clock cost, and a cold-cache random-access fetch still costs more per buffer than a page the sequential scan was already streaming through. What matters for the argument is that both numbers dropped sharply once the index let the planner skip almost the entire table. This is exactly what a well-targeted data-access change looks like: cheap, no new infrastructure, and it moved the real measured numbers.
Trade-offs and pitfalls
Reaching for caching or a materialized view (a query's result saved as its own table and refreshed on a schedule, instead of recomputed on every read) before confirming the query is not just missing an index over-engineers an expensive fix for a cheap problem. Adding an index reflexively without checking its write-side cost is its own mistake, every index slows down inserts, updates, and deletes on that table and adds storage, so index the columns real queries actually filter or join on, not every column that might someday be useful. Treating "add a read replica" as a data-access fix is a category error, it is a capacity fix, more read throughput across more machines, and it does not help if the problem is a single query's per-execution cost, that same slow query just runs slowly on more machines instead of one.
You need to add a column to a production table with hundreds of millions of rows, and you cannot take a long lock or cause a visible outage. Describe a safe approach: what technique would you use to make the change incrementally, how would a dual-write-and-backfill strategy work if you needed one, and how would you monitor and cap the impact on live traffic while it runs, including a way to back out if something goes wrong?
Sample Answer
Direct answer
Adding the column itself should be near-instant: a plain ADD COLUMN with a constant (or
NULL) default is a metadata-only change on modern PostgreSQL and doesn't rewrite the
table. The actual risk is in populating it: run that as an explicit incremental backfill
in small batches with pacing between them, have the application dual-write the new column
on every new or updated row going forward before the backfill starts, and monitor replication lag (how far behind the primary a replica's copy of the data has fallen) and lock wait as hard caps that pause the job automatically.
The plan
The ADD COLUMN step itself. On PostgreSQL 11 and later, ADD COLUMN with a
constant default, including NULL, is metadata-only and does not rewrite existing rows.
That's a genuine change from older versions, and a common misconception is that any ADD COLUMN with a default rewrites the whole table; that's only true for a volatile,
non-constant default, or certain constraint additions. This step takes a brief exclusive
lock, but for milliseconds, a fundamentally different risk than the population step.
If a full backfill is needed (computing a derived value per row, the harder case this
question is really about):
- Incremental technique: iterate in small batches by primary-key range, each batch its
own short transaction, with a brief pause between batches so replication and any read
replicas can keep up, and so row locks don't accumulate into one long-held transaction. - Dual-write: ship the application change that writes the new column on every new or
updated row first, before the batch backfill starts, so the backfill isn't chasing a
moving target. The batch job then only needs to cover rows written before that deploy. - Monitoring and capping impact: watch replication lag (don't let a replica fall far
enough behind that a failover would lose data, or that it starts serving badly stale
reads); watch lock-wait time and set alock_timeoutso a batch that's queued too long
aborts rather than backs up behind unrelated traffic; watch dead-tuple growth (a dead tuple is the old copy of a row that Postgres leaves behind after anUPDATE, since it writes a new row version instead of editing the old one in place); the backfill itself generates a dead row version for every update, a real bloat event, meaning the table fills up with dead tuples faster than they get cleaned out, at this row count that needs its own vacuum planning (VACUUMis the background process that reclaims the space those dead tuples hold so the table doesn't just keep growing); and cap the overall rate (target
rows/sec) rather than running batches back to back as fast as possible. - Backing out: since the column can be added nullable with nothing yet depending on it,
backing out mid-backfill just means stopping the job. The column can sit partially
populated with no correctness impact as long as nothing treats it as authoritative
until the backfill and validation are complete, and the dual-write path can be
feature-flagged off independently of the backfill's progress.
Worked example: backfill timing (pinned inputs)
total_rows = 400_000_000
batch_size_rows = 5_000
sleep_between_batches_ms = 50
exec_time_per_batch_ms = 20 # illustrative per-batch UPDATE execution time
num_batches = total_rows / batch_size_rows
time_per_batch_s = (sleep_between_batches_ms + exec_time_per_batch_ms) / 1000.0
total_time_s = num_batches * time_per_batch_s
print(f"num_batches = {num_batches:,.0f}")
print(f"time_per_batch = {time_per_batch_s*1000:.0f}ms")
print(f"total wall time = {total_time_s:,.0f}s = {total_time_s/3600:.1f} hours")
num_batches = 80,000
time_per_batch = 70ms
total wall time = 5,600s = 1.6 hours
Halving the batch size to 2,500 rows roughly doubles num_batches to 160,000 and, for a
similar per-batch overhead, roughly doubles the wall time to about 3.2 hours, a direct
trade between how gentle each batch is on the live system and how long the whole backfill
takes. That trade is worth sizing explicitly against how much time the migration is
actually allowed to take, not defaulted to "as small as possible."
Trade-offs and pitfalls
- Assuming
ADD COLUMNitself is the risky step; on modern PostgreSQL with a constant
default it usually isn't. The real risk in this scenario is almost always the
subsequent full-table update to populate a derived value; don't spend the caution
budget on the step that doesn't need it. - Smaller batches and longer pauses are safer for the shared system but make the backfill
take proportionally longer, size this trade-off deliberately, not by habit. - Shipping the backfill before the dual-write path is verified working means the backfill is chasing a target that's still drifting from concurrent, uninstrumented writes: rows the backfill already processed keep getting changed by callers the dual-write path never reached, so the column silently falls out of sync with its source again after the batch job has already moved past that row. This is the same kind of drift failure that hits any derived copy of data, a cache, a summary table, a replica, whenever its update logic misses a write instead of catching every one.
- Not planning for the backfill's own bloat: 400 million updates generate 400 million
dead tuples somewhere, a genuinely large vacuum workload competing with the same system
the whole plan is trying not to disrupt.
A multi-tenant SaaS product runs many tenants on a shared database cluster. One tenant's workload suddenly causes I/O and CPU spikes that degrade performance for everyone else. Propose an architecture and operational plan to detect and isolate noisy tenants: logical isolation (schemas, row-level throttles), physical isolation (dedicated instances for the worst offenders), resource limits, and how you would think about cost and billing allocation.
Sample Answer
Direct answer
Layer the defense: detect resource usage per tenant, contain with logical throttles
first because they're cheapest and fastest to roll back, escalate to physical isolation
only for genuinely disproportionate tenants, and back the whole thing with per-tenant-
class SLA (service-level agreement, a stated performance guarantee) commitments and a
cost model that makes "move to dedicated capacity" a real lever rather than a punishment.
Architecture
flowchart TD
T[Tenant request] --> TAG[Tag query with tenant_id]
TAG --> MON[Per tenant resource monitor]
BATCH[Cross tenant reporting job] --> MON
MON --> CHECK{Over fair use threshold repeatedly?}
CHECK -->|No| SHARED[Shared cluster: standard tier]
CHECK -->|Yes| ESC[Offer dedicated tier]
ESC --> DEDICATED[Dedicated instance: premium tier]
SHARED --> THROTTLE[Per tenant rate limit and statement timeout]
DEDICATED --> SLA[Per tenant class SLA guarantee]
Detection. Attribute resource usage to a tenant, either by tagging every connection
or query with a tenant identifier (Postgres will show a tagged application_name in
pg_stat_activity, so usage can be joined back to a tenant) or by giving each tenant its
own connection pool (a shared, reusable set of already-open database connections that requests
borrow from and return instead of opening a new one each time) or role, so OS-level accounting can
separate them. A real, common
trigger scenario doesn't look like one loud tenant: it's a cross-tenant reporting job,
something like a nightly "generate this week's report for every account" batch that
scans and joins across many tenants in a single run. That job belongs to no single
tenant, so a naive per-request tenant tag won't flag it; detection has to separately
watch for platform jobs whose blast radius is everyone, not just single-tenant spikes.
Logical isolation. Schema-per-tenant or a shared schema with a tenant_id column,
paired with per-tenant limits enforced at the application or connection-pool layer:
capped concurrent queries, a tenant-specific statement_timeout. Cheapest to build and
revert. Limitation: throttling caps concurrency, not the shared resource's total
capacity, so a tenant staying within its concurrency cap can still exhaust shared I/O
bandwidth or the shared buffer pool (the database's in-memory page cache) if its queries
are individually expensive enough.
Physical isolation. Move the worst, persistently offending tenants to a dedicated
instance or read replica (a synced, read-only copy of the database that offloads read traffic from the
primary). Fully solves the shared-resource problem for everyone else, at
a real operational cost: more instances to patch and monitor, and a routing layer that
has to send each tenant's traffic to the right place.
Resource limits. Container/OS-level cgroup limits (CPU shares, I/O weight) if
tenants run in separate processes. Database-level equivalents: per-role
statement_timeout and work_mem caps, and, with an extension, I/O throttling. A
governor pattern where a tenant that exceeds its allotment gets queued or served a
degraded (e.g. cached, slightly stale) response rather than allowed to starve others.
Cost and billing. Tie the isolation tiers to pricing. A standard tier runs on the
shared cluster with logical throttles and a best-effort SLA. A premium tier gets a
dedicated instance and a per-tenant-class SLA guarantee, a specific stated p99
latency (the response time slower than 99% of requests, a worst-case-ish bound) and uptime number
that tenant is paying for. This gives you an objective,
revenue-positive trigger for escalation: a tenant consistently exceeding the shared
tier's fair-use threshold gets offered the paid isolated tier, instead of engineering
absorbing the cost of accommodating them for free indefinitely.
Worked example: sizing the fair-use threshold
fair_share=tenant_countcluster_IOPS=20020,000=100 IOPS/tenant
(IOPS: input/output operations per second, a measure of storage throughput.) A noisy
tenant observed consuming 14,000 IOPS is consuming 14000/20000 = 70% of the entire
cluster's capacity, 140 times its fair share, degrading the other 199 tenants who are
collectively left fighting over the remaining 30%. That gap, 140x fair share versus a
threshold you might set at, say, 5 to 10x, is what a detection rule alerts on; it's large
enough that a single noisy tenant is unambiguously the cause, not normal variance.
Trade-offs and pitfalls
- Logical throttles are fast to ship but don't protect against a tenant that stays within
its concurrency cap while running individually expensive queries; watch actual resource
consumption, not just request counts. - Physical isolation is the strongest guarantee and the most expensive to operate;
reserve it for tenants whose usage justifies it, not as a default. - A billing model that treats isolation purely as a cost center misses the upsell
opportunity; framing the premium tier's SLA as something tenants pay for turns
engineering effort into a product feature. - Detecting only per-request tenant spikes misses cross-tenant batch jobs entirely; a
reporting job that touches every tenant at once needs its own monitoring category.
You are running PostgreSQL with a 500GB dataset on a host with 200GB of total memory. How would you determine reasonable values for shared_buffers and work_mem? What would you monitor to know if you got it right (cache hit ratio, page read times, OS page cache usage), how would you avoid swapping, and how would you roll the change out safely in production?
Sample Answer
Direct answer
Since 500 GB doesn't fit in 200 GB of RAM, size shared_buffers to hold the actively hot
slice of the data, not the whole dataset, a widely used starting point is about 25% of
RAM (roughly 50 GiB here), leaving the rest to the operating system's own page cache,
which Postgres relies on as a second layer. Size work_mem conservatively, because it's
granted per sort or hash operation, not as one shared pool, and a single complex query
can request several of these at once, multiplied again by parallel workers. This matters most under
a typical OLTP workload (many short, concurrent transactions serving live application traffic, as
opposed to a few large analytical queries), since a high global default gets multiplied across every
concurrently active connection at once.
Reasoning through the two settings
shared_buffers: the 25% guideline is a documented starting point, not a law. Past
that point, growing it further competes with the OS page cache for the same physical
RAM (a page cached in both layers is partially wasted memory), while increasing
background-writer and checkpoint (the periodic point where Postgres flushes recently modified
in-memory pages to disk) overhead proportional to how much dirty data can
accumulate. It requires a full server restart to change, since it's allocated as shared
memory at startup, which shapes the rollout plan below.
work_mem: budgeted per sort, hash join, or merge operation. A single query with two
sorts and a hash join can request three allotments simultaneously; parallel workers each
get their own. The real ceiling isn't total_RAM / connection_count, it's headroom
divided by realistic peak concurrent operations. Setting it too high is a classic
out-of-memory or swap risk under concurrent load, not a free win for one query's sort.
Unlike shared_buffers, it can be changed per session or per role without a restart.
Rollout plan: because work_mem is a session/role-level setting, trial a higher value
on just the reporting role first, watch memory behavior under real concurrency, then
promote it to the default once validated. Because shared_buffers needs a restart, plan
a maintenance window, keep the previous value on hand, and know reverting also needs a
restart, it isn't a hot config reload either direction.
What to monitor: the buffer cache hit ratio (below), whether sorts are spilling to
disk (the query plan will say so explicitly, see below), OS-level free memory and swap
activity (a host that starts swapping is making everything slower, not just one
under-provisioned query), and page read/write timing if track_io_timing is enabled.
Avoiding swapping: never let shared_buffers + (work_mem × realistic peak concurrent operations) + other server overhead approach physical RAM. A work_mem value that's
safe at 5 concurrent users can be the thing that swaps the box at 50; set a deliberate
ceiling per role rather than raising it until one report "feels fast."
Worked example (executed on a real PostgreSQL 16 instance)
Defaults on a fresh container (small demo values, shown only to illustrate the commands,
not as a sizing recommendation for a 200 GB host):
SHOW shared_buffers; -- 128MB
SHOW work_mem; -- 4MB
SHOW effective_cache_size; -- 4GB
The concrete signal a query plan gives you when work_mem is too low, run against a
500,000-row, 87 MB table, ordering by a text column:
SET work_mem = '64kB';
EXPLAIN (ANALYZE, BUFFERS) SELECT * FROM bench ORDER BY payload;
Sort Method: external merge Disk: 26008kB
Worker 0: Sort Method: external merge Disk: 26448kB
Worker 1: Sort Method: external merge Disk: 27024kB
SET work_mem = '64MB';
EXPLAIN (ANALYZE, BUFFERS) SELECT * FROM bench ORDER BY payload;
Sort Method: quicksort Memory: 36981kB
Worker 0: Sort Method: quicksort Memory: 32689kB
Worker 1: Sort Method: quicksort Memory: 34701kB
"External merge" with a Disk: size means the sort didn't fit in its work_mem
allotment and spilled to disk, extra I/O the query wouldn't otherwise pay. "Quicksort"
with a Memory: size means it stayed fully in memory. That line in a plan is a direct,
observable signal that a hot query's work_mem is undersized; it's what you'd grep for
in a slow-query log, not just a number you infer from theory.
Cache hit ratio, executed against the same instance:
SELECT sum(heap_blks_hit) AS hit, sum(heap_blks_read) AS miss,
round(100.0*sum(heap_blks_hit)/nullif(sum(heap_blks_hit)+sum(heap_blks_read),0),2) AS pct
FROM pg_statio_user_tables;
-- hit=555554 miss=0 pct=100.00
The demo database fits entirely in this small container's memory, so 100% is expected
here and isn't the target for the 500 GB scenario; on the actual host you'd track this
ratio against a target you've validated for your workload, since 500 GB genuinely
cannot be fully cached in 200 GB of RAM, that gap is exactly why this sizing decision is a
real trade-off rather than a solved problem.
Trade-offs and pitfalls
- Treating "25% of RAM" as a formula rather than a starting point: a mostly-sequential-
scan analytical workload wants a smallershared_buffersand a larger
effective_cache_size(a planner hint about assumed total cache, it allocates no memory
itself) instead. - Setting
work_memglobally high enough to satisfy the worst report risks an OOM under
real peak concurrency; a lower global default plus a per-role override for the specific
workload that needs bigger sorts is the safer shape. effective_cache_sizebeing misconfigured (too low or too high) silently misleads the
query planner's cost estimates even though it allocates no memory itself, worth being
explicit about since it's a common point of confusion withshared_buffers.shared_bufferschanges need a restart to apply and to revert; roll it out in a
maintenance window with the previous value ready, not as a live config reload.
A service has an end-to-end p99 latency target of 100 milliseconds. How would you decompose that latency budget across network, application, and database components? Which database-specific metrics would you instrument to find the database's share of the tail latency, and what would you do if the database were blowing the budget?
Sample Answer
Direct answer
Decompose the 100 ms p99 target (the latency below which 99% of requests finish; the
slowest 1 in 100 is what you're budgeting against, not the typical request) as a chain
of sequential legs: network, application processing, and database, and treat the sum of
each leg's own worst case as the budget ceiling. On the database side, instrument
percentile latency per query, wait-event breakdown, and connection-pool queue time
separately, because "the database is slow" and "the database queue is slow" call for
different fixes. If the database is blowing its share, diagnose which of those three
before touching anything.
How to decompose the budget
A useful mental model is a pipeline of legs, each contributing its own tail latency:
Tp99≈Tnet+Tapp+Tdb
This is deliberately conservative: percentiles don't sum linearly the way averages do.
The true p99 of a sum of independent legs is usually less than the sum of each leg's
own p99 (their worst cases rarely all land on the same request), but budgeting as if they
do stack gives you headroom you don't have to prove exists, which is the safer default
when you don't fully control every leg (a third-party network path, a shared database).
What to instrument on the database side specifically:
- Per-query latency percentiles, not the mean. Postgres's built-in query-stats
extension,pg_stat_statements, tracksmean_exec_time,min_exec_time,
max_exec_time, andstddev_exec_time, but has no percentile column. A 5 ms mean can
still blow a 60 ms p99 budget if 1% of executions hit a slow plan. True percentiles need
either sampled slow-query logging (auto_explainwithlog_min_duration_statement)
fed through an external percentile calculation, or an APM/tracing tool that tags
database calls as spans. - Wait-event breakdown, not just elapsed time. Postgres exposes
wait_eventand
wait_event_typeinpg_stat_activity, telling you whether a slow query is waiting on
a lock, on I/O, or on the client, versus genuinely doing CPU work. "Slow" and "blocked"
need different fixes. - Connection-pool queue wait, separated from execution time. A request can spend 40 ms
waiting for a pooled connection and 2 ms actually running the query; conflating the two
makes you tune the wrong thing. - Replication lag, if reads are served from a replica (a read replica: a continuously updated copy of the primary database that serves read-only queries to take load off the primary): a lagging replica can force a
fallback to the primary for a consistency-sensitive read, spiking that request's
database leg specifically. - Maintenance overlap: checkpoint activity (the periodic point where the database
flushes all recently modified in-memory pages to disk so crash recovery has less to
replay) and autovacuum (the background process that reclaims space left by deleted or
updated rows) both compete for I/O and can push an otherwise-fast query into the tail
during their run window.
Worked example: building the budget
| Leg | Budgeted p99 |
|---|---|
| Client to edge/load balancer (mobile-ish network) | 18 ms |
| Load balancer to application | 2 ms |
| Application processing, excluding the database call | 12 ms |
| Application to database round trip (network only) | 2 ms |
| Sum of fixed legs | 34 ms |
That leaves 100 - 34 = 66 ms for the database's own execution time. Because the
additive model above is a conservative approximation, not an exact composition, it's
reasonable to hold back a margin rather than spend the full 66 ms: budgeting the database
at roughly 55 ms (about a 10 ms cushion) leaves room for the fact that these legs aren't
perfectly independent in practice (a loaded application server slows its own database
calls too).
If the database is blowing its 55 ms share, work through the instrumentation above in
order:
- Check the wait-event distribution first. If it's I/O-bound (high
shared_blks_read
relative toshared_blks_hitinpg_statio_user_tables), check whether the working
set fits the buffer pool (the database's in-memory page cache) or the query is missing
an index. - If it's lock-bound, look for concurrent DDL, a long-running transaction holding a
lock, or contention on a specific hot row. - If the queue wait dominates rather than execution time, the fix is connection-pool
sizing or astatement_timeoutto stop one runaway query from starving the pool, not
query tuning. - If execution time genuinely dominates after ruling out the above, consider moving the
query off the primary's critical path: a read replica, a cache, or an asynchronous
pattern where the client gets an acknowledgment once the write is durable and
non-critical work happens after the response.
Trade-offs and pitfalls
- Using the database's mean latency instead of its p99 in the budget hides exactly the
failure this exercise exists to catch: a fast-on-average query that occasionally hits a
slow plan or a lock wait. - Treating additive percentile budgeting as mathematically wrong and cutting the
database's share too tight (because "the legs are independent, so it won't really
stack") removes headroom you're likely to need the first time something outside your
control gets slower. - Mistaking pool queue wait for database slowness leads to scaling the database when the
real fix is resizing the pool. - Chasing a database-side p99 win, for example a new index, can worsen a different
metric (write latency, storage), so the win has to be weighed against what it costs
elsewhere, not celebrated in isolation.
Unlock Full Question Bank
Get access to all 38 Database Performance Tuning and Scaling interview questions and detailed answers.
Sign in to ContinueJoin thousands of developers preparing for their dream job.