Why This Matters#
PostgreSQL is not the "boring relational option" you pick before the interesting part of the design. It is the default system of record for most of the industry, and the real Staff question is never "should we use Postgres?" — it is "what is the event that forces us off a single Postgres primary, and how far away is it?" Most systems never reach it. The ones that do reach it through one of four doors: write throughput on one primary, table size that vacuum can't keep up with, connection count, or a multi-region requirement. Knowing which door you'll hit first, and when, is the design.
It keeps showing up in interviews because candidates reach for NoSQL too early. Interviewers want to see that you can get 10–50K write TPS and 100K+ read QPS out of one well-run Postgres cluster before introducing a distributed database and its operational tax. They also want to see that you know Postgres's own tax: MVCC bloat, autovacuum, transaction ID wraparound, connection limits, and DDL locks that can take a site down with a one-line migration.
The L5 answer says "Postgres with read replicas." The L6 answer says "single primary on 32 vCPU handles our 8K write TPS with 3× headroom; PgBouncer in transaction mode caps backend connections at 200; orders is range-partitioned by month so retention is DROP PARTITION, not DELETE; replicas serve the catalog with a 1 s lag budget; checkout reads its own writes from the primary." The L7 answer asks how many teams are running unowned Postgres clusters with statement_timeout = 0, and what the org's sharding story is before the first team needs it.
The 60-Second Pitch#
"Postgres gives us ACID transactions, constraints that enforce correctness in the database rather than in every service, and a query planner that answers questions we haven't thought of yet. A single primary on modern hardware does 10–50K write transactions per second and holds a few TB comfortably; replicas scale reads. I'd use it as the system of record, enforce invariants with unique and foreign-key constraints, publish changes through an outbox table and logical replication, and only shard when a specific ceiling — write TPS, table size, or connection count — is within 12–18 months. When we shard, I'd shard by tenant so 99% of queries stay single-shard."
The Staff-level insight: Postgres's flexibility is the asset you're buying. A key-value store forces you to know your queries up front; Postgres lets you discover them. That makes it the right default whenever access patterns are still evolving — which is almost always the case in year one. You give that flexibility up only when you can name the ceiling that forces you to.
| Intent | Constraint | Postgres Strategy | Failure Mode | Correctness Bar |
|---|---|---|---|---|
| System of record (OLTP) | Invariants must hold | Normalized schema, constraints, Read Committed + targeted SELECT ... FOR UPDATE | Lock contention on hot rows | Zero lost or double-applied writes |
| Read-heavy product surface | p99 < 50 ms at 50K QPS | Replicas + covering indexes + cache | Replica lag shows stale data | Bounded staleness, read-your-writes where needed |
| Multi-tenant SaaS | Noisy tenants, isolation | tenant_id in every PK, RLS, later Citus / app sharding by tenant | One tenant's query plan starves others | Tenant data never crosses |
| Analytics / reporting | Big scans | Replica or CDC to warehouse — not the primary | Long queries block vacuum on primary | Approximate is fine; freshness minutes |
🎯 Staff Move: "I'll treat Postgres as the system of record and commit to one primary until we can name the ceiling. At 2K writes a second today and 3× annual growth, we have roughly two years before write TPS matters — so I'll choose the shard key now and not build sharding yet."
Architecture & Internals#
Only the internals that change design decisions.
Process-per-Connection#
Every client connection is a forked OS process with its own memory (~5–10 MB baseline, more with work_mem-heavy queries). Past a few hundred active backends, context switching, lock-manager contention, and snapshot computation (GetSnapshotData scales with connection count) degrade throughput. The practical rule: ~2–4 active connections per CPU core is the throughput peak; everything beyond that is queueing inside the database instead of outside it.
This is why every production Postgres at scale sits behind a pooler — PgBouncer, RDS Proxy, or pgcat — in transaction mode: 5,000 client connections multiplexed onto 100–300 server connections. Transaction mode breaks session state (session-level SET, advisory locks held across transactions, LISTEN, and — before PgBouncer 1.21 — protocol-level prepared statements), which is the tradeoff you name.
| Pooling Mode | Server Connection Held For | Multiplexing Ratio | Breaks | Pick When |
|---|---|---|---|---|
| Session | The whole client session | ~1:1 | Nothing | Legacy apps; admin tools |
| Transaction | One transaction | 10–50:1 | Session SET, cross-txn advisory locks, LISTEN/NOTIFY, older prepared-statement handling | Default for web/API tiers |
| Statement | One statement | Highest | Multi-statement transactions | Rare — autocommit-only workloads |
🎯 Staff Move: "I'll size the server pool at about 3× cores — ~100 connections on a 32-vCPU primary — and let 4,000 app connections queue in PgBouncer. Queueing in the pooler is cheap and visible; queueing inside Postgres is lock contention nobody can see."
MVCC, Dead Tuples, and Vacuum#
Postgres never updates a row in place. An UPDATE writes a new tuple version and marks the old one dead (via xmax). Readers see the version valid for their snapshot, so readers never block writers and writers never block readers. The bill arrives later:
- Dead tuples accumulate until
VACUUMmarks their space reusable. An update-heavy table with lagging vacuum bloats 2–10× its live size; indexes bloat too. - Long-running transactions hold back the xmin horizon. A single 6-hour analytics query or an
idle in transactionsession prevents vacuum from removing any tuple newer than its snapshot — database-wide. - HOT updates (heap-only tuples) avoid index writes when no indexed column changes and the page has free space. Setting
fillfactor = 80–90on hot-update tables keeps room on the page; every extra index on an updated column kills HOT for that update. - Transaction ID wraparound. XIDs are 32-bit; after ~2 billion transactions old rows must be "frozen." If autovacuum can't freeze fast enough, Postgres enters a protective mode and stops accepting writes until a manual vacuum completes — historically hours to days on multi-TB tables. At 10K TPS you consume 2B XIDs in ~2.3 days, so anti-wraparound vacuum must run continuously and never be blocked.
WAL, Checkpoints, and Durability#
Every change is written to the write-ahead log before the data page. A commit is durable when its WAL record is flushed (fsync) — typically 0.1–1 ms on NVMe, 1–3 ms on network block storage. Dirty data pages are written later at checkpoints (checkpoint_timeout 5–15 min, max_wal_size 4–64 GB on busy systems). The first modification of a page after a checkpoint writes the full page image to WAL (8 KB), which is why random-key inserts (UUIDv4 primary keys) can multiply WAL volume several-fold compared to sequential keys — each insert touches a different B-tree leaf page.
Replication#
| Mode | Mechanism | Lag | Use |
|---|---|---|---|
| Streaming (physical) | Ships WAL bytes; replica is a byte-identical copy | ms–seconds; spikes under heavy writes or replica load | HA standbys, read replicas |
| Synchronous streaming | Commit waits for standby flush (synchronous_standby_names) | Adds 1 RTT (0.5–2 ms same region, 30–80 ms cross-region) | Zero-loss failover within a region |
| Logical replication / decoding | Decodes WAL into row changes per table (publications/slots) | Seconds | CDC to Kafka (Debezium), major-version upgrades, partial replication |
Replicas can serve reads but face query conflicts: vacuum on the primary removes tuples a long replica query still needs. You choose between canceling the replica query (max_standby_streaming_delay) or letting replica lag grow, or hot_standby_feedback = on — which makes the replica's long queries hold back vacuum on the primary. There is no free option; name the one you pick.
The Planner Is a Dependency#
The cost-based planner chooses plans from table statistics (ANALYZE, default_statistics_target 100). A stale or skewed statistic flips a 2 ms index scan into a 40 s sequential scan — with no code change and no deploy. Plans also flip when a table crosses a size threshold or when a parameterized query switches from custom to generic plan after 5 executions. Plan instability is a production risk you own: monitor pg_stat_statements for mean-time regressions, pin critical statistics with extended statistics (CREATE STATISTICS) for correlated columns, and test migrations against production-sized data.
Data Modeling / Core Usage — "The Entire Game"#
In Postgres, the entire game is pushing invariants into the database and keeping the hot path on an index. Constraints are cheaper than distributed coordination, and an index scan is cheaper than any cache you'd put in front of a sequential scan.
Model Invariants as Constraints#
CREATE TABLE accounts (
id BIGINT GENERATED ALWAYS AS IDENTITY PRIMARY KEY,
tenant_id BIGINT NOT NULL,
email CITEXT NOT NULL,
balance_cents BIGINT NOT NULL CHECK (balance_cents >= 0), -- no overdraft, ever
version INT NOT NULL DEFAULT 0,
UNIQUE (tenant_id, email) -- dedup at the source
);
CREATE TABLE ledger_entries (
id BIGINT GENERATED ALWAYS AS IDENTITY PRIMARY KEY,
account_id BIGINT NOT NULL REFERENCES accounts(id),
idempotency_key TEXT NOT NULL,
amount_cents BIGINT NOT NULL,
created_at TIMESTAMPTZ NOT NULL DEFAULT now(),
UNIQUE (account_id, idempotency_key) -- retries become no-ops
);
CREATE INDEX ON ledger_entries (account_id, created_at DESC); -- FK + history query
The UNIQUE (account_id, idempotency_key) constraint replaces an idempotency service, a Redis lock, and a reconciliation job. The CHECK replaces a read-then-write race. Every constraint you move into the schema is a distributed-systems problem you no longer have.
Index Types and When Each Wins#
| Index | Good For | Cost / Caveat |
|---|---|---|
| B-tree (default) | Equality, range, ORDER BY, uniqueness | Every index adds ~1 extra write per insert/update of that column |
Composite (a, b) | Queries filtering a then b or sorting by b within a | Column order matters: leftmost prefix rule |
Covering INCLUDE (c) | Index-only scans, no heap fetch | Bigger index; needs visibility map (vacuum) to be effective |
Partial WHERE status = 'pending' | Queues, soft-deletes, hot subsets | Tiny index for a tiny hot set — often 100× smaller |
| GIN | JSONB containment, arrays, full-text (tsvector) | Slower writes; pending-list flushes cause latency spikes |
| GiST / SP-GiST | Geospatial (PostGIS), ranges, exclusion constraints | Larger, slower than B-tree for equality |
| BRIN | Huge append-only tables ordered by time | ~1000× smaller than B-tree; useless if data isn't physically ordered |
Keys: The Choice You Can't Easily Undo#
| Key | Pros | Cons | Pick When |
|---|---|---|---|
BIGINT IDENTITY | 8 bytes, sequential, B-tree friendly | Leaks volume; needs coordination across shards | Single-cluster tables, internal IDs |
| UUIDv4 | Globally unique, no coordination | 16 bytes, random inserts → page splits, cache misses, WAL full-page bloat | Only when IDs must be unguessable and table is small |
| UUIDv7 / ULID | Globally unique and time-ordered | 16 bytes | Default for new distributed-friendly schemas |
| Snowflake-style 64-bit | Time-ordered, embeds shard ID, 8 bytes | Needs an ID-generation scheme | Pre-sharded designs (Instagram's approach) |
JSONB: Flexible Columns, Not a Document Database#
JSONB with a GIN index lets you store sparse attributes without schema changes. Use it for truly variable attributes (product specs, webhook payloads), not to avoid modeling core fields. Rules: anything you filter, join, or constrain on becomes a real column; keep documents under ~2 KB to avoid TOAST (out-of-line storage above ~2 KB costs an extra fetch and rewrites the whole value on any update).
Postgres as a Queue: SKIP LOCKED#
-- Worker claims up to 10 jobs without blocking other workers
WITH next AS (
SELECT id FROM jobs
WHERE status = 'ready' AND run_at <= now()
ORDER BY run_at
LIMIT 10
FOR UPDATE SKIP LOCKED
)
UPDATE jobs SET status = 'running', locked_at = now(), attempts = attempts + 1
FROM next WHERE jobs.id = next.id
RETURNING jobs.*;
With a partial index ON jobs (run_at) WHERE status = 'ready', this sustains ~1–5K jobs/sec with transactional enqueue (the job commits atomically with the business write — the outbox property for free). Past ~5K/sec, or when finished jobs create more dead tuples than vacuum can clear, move to Kafka or a dedicated queue. See Distributed Job Scheduler.
The Outbox Pattern#
Never dual-write to Postgres and Kafka from the app — one will fail and they'll diverge. Write the event to an outbox table in the same transaction as the state change; Debezium tails the WAL through a logical replication slot and publishes to Kafka. Delivery is at-least-once; consumers dedupe by event ID. See Data Pipeline Patterns.
The Tunable Tradeoff — Isolation × Commit Durability#
Postgres gives you two independent dials. Most teams never touch either, which means they are running Read Committed with fully synchronous local commit — a fine default, but one they didn't choose.
Dial 1: Isolation Level#
| Level | What You Get | Anomalies Still Possible | Cost | Pick When |
|---|---|---|---|---|
| Read Committed (default) | Each statement sees data committed before it started | Non-repeatable reads, lost updates on read-modify-write, write skew | Cheapest | 90% of OLTP — with explicit row locks or conditional updates where it matters |
| Repeatable Read (snapshot isolation) | Whole transaction sees one snapshot; concurrent update of same row → error 40001 | Write skew (two txns read overlapping sets, write disjoint rows) | Retries on conflict | Reports needing a consistent view; multi-statement reads |
| Serializable (SSI) | Equivalent to some serial order | None | Predicate-lock tracking; 1–10% abort rate under contention, must retry | Invariants spanning rows ("at most 3 on-call doctors") where a constraint can't express it |
The L6 point: you usually don't need a higher isolation level — you need the right write pattern.
-- Lost-update-proof decrement without raising isolation:
UPDATE inventory
SET available = available - 1
WHERE sku = 'SKU-42' AND available >= 1 -- condition evaluated on the latest row version
RETURNING available; -- 0 rows returned = sold out
-- Optimistic concurrency with a version column:
UPDATE accounts SET balance_cents = $new, version = version + 1
WHERE id = $id AND version = $expected_version; -- 0 rows = someone else won, retry
Under Read Committed, a conditional UPDATE ... WHERE available >= 1 re-checks the condition against the newest committed version after waiting on the row lock — correct without Serializable. A hot row (one SKU in a flash sale) serializes at ~1–5K updates/sec because every writer waits for the previous row lock; that's a data-modeling problem (split inventory into N buckets), not an isolation problem. See Dealing with Contention and Flash Sales.
Dial 2: Commit Durability (synchronous_commit)#
| Setting | Commit Returns After | Loss on Primary Crash | Loss on Failover | Latency Added |
|---|---|---|---|---|
off | WAL in memory; flushed within ~3× wal_writer_delay (≈ 600 ms) | Up to ~600 ms of commits (no corruption) | Replication lag | None — often 2–5× write throughput on slow disks |
local | Local WAL flushed | 0 | Replication lag | Local fsync (0.1–3 ms) |
on (default) + no sync standby | Local WAL flushed | 0 | Replication lag | Local fsync |
remote_write + sync standby | Standby received WAL (OS buffer) | 0 | ~0 unless both nodes crash together | +1 RTT |
on + sync standby | Standby flushed WAL | 0 | 0 | +1 RTT + standby fsync |
remote_apply | Standby applied WAL — visible on standby reads | 0 | 0 | Highest; gives read-your-writes on that standby |
The setting is per transaction. The Staff pattern: SET LOCAL synchronous_commit = off for high-volume, low-value writes (analytics events, view counts, session touches) and full sync for ledger writes — in the same database.
effective_write_loss_on_failover ≈ replication_lag_at_crash (async standby)
≈ 0 (sync standby, quorum ANY 1 of 2)
commit_latency ≈ local_fsync + (sync ? RTT_to_standby + standby_fsync : 0)
🎯 Staff Move: "I'll run synchronous replication to one of two same-region standbys —
ANY 1 (s1, s2)— so a single standby failure doesn't stall writes, and we get zero acknowledged-write loss on failover for about 1 ms per commit. Clickstream writes opt out withsynchronous_commit = off."
Anti-Patterns — What Kills PostgreSQL Deployments#
1. The Migration That Takes the Site Down#
ALTER TABLE orders ADD COLUMN ... or CREATE INDEX (without CONCURRENTLY) takes an ACCESS EXCLUSIVE or SHARE lock. If a long-running query holds a conflicting lock, the DDL waits — and every new query on that table queues behind the waiting DDL. A 30-second analytics query plus one migration equals a 30-second full outage on the table.
Fix: every migration sets lock_timeout = '3s' and retries; indexes use CREATE INDEX CONCURRENTLY; constraints are added NOT VALID then VALIDATE CONSTRAINT (which takes a weaker lock); column type changes go through expand/contract (new column, dual-write, backfill in batches of 1–10K rows, switch reads, drop old). Adding a column with a constant default is metadata-only since Postgres 11; a volatile default rewrites the table.
2. Connection Storms#
Autoscaling adds 200 pods × pool of 20 = 4,000 connections. max_connections = 5000 "fixes" it until throughput collapses from backend contention. Fix: PgBouncer in transaction mode; server pool ≈ 2–4× cores; client pools of 5–10 per pod.
3. Idle-in-Transaction and Long Transactions#
An app opens a transaction, calls an external HTTP API for 30 s, then commits. Or an ORM leaves a transaction open between requests. Each holds row locks and the xmin horizon. Fix: idle_in_transaction_session_timeout = '60s', statement_timeout = '30s' for app roles (longer for batch roles), never do network I/O inside a transaction.
4. OFFSET Pagination on Big Tables#
OFFSET 1000000 LIMIT 50 reads and discards 1M rows. Fix: keyset pagination — WHERE (created_at, id) < ($last_ts, $last_id) ORDER BY created_at DESC, id DESC LIMIT 50 on a matching composite index; constant cost at any depth. See API Design Patterns.
5. Analytics on the Primary#
A dashboard runs a 10-minute GROUP BY across 500M rows on the primary: it evicts the hot working set from shared_buffers, holds back vacuum, and competes for I/O. Fix: analytics goes to a replica with hot_standby_feedback = off and tolerant query cancellation — or better, CDC into a warehouse.
6. Unbounded Growth With DELETE-Based Retention#
DELETE FROM events WHERE created_at < now() - interval '90 days' nightly on a 2 TB table generates hundreds of millions of dead tuples, WAL spikes, and replica lag. Fix: declarative range partitioning by day/month; retention becomes DROP TABLE events_2026_06 — instant, no dead tuples.
7. Missing Indexes on Foreign Keys#
Postgres does not auto-index the referencing column. DELETE FROM accounts WHERE id = 7 sequentially scans ledger_entries to check references — and holds locks while doing it. Fix: index every FK column used in joins or cascades; lint for it in CI.
| Anti-Pattern | Detection Signal | Blast Radius | Who Pays |
|---|---|---|---|
| Blocking DDL | pg_locks waiters > 50; lock wait graph rooted at DDL | Every query on the table | Customer-facing services, on-call |
| Connection storm | numbackends > 3× cores; CPU in LWLock waits | Whole database | All tenants of the cluster |
| Long transactions | max(now() - xact_start) > 5 min; n_dead_tup climbing | Database-wide bloat | Future query latency, storage bill |
| OFFSET pagination | pg_stat_statements mean time scales with page | One endpoint, heavy I/O | Primary CPU |
| Analytics on primary | Buffer hit ratio drops < 95% during reports | OLTP latency | Product SLOs |
| DELETE retention | WAL rate spikes nightly, replica lag > 30 s | Replicas and vacuum | Read path, backups |
| Unindexed FK | Seq scans on child tables during parent delete | Lock contention | Writers of both tables |
The Technology Landscape / Head-to-Head Comparison#
| Dimension | PostgreSQL (self/RDS) | Aurora PostgreSQL | MySQL / InnoDB | Citus | CockroachDB / YugabyteDB / Spanner | DynamoDB |
|---|---|---|---|---|---|---|
| Write scaling | One primary; ~10–50K TPS typical OLTP | One writer; storage scales to 128 TB | One primary; similar | Sharded across workers | Horizontally scalable writes | Unlimited with good keys |
| Transactions | Full ACID, SSI | Same | ACID, weaker default isolation semantics | Single-shard fast; cross-shard 2PC | Distributed serializable | Up to 100 items per TransactWriteItems |
| Latency (point write) | 1–3 ms | 2–5 ms (6-way quorum storage) | 1–3 ms | 2–10 ms | 5–20 ms in-region (consensus) | 5–10 ms |
| Failover | Patroni / Multi-AZ: 30–120 s | ~30 s typical | Similar to Postgres | Per-node HA | Automatic, seconds | Transparent |
| Replica lag | ms–s | Typically < 20–100 ms (shared storage) | ms–s | Per worker | Consensus — no lag for leaseholder reads | Eventually consistent reads < 1 s |
| Ecosystem | Richest: PostGIS, pgvector, extensions | Most extensions | Huge, simpler | Postgres extension | Postgres wire-compatible, subset | AWS-native |
| Pick when | Default system of record | Want managed + faster replicas, accept AWS lock-in and I/O pricing | Team expertise, Vitess sharding path | Multi-tenant SaaS outgrowing one node | Multi-region writes with strong consistency | Known access patterns, extreme scale, zero ops |
The Postgres-to-distributed-SQL decision is about the one-way door of latency: every write in a consensus-based database pays at least one quorum round trip. If your p99 budget is 50 ms and you're multi-region, that's acceptable. If you're single-region at 5K TPS, you're paying 5–10× write latency and new failure modes for scale you don't need.
🎯 Staff Insight: "Postgres or NoSQL?" is the wrong question. The right one is "which ceiling will we hit first — write TPS, table size, connections, or geography — and when?" If none is within 18 months, Postgres wins on flexibility and team familiarity. If geography is the ceiling, a distributed SQL database beats hand-sharding.
Patterns#
Pattern 1: Primary + Replicas with Routed Reads#
Writes and read-your-writes to the primary; everything else to replicas. The subtle part is routing: after a user writes, pin their reads to the primary for a window (e.g., 5 s), or read from a replica only if its pg_last_wal_replay_lsn() ≥ the LSN returned by the write. See Scaling Reads.
Pattern 2: Declarative Partitioning for Time-Series and Retention#
CREATE TABLE events (
tenant_id BIGINT NOT NULL,
event_id UUID NOT NULL, -- UUIDv7
created_at TIMESTAMPTZ NOT NULL,
payload JSONB,
PRIMARY KEY (tenant_id, created_at, event_id) -- partition key must be in every unique constraint
) PARTITION BY RANGE (created_at);
CREATE TABLE events_2026_09 PARTITION OF events
FOR VALUES FROM ('2026-09-01') TO ('2026-10-01');
-- pg_partman creates future partitions and drops expired ones
Partition pruning keeps queries that filter on created_at to 1–2 partitions. Keep partition count in the low hundreds; thousands of partitions slow planning. For heavier time-series needs see Time Series DBs.
Pattern 3: Multi-Tenancy#
| Model | Isolation | Tenant Count Ceiling | Ops Cost | Pick When |
|---|---|---|---|---|
Shared tables, tenant_id column + RLS | Logical | Millions | Lowest | Default SaaS |
| Schema per tenant | Namespace | ~1–5K per cluster (catalog bloat, migration fan-out) | Migrations × N | Regulated mid-market, per-tenant customization |
| Database per tenant | Strong | Hundreds per cluster | Highest | Enterprise contracts demanding isolation |
| Cell per group of tenants | Blast radius | Unlimited via more cells | Fleet tooling needed | Large SaaS at scale |
Lead every primary key and index with tenant_id from day one. It costs nothing now and makes tenant-based sharding a data move instead of a schema rewrite.
Pattern 4: Outbox + CDC#
Business write and outbox insert in one transaction; Debezium streams the WAL via a logical slot to Kafka. Watch the slot: an unconsumed replication slot retains WAL indefinitely and can fill the primary's disk (max_slot_wal_keep_size caps it at the cost of breaking the slot).
Pattern 5: Application-Level Sharding#
Route by shard key in the app or a proxy layer; each shard is an ordinary Postgres cluster. Pre-split into many logical shards (e.g., 512–4,096) mapped onto few physical clusters so growth is moving logical shards, not rehashing keys. Instagram's published ID scheme (timestamp + logical shard ID + sequence in 64 bits) and Notion's 480 logical shards over 32 physical databases are the canonical public examples. See Database Sharding.
Pattern 6: Postgres as the Everything-Store (Until It Isn't)#
Queue (SKIP LOCKED), search (tsvector + GIN), vectors (pgvector + HNSW), geospatial (PostGIS), cache (UNLOGGED tables). Each is good enough up to a threshold — roughly 1–5K jobs/s, 10M documents with simple relevance, 1–10M vectors, moderate geo QPS. The Staff move is naming the threshold at which each moves out.
| Pattern | Contract | Ceiling Signal | Next Step | Owner |
|---|---|---|---|---|
| Replicas | Bounded staleness | Replica lag p99 > budget | Cache, or remove replica-heavy query | Service team |
| Partitioning | Cheap retention | > 500 partitions or cross-partition queries dominate | Time-series DB | Service team |
| RLS multi-tenant | Logical isolation | One tenant > 20% of load | Move tenant to own cell | Platform + account team |
| Outbox + CDC | At-least-once events | Slot lag > 10 min | Scale consumers; cap slot WAL | Data platform |
| App sharding | Horizontal writes | Cross-shard queries > 5% | Re-pick shard key (painful) | Platform |
| Everything-store | Good enough | Feature's p99 or CPU share > 20% | Dedicated system | Service team |
Scaling#
The Four Ceilings#
| Ceiling | Typical Threshold (single primary, modern hardware) | Early Warning | First Response |
|---|---|---|---|
| Write TPS | ~10–50K simple OLTP write txns/s; less with many indexes | Primary CPU > 60% sustained; WAL generation > 100–200 MB/s | Batch writes, drop unused indexes, move low-value writes out |
| Table size | Single tables > 1–2 TB get painful: vacuum takes hours, index builds take hours | Autovacuum runs never finish; n_dead_tup monotonic | Partition; archive cold data |
| Connections | ~300–500 active backends | LWLock waits, CPU in kernel | Pooler; fewer, busier connections |
| Geography | One writer region | Remote-region write p99 > 100 ms | Regional read replicas; then distributed SQL or cells |
Back-of-Envelope Sizing#
| Input | Rule of Thumb | Example |
|---|---|---|
| Row size on disk | Payload + ~24-byte tuple header + ~4-byte item pointer, rounded to 8 bytes | 200-byte payload ≈ 230 bytes |
| Index size | ~key width + 8–16 bytes per entry, ~70–90% page fill | 1B rows × BIGINT index ≈ 25–30 GB |
| Table growth | rows/day × row size × (1 + index overhead ~0.5–1.0) | 50M rows/day × 230 B × 1.8 ≈ 20 GB/day |
| WAL volume | ~1–3× the logical write volume (more with random keys) | 20 GB/day data ≈ 40–60 GB/day WAL |
| Working set in RAM | Hot indexes + hot rows should fit shared_buffers + OS cache | Keep cache hit ratio > 99% |
| Replica count | Read QPS ÷ ~20–50K simple indexed reads/s per replica | 150K QPS ≈ 4–6 replicas with headroom |
Vertical Scaling Is Underrated#
Cloud instances reach 96–192 vCPU and 768 GB–2 TB+ RAM. If the working set fits in RAM, a single primary on one of these serves very large products. Vertical scaling is a config change and a failover (~1–2 min); sharding is a multi-quarter project. Exhaust vertical headroom with a 12-month runway plan before sharding — but start sharding work when the biggest available instance is < 2× your projected 18-month peak.
Read Scaling#
- Replicas: up to 5–15 per primary is common; each adds WAL-send load on the primary.
- Cache in front for hot keys (see Redis); Postgres's own
shared_buffers+ OS cache already make hot index lookups ~0.1 ms server-side — a cache mostly saves connections and CPU, not disk. - Materialized views for expensive aggregates;
REFRESH MATERIALIZED VIEW CONCURRENTLYon a schedule.
Write Scaling Before Sharding#
- Batch: multi-row
INSERT/COPY— 10–50× more rows/s than single-row commits. - Remove indexes nobody uses (
pg_stat_user_indexes.idx_scan = 0over 30 days). Each index is a write. - Move append-heavy, low-value data out (events, logs, metrics) to Kafka → warehouse or a time-series store.
- Split by domain (vertical partitioning): separate clusters for billing, catalog, messaging. Figma's public write-up describes doing this first before horizontally sharding.
- Then shard horizontally by the key most queries already filter on — usually
tenant_idoruser_id.
Multi-Region#
| Option | Writes | Reads | RPO / RTO | Pick When |
|---|---|---|---|---|
| Single region + cross-region async replica | One region | Local replica in each region | RPO = lag (seconds); RTO 5–30 min with promotion runbook | Default DR |
| Aurora Global Database | One region | Replicas in up to 5 secondary regions (typical lag < 1 s) | RPO ~1 s; managed failover ~1 min | AWS shops wanting managed DR |
| Regional cells (users homed to a region) | Each region writes its own users | Local | Per-cell | Data residency, most traffic region-local |
| Distributed SQL (Spanner, CockroachDB) | Any region | Local or leaseholder | RPO 0 | True global writes with strong consistency |
🎯 Staff Move: "For multi-region I'd home each tenant to a region and run an ordinary Postgres cell there. 95% of requests stay local, residency is solved for free, and I avoid paying consensus latency on every write for the 5% of cross-region traffic."
Failure Modes & Recovery#
1. Transaction ID Wraparound Shutdown#
Symptom: Warnings database "x" must be vacuumed within N transactions, then writes refused with a wraparound-protection error. Reads still work; the product is effectively down for writes.
Root cause: Anti-wraparound autovacuum couldn't freeze old tuples fast enough — blocked by a long transaction, an abandoned replication slot, or a prepared transaction; or throttled by default cost limits on a multi-TB table at 10K+ TPS.
Detection: age(datfrozenxid) per database and age(relfrozenxid) per table. Alert at 500M (≈ 25% of the limit), page at 1B.
Fix: Kill the xmin-holder (pg_stat_activity, pg_replication_slots, pg_prepared_xacts); run manual VACUUM (FREEZE, VERBOSE) on the oldest tables with raised maintenance_work_mem; in the worst case, single-user mode vacuum — hours to days of write downtime on large tables.
Prevention: Per-table autovacuum tuning for large hot tables (autovacuum_vacuum_cost_limit 2,000–10,000, autovacuum_freeze_max_age sized to churn); partitioning so frozen old partitions stop aging; alert on xmin age. Owner: DB platform / DBRE; service team owns the long transactions.
2. Lock-Queue Outage from a Migration#
Symptom: Error rate jumps to 100% on one table's endpoints within seconds of a deploy; connection count spikes to max_connections.
Root cause: DDL waiting for ACCESS EXCLUSIVE behind a long SELECT; all new queries queue behind the DDL; pooled connections exhaust.
Detection: pg_locks with granted = false > 20; pg_blocking_pids() tree rooted at a DDL statement; connection saturation.
Fix: Cancel the DDL (pg_cancel_backend) — queued queries drain in seconds.
Prevention: Migration framework enforces lock_timeout 2–5 s with retry, CONCURRENTLY for indexes, expand/contract for type changes; CI lints unsafe DDL. Owner: service team writes migrations; platform owns the linter and framework.
3. Replica Lag and Stale Reads#
Symptom: Users update settings and see old values; background jobs process a row "that doesn't exist yet." Replica lag grows to minutes during bulk loads.
Root cause: Replay is single-threaded per replica; a large batch update or index build on the primary produces WAL faster than one replay process applies it. Or replica queries conflict with replay and max_standby_streaming_delay lets lag grow.
Detection: pg_stat_replication.replay_lag, now() - pg_last_xact_replay_timestamp() on replicas; alert when p99 > lag budget (e.g., 5 s).
Fix: Route affected reads to primary; throttle the batch job (commit every 5–10K rows with pauses); LSN-aware routing.
Prevention: Batch jobs have a lag-aware throttle; replica routing checks lag and falls back to primary. Owner: service team for routing; data/batch team for throttling.
4. Bloat and Vacuum Debt#
Symptom: Table is 800 GB but live data is 150 GB; query latency drifts up over months; index-only scans stop being index-only; storage bill climbs.
Root cause: Update-heavy table with default autovacuum thresholds (scale_factor 0.2 = vacuum after 20% of rows change — 200M rows on a 1B-row table), long transactions holding the horizon, and HOT updates defeated by indexes on updated columns.
Detection: n_dead_tup / n_live_tup > 0.2; pgstattuple bloat estimates; last_autovacuum older than a day on hot tables.
Fix: pg_repack (online, needs ~2× table space) to reclaim; per-table autovacuum_vacuum_scale_factor = 0.01–0.05; drop indexes that break HOT.
Prevention: Autovacuum settings per table class in the standard; bloat dashboards. Owner: DBRE.
5. Failover That Loses Writes or Splits Brain#
Symptom: After primary failure, the promoted standby is missing the last N seconds of orders; or two nodes accept writes briefly.
Root cause: Async standby promoted with lag; or HA tooling without a proper consensus-backed leader lock and fencing, so the old primary keeps accepting writes after a network partition.
Detection: LSN comparison at promotion (logged by Patroni); application-level reconciliation; duplicate primary detection alerts.
Fix: Reconcile from upstream logs (payment provider, Kafka); rewind the old primary with pg_rewind.
Prevention: Patroni (etcd/Consul leader lease) or managed Multi-AZ; synchronous standby quorum for tier-0; watchdog / STONITH fencing. Owner: DB platform.
Operational Reality Matrix#
| Failure | Detection Signal | Blast Radius | Mitigation | Owner |
|---|---|---|---|---|
| XID wraparound | age(datfrozenxid) > 1B | All writes on the database | Kill xmin holder, freeze vacuum | DBRE |
| Migration lock queue | pg_locks waiters, blocking tree | One table → whole service | Cancel DDL; lock_timeout | Service team |
| Replica lag | replay_lag > budget | Stale reads | Route to primary, throttle batch | Service + batch owners |
| Bloat | Dead-tuple ratio > 20% | Latency, cost | pg_repack, per-table autovacuum | DBRE |
| Failover loss / split-brain | LSN gap at promotion | Recent writes | Sync standby, fencing | DB platform |
| Connection exhaustion | numbackends ≈ max | Whole cluster | Pooler, kill idle | Service team |
| Replication slot WAL pile-up | pg_replication_slots retained WAL > 50 GB | Primary disk full → outage | Drop/advance slot; cap WAL | Data platform |
| Plan regression | pg_stat_statements mean time 10× | One query → CPU for all | ANALYZE, extended stats, hint via index | Service team |
When to Use vs. Alternatives#
| Requirement | Pick | Why | Not Postgres Because |
|---|---|---|---|
| System of record with invariants | PostgreSQL | Constraints, ACID, flexible queries | — |
| Evolving access patterns, early product | PostgreSQL | Planner answers new questions | — |
| Multi-tenant SaaS < ~10 TB | PostgreSQL (+ Citus later) | RLS, tenant-leading keys | — |
| > 100K sustained writes/s, known access patterns | DynamoDB / Cassandra | Horizontal writes without sharding work | Single-writer ceiling |
| Multi-region active-active with strong consistency | Spanner / CockroachDB | Consensus replication | One writer region |
| Full-text relevance, facets, fuzzy search at scale | Elasticsearch | Inverted index, BM25, aggregations | tsvector is fine to ~10M docs, weak relevance tooling |
| Event log with replay | Kafka | Retention + offsets | Table-as-log bloats and doesn't fan out |
| Sub-ms hot-key reads at 500K QPS | Redis in front | RAM, no connection cost | Connection and CPU cost per query |
| Metrics / time-series at millions of points/s | Time Series DBs | Compression, downsampling | Row overhead ~24 bytes/tuple + indexes |
When NOT to use Postgres: as the only store for a firehose of append-only events at > 50K rows/s; as a message bus with many independent consumers; when geography demands multi-writer; when the team is going to run it unmanaged without anyone who understands vacuum.
Operational Concerns#
The Config Every Production Cluster Should Have#
shared_buffers 25% of RAM (8–64 GB typical)
effective_cache_size ~75% of RAM (planner hint only)
work_mem 16–64 MB (per sort/hash node, per query — multiply by concurrency)
maintenance_work_mem 1–4 GB (vacuum, index builds)
max_connections 300–500, behind PgBouncer transaction pooling
statement_timeout 30s for app roles; per-role overrides for batch
idle_in_transaction_session_timeout 60s
lock_timeout 5s default for migration roles
autovacuum_vacuum_scale_factor 0.02–0.05 on large hot tables (default 0.2 is too lazy)
autovacuum_max_workers 5–10; cost_limit raised to 2000+
checkpoint_timeout / max_wal_size 15 min / 16–64 GB on write-heavy systems
random_page_cost 1.1 on SSD/NVMe (default 4.0 assumes spinning disks)
log_min_duration_statement 500ms; enable pg_stat_statements everywhere
Backups and Recovery#
Physical base backups plus continuous WAL archiving (pgBackRest or WAL-G to object storage) give point-in-time recovery to any second in the retention window. The public 2017 GitLab incident — an accidental delete of the primary's data directory followed by discovering several backup mechanisms weren't working — is the canonical reminder: a backup you haven't restored is a hypothesis. Standard: automated weekly restore test to a scratch instance, measured RTO (a 2 TB restore + WAL replay commonly takes 2–6 hours), alert if the newest archived WAL is > 5 min old.
Upgrades#
- Minor versions: restart; with HA, fail over to an upgraded standby (~30–60 s disruption).
- Major versions:
pg_upgrade --link(minutes of downtime, but no easy rollback) or logical replication to a new-version cluster, then switch over (seconds of downtime, full rollback possible, but DDL and sequences need manual handling). Tier-0 databases use logical replication; the rest usepg_upgrade.
The Dashboard the On-Call Actually Uses#
| Metric | Healthy | Page |
|---|---|---|
| Primary CPU | < 60% | > 85% for 10 min |
| Active connections / pool size | < 70% | > 90% |
age(datfrozenxid) | < 200M | > 1B |
Replica replay_lag p99 | < 1 s | > lag budget (e.g., 10 s) |
| Longest transaction age | < 1 min | > 15 min (non-batch roles) |
| Lock waiters | 0–5 | > 50 |
| Buffer cache hit ratio | > 99% OLTP | < 95% |
| Replication slot retained WAL | < 5 GB | > 50 GB |
| Disk free | > 30% | < 15% |
What the On-Call Actually Does#
SELECT pid, state, now() - xact_start, wait_event, query FROM pg_stat_activity ORDER BY xact_start— find the oldest transaction and what's waiting.pg_blocking_pids()— find the root blocker; cancel before terminate.pg_stat_statementsordered by total time — find the query that changed.- Check replication slots and replica lag before blaming the primary.
- Never restart the primary to "clear it" — you lose the cache, trigger failover, and hide the evidence.
Interview Application — Staff-Level Plays#
Which Case Studies Use PostgreSQL#
| Case Study | How Postgres Is Used | Key Pattern |
|---|---|---|
| Payment Processing | Ledger and payment state machine | Double-entry ledger, UNIQUE idempotency keys, sync standby |
| Reservation Systems | Inventory holds and bookings | Conditional UPDATE, exclusion constraints on time ranges |
| Flash Sales | Authoritative order record behind a Redis gate | Bucketed inventory rows, outbox |
| Database Sharding | The thing being sharded | Logical shards over physical clusters, tenant keys |
| URL Shortener | Mapping store at modest scale | Sequence-based IDs, replicas + cache |
| Distributed Job Scheduler | Job table | FOR UPDATE SKIP LOCKED, partial indexes |
| Database Selection | The default to beat | Four ceilings framework |
| Maps & Geospatial | PostGIS for spatial queries | GiST indexes, geohash columns |
Every System Design Question Has a Postgres Moment#
- Chat: users, conversations, and memberships live in Postgres (relational, low volume); message bodies go to Cassandra once writes pass ~50K/s.
- Ticketing: an exclusion or unique constraint on
(event_id, seat_id)makes double-booking impossible — even if every cache and lock above it fails. - Notifications: preferences and templates in Postgres; the delivery pipeline in Kafka; the outbox connects them.
- Feed: the social graph starts in Postgres with an index on
(follower_id)and(followee_id); it moves out when fan-out reads exceed what replicas can serve.
L5 → L6 → L7 Responses#
| Scenario | Senior (L5) | Staff (L6) | Principal (L7) |
|---|---|---|---|
| "Pick a database" | "Postgres — it's reliable and supports transactions." | "Postgres until a named ceiling. We're at 3K write TPS with 3× annual growth; the ceiling is ~2 years out, so I pick tenant_id as the future shard key now and lead every PK with it." | "What's the org's paved road? If 80% of services should be on managed Postgres with standard configs, this one should be too — and we fund the sharding capability once, centrally, before the first team needs it." |
| "Reads are slow" | "Add read replicas and Redis." | "Check pg_stat_statements first — usually one query, a missing composite index, or OFFSET pagination. Then replicas with LSN-aware routing for read-your-writes." | "Slow-query regressions recur because nobody owns query review. I'd add plan-regression detection to CI against a prod-sized snapshot." |
| "Deploy a schema change" | "Run the migration in the deploy." | "lock_timeout 3 s with retry, CREATE INDEX CONCURRENTLY, expand/contract for type changes, batched backfill with lag-aware throttling." | "Migration outages are an org pattern. The migration framework enforces these rules; unsafe DDL fails CI. Teams can't opt out without a review." |
| "Primary fails" | "Fail over to a replica." | "Patroni with an etcd lease, sync standby quorum ANY 1 for tier-0 — zero acknowledged loss, ~30 s write unavailability. Tier-2 runs async and accepts seconds of loss." | "Which databases are tier-0, and when did we last test failover under load? I'd make quarterly failover game days a requirement and publish RPO/RTO per tier." |
| "We need to scale 10×" | "Shard the database." | "Vertical headroom first, then domain split, then shard by tenant into 1,024 logical shards on 8 physical clusters." | "Sharding costs ~2–4 engineer-quarters plus permanent complexity. Price it against a distributed SQL migration and against the biggest instance's runway — then pick the one-way door deliberately." |
Why "Pick a database" separates levels
"Postgres, it's reliable" is correct and competent — it's the answer most senior engineers give and it's usually right. The Staff answer turns the choice into a timeline: which ceiling, how far away, and what cheap decision today (tenant-leading keys) keeps the expensive decision (sharding) cheap later. The Principal answer recognizes that this choice is being made independently by 30 teams and that the org, not each team, should own the path past the ceiling.
Why "Deploy a schema change" separates levels
Every experienced engineer has run migrations; few have caused a lock-queue outage and learned why a harmless ALTER TABLE can take down a table. The Staff-level signal is naming the lock-queue mechanism and the specific guardrails. The Principal signal is noticing that guardrails in one engineer's head don't scale — they belong in the framework, enforced by CI.
The Staff PostgreSQL Checklist#
- Name the intent and the ceiling: "System of record, 3K write TPS, ceiling ~2 years out on write throughput."
- Push invariants into the schema: "Unique constraint on the idempotency key, check constraint on balance, FK with an index."
- Design the hot path's index: "Composite
(tenant_id, created_at DESC)covering the list endpoint — keyset pagination." - State isolation and durability per write class: "Read Committed with conditional updates; sync standby for the ledger,
synchronous_commit = offfor clickstream." - Plan retention and growth: "Monthly partitions, retention by
DROP PARTITION, autovacuum tuned per table." - Name the operational guardrails: "PgBouncer transaction pooling,
statement_timeout,idle_in_transaction_session_timeout,lock_timeouton migrations, xmin-age alert."
🎯 Staff Insight: What NOT to use Postgres for: a high-volume append-only firehose, a pub/sub bus with many consumers, or multi-region active-active writes. Saying "Postgres holds the state; Kafka carries the events; the outbox joins them" shows you know where its contract ends.
The Principal Lens#
Why L7 Sees This Problem Differently#
At Staff level, Postgres is a database you design well. At Principal level, it is the substrate most of the company's correctness depends on — and it fails in the same five ways in every team that runs it: blocking migrations, connection storms, vacuum debt, untested failover, and an unplanned sharding cliff. The L7 job is to convert those recurring failures into platform defaults and to decide, once, how the org gets past the single-writer ceiling — so that no individual team discovers it during a peak-season outage.
The Org-Level Fault Line#
Shared database vs database-per-service vs a managed Postgres platform.
| Option | What Works | What Breaks | Who Pays |
|---|---|---|---|
| One shared monolith database | Joins across domains, one thing to operate | Coupled schema changes, one team's query hurts all, one failover takes everything | Every team, on every incident |
| Database-per-service, self-run | Autonomy, blast-radius isolation | 60 clusters, 60 vacuum configs, 60 untested backups | Each team's on-call; DBRE knowledge is diluted |
| Managed Postgres platform with per-service clusters | Isolation + standard configs, backups, upgrades, migration tooling | Platform team is a dependency; needs a clear SLA | 3–6 platform engineers, repaid by incident reduction |
The Principal default: per-service clusters on a paved-road platform, with a central capability for the rare team that needs sharding — rather than each team building its own.
Cost Model#
Assumptions: managed Postgres on memory-optimized instances in a major US region, on-demand, approximate — 4 vCPU/32 GB ≈ $330/month, 32 vCPU/256 GB ≈ $2,700/month per node; storage ~$0.12/GB-month plus provisioned IOPS; backups ~$0.02–0.10/GB-month. Loaded engineer ~$25K/month.
| Scale | Footprint | Infra $/month | People | On-call Load | Dominant Risk |
|---|---|---|---|---|---|
| Startup — 1 app, 200 GB, 1K TPS | 4 vCPU primary + standby, 1 replica | ~$1–1.5K | ~0.1 FTE | Rare | Unsafe migrations, no restore test |
| Growth — 10 services, 5 TB, 15K TPS peak | 3–4 clusters of 32 vCPU primary + standby + 2 replicas | ~$30–45K | 1–2 DBREs | 2–4 pages/month | Vacuum debt, lock-queue outages, replica lag |
| Enterprise — 60 clusters + 1 sharded domain (64 logical → 16 physical) | ~200 nodes total | ~$250–400K | 4–6 platform/DBRE + sharding team of 3–4 | 10–20 pages/month fleet-wide | Correlated failure: bad default pushed fleet-wide, major-version upgrade debt, sharding cliff |
Enterprise levers: reserved instances (30–50% off), right-sizing replicas nobody reads (commonly 20–30% of replicas), moving append-only data to object storage + warehouse. A migration from sharded Postgres to distributed SQL typically costs 4–8 engineer-quarters — price it before choosing.
The 3-Year Evolution Path#
One-Way Doors vs Two-Way Doors#
| Decision | Door | Reversal Cost | Why |
|---|---|---|---|
Config knobs (autovacuum, work_mem, timeouts) | Two-way | Minutes | Reload-only for most |
| Adding a read replica / cache | Two-way | Days | Remove and re-route |
| Primary key type (UUIDv4 vs sequential vs UUIDv7) | One-way at scale | Table rewrite + every FK: months | Every referencing table and client encodes it |
| Shard key | One-way | Full data re-distribution: quarters | All queries and routing assume it |
| Tenancy model (shared vs schema vs DB per tenant) | One-way | Re-platform all tenants | Isolation promises in contracts |
| Aurora-specific or vendor features | Mostly one-way | Re-test performance on vanilla; lose storage behavior | Lock-in to storage engine semantics |
| Extensions (PostGIS, pgvector, Citus) | One-way-ish | Rewrite feature on another store | Data and queries depend on them |
The Standard I'd Write#
RFC: Relational Database Standard (v1)
Scope: All PostgreSQL databases serving production traffic.
MUST
- Run on the managed platform with PITR enabled and an automated restore test at least weekly; measured RTO is published per database.
- Connect through the platform pooler; app roles have
statement_timeout≤ 30 s andidle_in_transaction_session_timeout≤ 60 s.- Ship schema changes through the migration framework:
lock_timeoutenforced,CONCURRENTLYfor indexes, unsafe DDL fails CI.- Declare a tier; tier-0 uses a synchronous standby quorum and passes a quarterly failover game day under load.
- Lead every primary key on multi-tenant tables with
tenant_id.SHOULD
- Use UUIDv7 or sequential keys, never UUIDv4, for tables expected to exceed 100M rows.
- Send analytics to the warehouse via CDC, not to the primary.
- Request the central sharding capability when projected 18-month peak exceeds 50% of the largest instance.
Exceptions: Reviewed by the database platform team; time-boxed to 2 quarters; tier-0 exceptions need VP approval.
Success metrics: Migration-caused incidents → 0 within 2 quarters; 100% of databases with a passing restore test in the last 7 days; xmin-age pages → 0; fleet cost per transaction down 15% year over year.
What I'd Tell the VP#
Nearly every product we have depends on Postgres, and last year about a third of our database incidents came from the same few mistakes repeated in different teams — unsafe schema changes and untested backups being the worst. I'm proposing we make the safe way the default: a managed database platform that enforces safe migrations, tests restores weekly, and runs failover drills for critical systems. It needs four engineers and pays for itself in avoided outages and roughly 20% infrastructure savings from right-sizing. Separately, our fastest-growing product will outgrow a single database in about 18 months; I want to start that scaling work now as a shared capability rather than as an emergency.
Principal Interview Signals#
| Signal | What It Sounds Like |
|---|---|
| Ceilings as a timeline | "We have ~20 months of write headroom. The one-way door I'll open now is the shard key; the expensive work waits." |
| Prices migrations | "Sharding is 3 engineer-quarters plus a permanent routing layer. Distributed SQL is 6 quarters and 5× write latency. I'd pay for sharding here because 99% of queries are tenant-scoped." |
| Converts incidents into defaults | "Two teams took lock-queue outages this year. The fix is a migration linter, not a wiki page." |
| Correlated failure awareness | "A bad autovacuum default pushed fleet-wide is our largest blast radius. Config changes roll out by tier, tier-2 first." |
| Knows when not to centralize | "I wouldn't consolidate into one shared database for cost. Isolation is worth the extra clusters." |
Staff answers that L7 interviewers find insufficient:
- "We'll use Patroni with a sync standby." — Right for one cluster; silent on who tests failover across 60.
- "We'll shard by tenant_id when we need to." — Correct key, but no trigger, no owner, no cost, no plan for the teams that will each try to build it themselves.
- "Postgres scales fine with replicas." — True for reads; doesn't price the write ceiling or the org's exposure to it.
🧭 Principal Move: "Before we design this database, I'd check what the platform already guarantees. If safe migrations, PITR, and failover drills aren't defaults, that's the design work with the most leverage — this schema is the easy part."
In the Wild#
Instagram — Sharded Postgres with Embedded IDs#
Instagram publicly described sharding Postgres into thousands of logical shards mapped onto far fewer physical servers, with IDs generated inside Postgres (PL/pgSQL) that pack 41 bits of timestamp, 13 bits of logical shard ID, and 10 bits of per-shard sequence into a 64-bit integer. IDs are time-sortable and self-routing.
Staff insight: Logical-over-physical sharding plus shard-in-the-ID makes rebalancing a matter of moving logical shards — the technique to name when an interviewer asks how you'd reshard without rewriting keys.
Notion — 480 Logical Shards#
Notion's engineering blog describes migrating its monolithic Postgres to 480 logical shards across 32 physical databases, partitioned by workspace ID, with a double-write and backfill migration plan and careful verification before cutover. They later described further splitting physical databases as load grew.
Staff insight: The shard key matched the product's natural tenant boundary (workspace), so almost every query stayed single-shard — the precondition that makes application-level sharding worth its cost.
Figma — Vertical Then Horizontal#
Figma has publicly written about scaling its Postgres stack first by vertical partitioning (moving groups of tables into separate databases) and later building horizontal sharding with a query-routing proxy layer, rolled out incrementally table by table.
Staff insight: Vertical split first buys years cheaply; horizontal sharding is reserved for the domains that still outgrow one instance. That ordering is the Staff answer to "how do you scale Postgres."
Practice Drill#
Prompt: "Your multi-tenant SaaS runs on one Postgres primary (64 vCPU). Primary CPU is at 75% at peak, growing 8% per month. The biggest tables are
events(3 TB, append-mostly) anddocuments(800 GB, update-heavy). Leadership asks whether to migrate to a distributed database. What do you recommend?"
Staff Answer
Not yet — first buy runway, then shard by tenant on our own schedule. At 8% monthly growth, 75% CPU reaches ~100% in about 4 months, so action is urgent but migration to distributed SQL (5–8 engineer-quarters, higher write latency) doesn't land in time anyway. Step 1 (weeks): pg_stat_statements top-10 by total time — typically 2–3 queries are 50%+ of CPU; fix indexes and pagination. Move events out: it's append-mostly, so stream it via CDC/Kafka to a warehouse or time-series store and keep only 30 days in monthly partitions — this likely removes 20–30% of write load and most vacuum pressure. Tune documents for HOT updates (fillfactor = 85, drop indexes on frequently updated columns, per-table autovacuum at 2%). Step 2 (1–2 months): scale vertically to the next instance size — a failover of ~1 minute buys ~12 months at current growth. Step 3 (2–3 quarters): tenant-based sharding — 1,024 logical shards on 4 physical clusters, routing in a thin data-access layer, largest tenants isolated first. Metrics: primary CPU p95, WAL MB/s, dead-tuple ratio on documents, share of cross-tenant queries (must be < 1% to shard by tenant). Owners: service team for query fixes; data platform for the events move; DB platform for sharding.
Why this is L6:
- Converts growth rate into a deadline and sequences work to beat it.
- Separates the two tables by workload shape rather than treating "the database" as one problem.
- Picks a shard key from query evidence and names the metric that validates it.
What L7 adds:
- Builds the sharding layer as a reusable platform capability, since the next two products are on the same curve.
- Prices the options for leadership: sharding (~3 quarters, ~4 engineers) vs distributed SQL (~6–8 quarters, plus 2–5× write latency) vs vertical-only (runway ends in ~16 months).
- Uses the moment to set an org rule: any table projected past 1 TB needs a partitioning and retention plan at design review.
Quick Reference Card#
Write ceiling: ~10–50K OLTP write TPS per primary; start sharding plan at 50% of biggest instance
Connections: process per connection (~5–10 MB); 2–4 active per core; PgBouncer transaction mode
MVCC: UPDATE = new tuple; vacuum reclaims; long txns block vacuum database-wide
Wraparound: 32-bit XIDs, ~2B limit; alert age(datfrozenxid) > 500M, page > 1B
Isolation: Read Committed default; conditional UPDATE beats raising isolation; SSI aborts need retry
Durability: synchronous_commit per txn; sync standby ANY 1 of 2 = zero-loss failover, +1 RTT
Replication: physical streaming (HA/replicas), logical (CDC, upgrades); watch slot WAL retention
Indexes: B-tree default; partial for queues; GIN for JSONB/FTS; BRIN for huge time-ordered
Keys: BIGINT identity or UUIDv7; avoid UUIDv4 on 100M+ row tables
Migrations: lock_timeout 3–5s + retry; CREATE INDEX CONCURRENTLY; expand/contract
Retention: partition by time, DROP PARTITION — never mass DELETE
Timeouts: statement_timeout 30s, idle_in_transaction_session_timeout 60s
Failover: Patroni / Multi-AZ 30–120 s; test quarterly under load
Red flags: analytics on primary, OFFSET pagination, unindexed FKs, max_connections 5000,
untested backups, dual-writing Postgres + Kafka without an outbox