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.
Your OLTP database is seeing p99 write latency spikes traced to bursts of fsyncs. What mitigations would you consider at the OS, filesystem, database-engine, and application layers to reduce that tail latency without giving up durability, and what are the trade-offs of each?
Sample Answer
Direct answer
fsync is the system call that forces the write-ahead log (WAL, the durability log every
committed transaction must be flushed to before the commit is acknowledged) to durable
storage, so every mitigation here is really trading durability risk for latency, not a
free performance win. The highest-leverage, lowest-risk lever is grouping many commits
behind fewer fsync calls (group commit) before reaching for anything that weakens the
durability guarantee itself.
Mitigations, layer by layer
OS and filesystem layer. Avoid a filesystem or mount configuration that adds its own
extra journal or flush overhead on top of the database's own fsync (double-journaling).
Confirm the storage device honors fsync correctly: a virtualized or cloud disk with a
write cache that isn't battery- or power-loss-protected can return from a flush call
before data is actually durable, which is a correctness risk wearing a performance-win
disguise, not a real mitigation.
Database-engine layer (verified against a live PostgreSQL 16 instance's actual
parameter names and current defaults):
commit_delay(default0) andcommit_siblings(default5): lets the WAL writer
wait a short configurable delay before flushing, if at leastcommit_siblingsother
transactions are also waiting to commit, so onefsynccovers several commits instead
of one each. This is literally group-commit tuning, and it only helps when there's
actually something to group with.synchronous_commit(defaulton): can be relaxed tooff(the client is told the
transaction committed before its WAL record is confirmed flushed) or tolocal/
remote_writein a replicated setup. This directly trades durability, a crash inside
the risk window can lose the last fraction of a second of commits, for removing the
fsyncwait from the commit's critical path entirely. It's a named, explicit trade, not
free performance, and should be scoped to data that's been deliberately decided doesn't
need commit-level durability.wal_writer_delay(default200ms) andwal_writer_flush_after(default1MB):
control how eagerly a background process flushes WAL even without an explicit commit
waiting, smoothing out bursts before they reach the commit path at all.full_page_writes(defaulton): protects against a torn page (a page whose write a crash interrupted partway through, leaving it part old data and part new data on disk) by writing a full page image to WAL the first time a page is modified after a checkpoint (the
periodic point where all recently modified in-memory pages are flushed to disk).
Disabling it reduces WAL volume, and thereforefsync-able bytes, but reopens a real
corruption risk on crash unless the storage layer independently guarantees atomic page
writes; generally not worth disabling on typical cloud block storage (a virtual disk volume attached to the database server over the network, like AWS EBS or its equivalent on other clouds).
Application layer. Batch logically related writes into fewer, larger transactions
instead of many small autocommitted ones, fewer commits means fewer fsync calls by
definition. For workloads that can tolerate it, route non-critical writes (telemetry that
can be reconstructed, isn't itself a source of truth) to an asynchronous queue instead of
a synchronous commit on the hot path, reserving the fsync-bound path for writes that
genuinely need the guarantee.
Worked example: what group commit buys you
Without grouping, every commit pays its own fsync serially. With an illustrative
single-fsync latency of 4.0 ms:
TPSmax=Lms1000=4.01000=250 commits/sec
With group commit able to batch up to 50 transactions arriving within a 2.0 ms window
behind one fsync:
TPSgrouped=twindowbatch_size=0.002s50=25,000 commits/sec
This is an idealized upper bound: it assumes 50 concurrent committers genuinely arrive
within every 2 ms window. Group commit only helps when there's real concurrent commit
traffic to group; a single serial client committing one transaction at a time sees no
benefit from it at all, since there's never anyone else to batch with.
Trade-offs and pitfalls
- Reaching for
synchronous_commit = offfirst because it's the largest number on paper;
it should be the last lever, reserved for data explicitly decided not to need
commit-durability, not a blanket performance setting applied everywhere. - Mistaking a virtualized disk that returns from
fsyncquickly for a fast, durable disk;
verify the storage layer's actual durability guarantee, not just its latency. - Group commit helps throughput under concurrent load but can slightly increase the
latency of any individual commit under low concurrency, it may wait up to
commit_delayfor siblings that never show up, so validate it against your actual
traffic's concurrency profile rather than assuming it's a free win in every situation. - Disabling
full_page_writeswithout confirming the storage layer's atomic-write
guarantees trades a rare torn-page corruption risk for WAL savings that group commit or
WAL compression could usually get more safely.
You need a materialized view or aggregated table that stays fresh within 10 minutes over a dataset ingesting about 1TB per day. How would you implement incremental refresh so you are not recomputing the whole view every time, what would partition-level refresh buy you, and what do you do when the underlying data gets corrected after the fact through a backfill?
Sample Answer
Approach
Native Postgres materialized views do not support incremental refresh, REFRESH MATERIALIZED VIEW, with or without CONCURRENTLY, always recomputes the entire result. To hit a 10-minute freshness target over a dataset ingesting about 1 TB per day, build incremental refresh directly: a trigger-maintained summary table that updates only the rows affected by each write, so refresh cost scales with the volume of change since the last update, not with total dataset size.
The source table should itself be partitioned by time at this ingest rate (the same pattern used for retention). Combined with the trigger, a single write never touches a whole day's partition to update the summary, only the specific rows its own key maps to. If a full rebuild is ever needed, recovering from a bug, or onboarding this pattern onto existing history, doing it one partition at a time bounds the blast radius and lets the rebuild be resumed or parallelized instead of running one multi-terabyte query. Partition-level refresh is what makes that occasional full-rebuild cost tractable; it is not what keeps the steady-state refresh fast, the trigger is.
When a late-arriving correction changes a row that was already aggregated, the trigger's update branch must remove the old value's contribution before adding the new one. A naive trigger that only handles inserts will silently under- or over-count once corrections start arriving, a very easy defect to ship and not notice until the numbers are visibly wrong.
Code
CREATE TABLE daily_customer_totals (
customer_id int,
day date,
total_amount numeric(12,2) NOT NULL DEFAULT 0,
n_orders int NOT NULL DEFAULT 0,
PRIMARY KEY (customer_id, day)
);
CREATE OR REPLACE FUNCTION orders_incremental_refresh() RETURNS trigger AS $$
BEGIN
IF TG_OP = 'INSERT' THEN
INSERT INTO daily_customer_totals (customer_id, day, total_amount, n_orders)
VALUES (NEW.customer_id, NEW.created_at::date, NEW.amount, 1)
ON CONFLICT (customer_id, day) DO UPDATE
SET total_amount = daily_customer_totals.total_amount + EXCLUDED.total_amount,
n_orders = daily_customer_totals.n_orders + 1;
RETURN NEW;
ELSIF TG_OP = 'UPDATE' THEN
UPDATE daily_customer_totals
SET total_amount = total_amount - OLD.amount, n_orders = n_orders - 1
WHERE customer_id = OLD.customer_id AND day = OLD.created_at::date;
INSERT INTO daily_customer_totals (customer_id, day, total_amount, n_orders)
VALUES (NEW.customer_id, NEW.created_at::date, NEW.amount, 1)
ON CONFLICT (customer_id, day) DO UPDATE
SET total_amount = daily_customer_totals.total_amount + EXCLUDED.total_amount,
n_orders = daily_customer_totals.n_orders + 1;
RETURN NEW;
ELSIF TG_OP = 'DELETE' THEN
UPDATE daily_customer_totals
SET total_amount = total_amount - OLD.amount, n_orders = n_orders - 1
WHERE customer_id = OLD.customer_id AND day = OLD.created_at::date;
RETURN OLD;
END IF;
END;
$$ LANGUAGE plpgsql;
CREATE TRIGGER trg_orders_incremental
AFTER INSERT OR UPDATE OR DELETE ON orders
FOR EACH ROW EXECUTE FUNCTION orders_incremental_refresh();
Key points
The ON CONFLICT ... DO UPDATE upsert makes the insert branch safe under concurrent writers touching the same key. The update branch's subtract-then-add is what makes backfill corrections correct, the detail a naive implementation misses. An initial backfill over existing historical data must overwrite (DO UPDATE SET ... = EXCLUDED...), not skip on conflict, or any row already touched by the trigger before the backfill ran is left with an incorrect partial value forever, this exact bug was caught and fixed while building this demo: an initial ON CONFLICT DO NOTHING backfill left an undercount for every key the trigger had already partially updated.
Worked example
Real, executed output. A new $100.00 order inserted for an existing customer was picked up immediately with no full recompute, total_amount moved from 15,911.04 to 16,011.04 and n_orders from 55 to 56. Then a backfill correction, adding $50.00 to an existing order's stored amount:
UPDATE orders SET amount = amount + 50 WHERE id = 765;
The summary correctly reflected it, total_amount became 16,061.04 and n_orders stayed at 56, a correction to an existing order, not a new one. Verified directly against a fresh SUM/COUNT over the source rows for that same customer and day, which returned the identical 16061.04 / 56, confirming the incremental summary and a full recompute agree exactly.
Complexity
Each write costs O(1) additional work, one indexed upsert on the summary table's primary key, independent of total dataset size. A full rebuild, if ever needed, costs O(n) over whatever slice is rebuilt, which is why partition-level rebuilds bound that cost instead of paying O(n) over the entire multi-terabyte history at once.
Edge cases
Concurrent writers updating the same customer-and-day key simultaneously: ON CONFLICT DO UPDATE makes this safe at the row level since Postgres serializes conflicting upserts on the same key, but two concurrent update-branch subtract-then-add sequences on the same key are two separate statements each, not one atomic unit, so under heavy concurrent correction traffic on the same key this needs review for whether it should run inside a stricter transaction boundary. A delete of a row inserted before the trigger existed, with no corresponding insert having ever incremented the summary, would drive the summary negative, worth a floor check or a periodic reconciliation job. A correction that moves a row from one day bucket to another (changing created_at itself) is already handled correctly, since the trigger looks up the old day for the subtraction and the new day for the addition separately, but it is easy to get wrong if the day were stored as its own independent column instead of derived from created_at on each side of the trigger.
Trade-offs and pitfalls
This pattern is custom code per aggregation shape, unlike a native materialized view, every new summary needs its own trigger logic, that is the real cost of O(1) incremental refresh. Triggers add write-path latency and lock contention on the source table, every insert now also performs a summary-table upsert inside the same transaction, profile this under real write concurrency before assuming it is free. At very high write volume the trigger itself can become the bottleneck it was meant to avoid, since every source write now also pays for an extra upsert inside the same transaction; the usual next step up is a change-data-capture (CDC) pipeline, which reads the database's replication log to build the summary from a stream of row changes outside the source transaction entirely, trading the trigger's simplicity for a separate system to run and monitor. A bug in the trigger logic, as demonstrated above with the ON CONFLICT DO NOTHING backfill mistake, corrupts the summary silently while the source data stays correct, so this pattern needs its own periodic reconciliation check comparing the summary against a fresh aggregate on a sample, precisely because there is no database-native correctness guarantee the way a REFRESH-based materialized view has: however stale, a materialized view is always exactly correct as of some point in time, while a hand-rolled incremental one can be wrong at every point in time if the trigger logic has a bug.
Explain what VACUUM and ANALYZE do in PostgreSQL, why autovacuum exists, and what happens to a busy transactional table if neither runs. How would you detect that a table has become bloated, tune autovacuum's parameters for a high-write table, and safely run a manual VACUUM or VACUUM FULL on a production system without causing an outage?
Sample Answer
Direct answer
PostgreSQL keeps old row versions on disk under multi-version concurrency control (MVCC), a design where an UPDATE never overwrites a row in place; it writes a new version and leaves the old one behind, and a DELETE marks a row dead rather than removing it immediately. VACUUM reclaims the space held by those dead row versions so it can be reused by future writes, and ANALYZE refreshes the query planner's statistics (row counts and value distributions) so the optimizer keeps choosing good plans. Autovacuum exists because both of these were originally manual chores that got neglected under real workloads, so it runs both automatically in the background, triggered by how many rows have changed rather than by a fixed schedule. If neither runs on a busy transactional table, dead row versions pile up (bloat): the table and its indexes grow larger than the live data justifies, every scan has to skip past more dead space to find live rows, and in the extreme case of autovacuum being disabled for a very long time, PostgreSQL can be forced into emergency single-user vacuum mode to protect against transaction ID wraparound, a hard stop rather than a slow degradation.
Structured elaboration
Why dead rows accumulate: MVCC in one paragraph
Every row version carries hidden xmin and xmax columns recording which transaction created it and, if applicable, which transaction deleted or replaced it. A running transaction's snapshot determines which versions it is allowed to see, which is how PostgreSQL lets readers proceed without blocking writers and vice versa. The cost of that design is that old versions cannot simply disappear the instant they are superseded: some other transaction might still need to see them. VACUUM's job is to find versions that no currently running transaction could possibly still need, and mark that space reusable.
Detecting bloat
pg_stat_user_tables exposes n_live_tup and n_dead_tup per table, along with last_autovacuum and autovacuum_count, which together show whether autovacuum is actually keeping up. The pgstattuple extension gives a more precise, directly measured dead_tuple_percent and free_percent at the cost of a full scan. A large, growing gap between a table's on-disk size and what its live row count would predict is the practical symptom to watch for.
Tuning autovacuum for a high-write table
The default triggers (autovacuum_vacuum_threshold plus autovacuum_vacuum_scale_factor, which by default fire a vacuum once roughly 20% of a table's rows are dead) work fine on small and medium tables, but on a large, high-write table waiting for 20% of a huge row count to go dead means a huge amount of bloat accumulates between vacuum runs. Lower autovacuum_vacuum_scale_factor per table (for example to 0.05) so vacuum triggers on a smaller fraction, and consider autovacuum_analyze_scale_factor similarly for statistics freshness. If vacuum then runs so often that it competes with foreground I/O, autovacuum_vacuum_cost_delay and autovacuum_vacuum_cost_limit throttle how aggressively it consumes I/O per pass, trading vacuum speed for less foreground impact.
Running VACUUM and VACUUM FULL safely
Plain VACUUM is safe to run at any time against a live table: it does not block ordinary reads or writes. It reclaims dead-tuple space for reuse by future writes on the same table, but it does not shrink the file on disk; the reclaimed space stays allocated to that table for its own future use. VACUUM FULL is a different operation entirely: it rewrites the whole table into a new, compact file and requires an ACCESS EXCLUSIVE lock (the strongest lock PostgreSQL has) for the full duration, which blocks all reads and writes against the table, not just other writers. That is why VACUUM FULL is not a stronger version of routine maintenance, it is closer to a scheduled outage for that table. The safe pattern: avoid it for routine bloat control (a properly tuned autovacuum should prevent needing it), run it in a maintenance window when it genuinely is needed, and set lock_timeout before attempting it against a live system so a blocked attempt fails fast with a clear error instead of queuing up and blocking every new query behind it while it waits. For reclaiming space without an exclusive-lock window at all, the pg_repack extension rebuilds a table's storage online via a shadow copy and swap.
Worked example
Everything below was executed against a real PostgreSQL 16 instance. First, simulate a busy transactional table: 200,000 rows, then three full-table updates in a row (each UPDATE leaves the previous row version behind as dead space):
CREATE TABLE orders (id bigint PRIMARY KEY, status text NOT NULL, payload text);
INSERT INTO orders SELECT g, 'pending', repeat('x', 200) FROM generate_series(1, 200000) g;
ANALYZE orders;
UPDATE orders SET status = 'processing';
UPDATE orders SET status = 'shipped';
UPDATE orders SET status = 'delivered';
SELECT tuple_count AS live_tuples, dead_tuple_count AS dead_tuples, dead_tuple_percent,
pg_size_pretty(table_len) AS heap_only
FROM pgstattuple('orders');
SELECT pg_size_pretty(pg_total_relation_size('orders')) AS total_with_index;
pg_stat_user_tables.n_live_tup / n_dead_tup are asynchronous, incrementally-maintained
estimates, not a live physical count: right after a burst of writes they can be transiently
wrong (this table briefly reported 400,000 live tuples immediately after the load, double
the real 200,000, until the next ANALYZE or autovacuum corrected it) and they do not settle
to a trustworthy number on a predictable schedule. pgstattuple does a real scan and gives
an exact answer instead, so it is the more honest tool for a one-time bloat check like this
one; use n_live_tup/n_dead_tup for autovacuum's own trigger logic, not as ground truth
for a point-in-time reading.
Output (executed):
live_tuples | dead_tuples | dead_tuple_percent | heap_only
-------------+-------------+---------------------+-----------
200000 | 200000 | 23.83 | 195 MB
total_with_index
-------------------
204 MB
Only 200,000 dead tuples are physically present here, not 600,000 (one per UPDATE times
three updates), because Postgres opportunistically prunes a page's superseded row versions
during normal access once no transaction's snapshot can still need them; with three
UPDATEs run back to back and nothing holding an old snapshot open between them, each
UPDATE's write prunes the previous one's now-unreachable dead version on the way past.
Only the most recent update's dead versions are left standing by the time anything reads
the table. A long-running transaction started before the updates would hold those older
versions live and prevent this pruning, which is exactly the "a single long-open session
can defeat cleanup" pitfall below.
live, dead = 200000, 200000
heap_fresh_bytes, heap_now_bytes = 51200000, 204800000
print(f"{dead/(live+dead)*100:.0f}% of the table's own dead_tuple_percent by byte count (pgstattuple, exact scan)")
print(f"table's heap grew from {heap_fresh_bytes/1024/1024:.2f} MB freshly loaded to {heap_now_bytes/1024/1024:.2f} MB: {heap_now_bytes/heap_fresh_bytes:.2f}x, from dead versions plus page overhead")
Output (executed):
50% of the table's own dead_tuple_percent by byte count (pgstattuple, exact scan)
table's heap grew from 48.83 MB freshly loaded to 195.31 MB: 4.00x, from dead versions plus page overhead
(pgstattuple's own dead_tuple_percent reads 23.83%, not 50%: it is dead bytes over total
table bytes, not dead tuples over total tuples, so it is not the same ratio as the naive
tuple-count division above; both numbers are real, they just answer slightly different
questions, which is worth knowing before quoting either one out of context.)
Now a plain VACUUM (VERBOSE):
VACUUM (VERBOSE) orders;
Output (executed, trimmed to the relevant lines):
INFO: finished vacuuming "postgres.public.orders": index scans: 1
tuples: 200000 removed, 200000 remain, 0 are dead but not yet removable
WAL usage: 51095 records, 2 full page images, 5258435 bytes
SELECT n_dead_tup, pg_size_pretty(pg_relation_size('orders')) FROM pg_stat_user_tables WHERE relname = 'orders';
Output (executed):
n_dead_tup | pg_size_pretty
------------+-----------------
0 | 195 MB
n_dead_tup dropped to zero, but the table's on-disk size stayed at 195 MB, unchanged, exactly matching the heap-only reading taken right after the three updates and before this VACUUM ran. This is the core distinction: plain VACUUM freed the dead space for the table's own future reuse, it did not return a single byte to the operating system.
Next, VACUUM FULL while another session holds an open transaction against the table, with lock_timeout set so the attempt fails fast rather than hanging:
-- session A: opens a transaction and holds it open
BEGIN;
SELECT * FROM orders WHERE id = 1;
SELECT pg_sleep(6);
COMMIT;
-- session B, started while session A's transaction is still open:
SET lock_timeout = '2s';
VACUUM FULL orders;
Output (executed, session B):
ERROR: canceling statement due to lock timeout
That is the real, literal error text PostgreSQL returns: VACUUM FULL needed the ACCESS EXCLUSIVE lock, session A's open transaction (which had already run a SELECT, so it was still holding a lock on the table) was in the way, and the lock_timeout made the attempt fail cleanly after two seconds instead of blocking indefinitely, which itself would have started blocking every subsequent query behind it once it entered the lock queue. Once session A committed and no one else held a conflicting lock, the same command succeeded:
VACUUM FULL orders;
SELECT pg_size_pretty(pg_relation_size('orders')), pg_size_pretty(pg_total_relation_size('orders'));
Output (executed):
pg_size_pretty | pg_size_pretty
-----------------+-----------------
49 MB | 53 MB
The table shrank from 195 MB back down to 49 MB (53 MB total including its index), matching its freshly loaded size almost exactly, because VACUUM FULL rewrote it from scratch with none of the dead space plain VACUUM had only marked reusable.
Tuning autovacuum for this table for future high-write periods:
ALTER TABLE orders SET (autovacuum_vacuum_scale_factor = 0.05, autovacuum_vacuum_cost_delay = 2);
Trade-offs & pitfalls
VACUUM FULLis not "VACUUM but more thorough." It is a fundamentally heavier, blocking operation. Reaching for it as a first response to bloat, instead of first checking whether ordinary autovacuum is simply falling behind and needs its thresholds tuned, is the most common wrong turn here.- Very aggressive autovacuum tuning is a real trade-off, not a free win. A very low scale factor triggers vacuum more often, which is exactly what a high-write table needs, but without a matched
autovacuum_vacuum_cost_delayit can compete with foreground query I/O instead of staying in the background. - A single long-running or idle-in-transaction session can defeat vacuum entirely, no matter how well it is tuned: PostgreSQL cannot remove a dead row version that an old transaction's snapshot might still be able to see, so vacuum will keep running and keep finding nothing removable until that transaction ends. Check
pg_stat_activityfor long-open transactions before assuming a vacuum-tuning problem. - Disabling autovacuum on a noisy table is not a fix, it defers the failure mode. Left unchecked long enough, it converts a manageable, gradual bloat problem into a hard transaction ID wraparound stop, which is a sudden availability incident, not a performance degradation you can watch coming.
VACUUM FULLneeds roughly the table's own size again in free disk space while it builds the new compact copy; on a large table, checking for that headroom first is not optional.
You need to cut p95 read latency for a global user base from 120 milliseconds to under 40 milliseconds at 50,000 reads per second. Design a read-scaling architecture layering edge or CDN caching, in-memory caching, and read replicas. How would you decide what to cache and for how long, using a pattern like stale-while-revalidate, and what consistency guarantees would you be willing to give up to hit that target?
Sample Answer
Direct answer
Layer three caches in front of the primary: an edge / content delivery network (CDN) tier for data that is the same for everyone or close to it, a regional in-memory cache (Redis or similar) for personalized or fast-changing data, and regional read replicas as the fallback source once both caches miss. Whether this works comes down to one number: because a request from a distant user has to cross real physical distance to reach a single origin region, a cache miss alone can already exceed a 40 millisecond budget, so hitting p95 (95th-percentile) latency under 40 milliseconds at 50,000 reads per second requires the combined edge-plus-in-memory cache hit ratio to sit at roughly 95% or higher. The guarantee you give up to get there is strong, always-fresh reads for everyone except the user who just wrote the data: cached and replica-served reads become eventually consistent, bounded by a stated staleness window per data type, while the acting user's own write is routed around the cache so they never see their own action appear to have failed.
Structured elaboration
Why distance alone rules out an origin-only design
Signal propagation through fiber-optic cable travels at roughly two-thirds the speed of light in vacuum. That is a physical limit, not a database tuning problem: a round trip between two points thousands of kilometers apart costs real milliseconds no query optimization can remove. A single-origin-region architecture, however well-tuned the database is, cannot put a p95 under 40 milliseconds for users on the other side of the planet. That is the architectural reason edge and regional caching are not optional here, they are the only way to physically shorten the path for most requests.
Classifying data to decide what goes in which tier
| Data shape | Cache tier | Reasoning |
|---|---|---|
| Same for (almost) every user, changes rarely (catalog, public profile pages, static configuration) | Edge / CDN, long time-to-live (TTL) | Few distinct cache keys shared by many users means very high hit ratios are achievable at the point closest to the user |
| Personalized or computed per-user (a feed, an aggregate, a leaderboard) | Regional in-memory cache, shorter TTL | Too many distinct keys for edge caching to pay off; centralizing in one regional cache still avoids hitting the database on every read and is easier to invalidate correctly |
| The user's own just-written data | No cache; routed to the primary or a replica confirmed to have absorbed the write | This is the one case where staleness is not acceptable: a user should never read back something that looks like their own write silently failed |
| Anything that must be correct at the instant of read (for example, a balance immediately before authorizing a debit) | No cache, always the primary | An explicit, named exception, not something quietly caught by a general TTL |
Stale-while-revalidate (SWR) and why it, not active invalidation, should carry most of the load
Under stale-while-revalidate, a cache serves the value it already has immediately, even if its soft freshness window has passed, while kicking off a background refresh from the layer below (an in-memory cache refreshing from a replica, or an edge node refreshing from origin). The next request after the refresh completes gets the new value. The user never waits on the refresh. Compare that to actively pushing invalidation messages to every edge node whenever underlying data changes: that turns cache freshness into its own distributed consistency problem across every point of presence (PoP) in the network, and for most read-mostly data it is slower and more fragile than simply letting a short TTL expire. Reserve active invalidation for the small set of high-value keys where even a short TTL window of staleness is unacceptable (for example, a price that just changed).
Read replica routing and the freshness floor they set
Give each major region its own read replica, route reads to the nearest one, and continuously track each replica's replication lag (how far behind the primary it is). If a replica's lag exceeds the staleness budget promised by the cache tiers above it, fail toward the primary for that traffic rather than silently serving data staler than what was promised. Replication lag is also the floor beneath every promise this design makes: no cache TTL should be tighter than what the replica feeding it can actually guarantee.
Worked example
Two numbers actually gate this design: how close a hit has to be to beat 40 milliseconds, and how far a miss is from a distant origin region. The second one is physics, not opinion, so it is worth computing rather than assuming.
C_VACUUM = 299_792.458 # km/s, speed of light in vacuum
FIBER_INDEX = 1.47 # typical single-mode fiber refractive index
v_fiber = C_VACUUM / FIBER_INDEX # km/s, signal speed in fiber
routes_km = {
"cross-country (US east-west)": 4200,
"cross-continent (US-Europe)": 6000,
"antipodal (US-Singapore)": 15300,
}
for name, km in routes_km.items():
rtt_ms = 2 * km / v_fiber * 1000 # round trip: there and back
print(f"{name}: {rtt_ms:.1f} ms of propagation delay alone")
# p95 sizing: if a hit is comfortably under 40ms and a miss stays near today's
# 120ms baseline (the same WAN + database round trip that produced it), then
# for the 95th-percentile request to land under 40ms, at most 5% of requests
# may miss.
miss_budget_pct = 5.0
print(f"required combined edge+in-memory hit ratio >= {100 - miss_budget_pct:.0f}%")
read_qps = 50_000
target_hit_ratio = 0.97 # 2 points of margin above the 95% floor
miss_qps = read_qps * (1 - target_hit_ratio)
regions = 5
print(f"miss traffic reaching replicas: {miss_qps:,.0f} reads/sec, "
f"~{miss_qps/regions:,.0f} reads/sec per region across {regions} regions")
write_tps = 5_000
print(f"write target: {write_tps:,} TPS (transactions per second), unaffected by "
f"any of the caching above; it lands entirely on the primary")
Output (executed):
cross-country (US east-west): 41.2 ms of propagation delay alone
cross-continent (US-Europe): 58.8 ms of propagation delay alone
antipodal (US-Singapore): 150.0 ms of propagation delay alone
required combined edge+in-memory hit ratio >= 95%
miss traffic reaching replicas: 1,500 reads/sec, ~300 reads/sec per region across 5 regions
write target: 5,000 TPS (transactions per second), unaffected by any of the caching above; it lands entirely on the primary
Two things fall out of this. First, cross-continent and antipodal round trips already exceed 40 milliseconds on propagation delay alone, before a single database query runs, which is why no amount of origin-side tuning closes this gap for a global user base; the cache hit ratio has to do essentially all the work. Second, at a 97% hit ratio the replica tier only has to absorb about 300 reads per second per region, comfortably inside what a single well-indexed read replica handles, so the hard part of this design is the cache hit ratio, not replica capacity.
The companion write target matters here because it sets the freshness floor discussed above: 5,000 writes per second all land on the primary (caching and read replicas do nothing for writes), and the primary has to both sustain that rate and ship its write-ahead log (WAL, the durable record of every change) to every regional replica fast enough that replication lag stays inside whatever staleness window the cache tiers are promising. If write volume grows further, replication lag grows with it, and that is the point at which every cache TTL and every "how stale can a replica read be" assumption in this design needs to be re-validated, not treated as a one-time decision.
flowchart TD
U[Global user request] --> E["Edge / CDN PoP\ncache-only, long TTL"]
E -- hit --> RU[Response]
E -- miss --> R["Regional in-memory cache\nRedis, short TTL, SWR"]
R -- hit --> RU
R -- miss --> RR[Regional read replica]
RR -- fresh enough --> RU
RR -- lag over budget --> P[Primary]
P --> RU
W[Write / user action] --> P
P -- WAL replication --> RR
P -- fresh write --> OW["Read-your-own-write\nbypass cache, route to primary"]
Trade-offs & pitfalls
- Cache stampede on expiry. When a hot key's TTL lapses, many concurrent requests can all miss at once and hammer the layer below simultaneously. Stale-while-revalidate already helps (it keeps serving the old value during refresh), but pair it with request coalescing (only one in-flight refresh per key, other requests wait on that result) and TTL jitter (randomizing expiry slightly so many keys do not expire in the same instant).
- Over-caching mutable, high-stakes data. The temptation to cache everything for uniform simplicity is real; resist it for anything where staleness has a real cost (money, inventory at the moment of purchase, access control decisions). Naming the exceptions explicitly, as this design does, is safer than a blanket policy with silent edge cases.
- Edge caching does not help personalized data. If a page or response is unique per user, the edge tier will have as many cache keys as users and its hit ratio collapses; that traffic belongs in the regional in-memory tier, not forced into the edge tier for architectural tidiness.
- The 40ms budget math here is a simplified two-class model (hit versus miss), used to size the required hit ratio, not a promise about the real latency distribution. Validate the actual design against measured percentile telemetry once it is live; a real cache miss is not always exactly 120 milliseconds, and a real hit is not always a few milliseconds.
- Read-scaling work does not touch the write path. If the write target grows past 5,000 TPS, that is a primary-capacity and replication-fan-out problem, a different constraint than anything cache tiers or read replicas solve, and it should be sized and reasoned about on its own terms.
Tell me about a time you optimized a database's performance in production. What was slow, what metrics did you start with (latency, throughput), what did you actually do (indexing, partitioning, denormalization, query rewrites, or something else), how did you validate the improvement, and what trade-offs did you accept?
Sample Answer
Direct answer
Pick one real incident with one measurable starting symptom, latency or throughput are
the usual candidates, name specifically what you tried and why each step was chosen over
the alternatives, and close with how you validated the improvement and what it cost you
elsewhere. A grab-bag of "things I've optimized over the years" reads as unfocused; one
incident, narrated end to end, reads as judgment.
A scaffold for the story
Situation. State the symptom in a measurable way: "the checkout endpoint's p95
(95th-percentile) database time was regularly exceeding 400 ms during traffic peaks" beats "the database
was slow." Latency and throughput are the two most common starting metrics, latency
because users feel it directly, throughput because it caps how much load the system can
absorb before latency degrades for everyone.
Task. What target made the incident worth fixing: an SLA being missed, a
capacity ceiling being approached, a specific complaint.
Action. Narrate three concrete techniques in the order you tried them, not a list.
The available categories a real answer usually draws from are indexing, query rewrites,
partitioning, denormalization (storing redundant, precomputed copies of data to avoid recomputing it
on every read), or connection/configuration tuning; pick whichever three
actually applied to your incident. A representative shape:
- Confirmed the query's plan with
EXPLAIN ANALYZE(the command that runs the query and
reports the real plan and timing, not an estimate) and added an index matching the
filter and sort the query actually used, verifying the plan flipped from a sequential
scan to an index scan. - Found the query was also pulling unused joined columns (
SELECT *where three columns
were used); rewrote it to select only what was needed and pushed a filter earlier in
the plan. - Noticed via the database's activity view that many requests were queued waiting for a
pooled connection, not for query execution; resized the pool and added a
statement_timeoutso one runaway query couldn't starve it.
Result. Validate against the same metric you started with, measured under
comparable load, not a different, easier one. Watch it hold up over days in production,
not just at deploy time, since delayed effects (maintenance overhead, data growth) can
erode a fix that looked clean on day one.
Trade-off accepted. A senior answer names a real cost. For example: the new index
sped up reads but added overhead to a high-insert table's write path; that was monitored
afterward and accepted as a small, budgeted latency increase on writes rather than
blocking the read fix on eliminating it entirely.
Trade-offs and pitfalls
- Listing several unrelated wins instead of one incident reads as a highlight reel, not
evidence of how you reason through a problem under pressure. - Skipping the "why this technique, in this order" reasoning makes it sound like guessing
that happened to work, which is the opposite of what the question is testing for. - Claiming an improvement without describing how it was validated, and against what
baseline, is the single most common gap; interviewers will ask. - Not naming a trade-off accepted reads as either inexperience or a curated story; real
fixes cost something somewhere.
Unlock Full Question Bank
Get access to all 32 Database Performance Tuning and Scaling interview questions and detailed answers.
Sign in to ContinueJoin thousands of developers preparing for their dream job.