Hiring BarSupport

Design with PostgreSQL — Staff-Level Technology Guide

Technology guide41 min read5 diagrams

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.

IntentConstraintPostgres StrategyFailure ModeCorrectness Bar
System of record (OLTP)Invariants must holdNormalized schema, constraints, Read Committed + targeted SELECT ... FOR UPDATELock contention on hot rowsZero lost or double-applied writes
Read-heavy product surfacep99 < 50 ms at 50K QPSReplicas + covering indexes + cacheReplica lag shows stale dataBounded staleness, read-your-writes where needed
Multi-tenant SaaSNoisy tenants, isolationtenant_id in every PK, RLS, later Citus / app sharding by tenantOne tenant's query plan starves othersTenant data never crosses
Analytics / reportingBig scansReplica or CDC to warehouse — not the primaryLong queries block vacuum on primaryApproximate 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 ModeServer Connection Held ForMultiplexing RatioBreaksPick When
SessionThe whole client session~1:1NothingLegacy apps; admin tools
TransactionOne transaction10–50:1Session SET, cross-txn advisory locks, LISTEN/NOTIFY, older prepared-statement handlingDefault for web/API tiers
StatementOne statementHighestMulti-statement transactionsRare — 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 VACUUM marks 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 transaction session 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–90 on 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.
Diagram: MVCC, Dead Tuples, and Vacuum

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#

ModeMechanismLagUse
Streaming (physical)Ships WAL bytes; replica is a byte-identical copyms–seconds; spikes under heavy writes or replica loadHA standbys, read replicas
Synchronous streamingCommit 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 / decodingDecodes WAL into row changes per table (publications/slots)SecondsCDC 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.

Diagram: Replication

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#

IndexGood ForCost / Caveat
B-tree (default)Equality, range, ORDER BY, uniquenessEvery index adds ~1 extra write per insert/update of that column
Composite (a, b)Queries filtering a then b or sorting by b within aColumn order matters: leftmost prefix rule
Covering INCLUDE (c)Index-only scans, no heap fetchBigger index; needs visibility map (vacuum) to be effective
Partial WHERE status = 'pending'Queues, soft-deletes, hot subsetsTiny index for a tiny hot set — often 100× smaller
GINJSONB containment, arrays, full-text (tsvector)Slower writes; pending-list flushes cause latency spikes
GiST / SP-GiSTGeospatial (PostGIS), ranges, exclusion constraintsLarger, slower than B-tree for equality
BRINHuge 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#

KeyProsConsPick When
BIGINT IDENTITY8 bytes, sequential, B-tree friendlyLeaks volume; needs coordination across shardsSingle-cluster tables, internal IDs
UUIDv4Globally unique, no coordination16 bytes, random inserts → page splits, cache misses, WAL full-page bloatOnly when IDs must be unguessable and table is small
UUIDv7 / ULIDGlobally unique and time-ordered16 bytesDefault for new distributed-friendly schemas
Snowflake-style 64-bitTime-ordered, embeds shard ID, 8 bytesNeeds an ID-generation schemePre-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#

LevelWhat You GetAnomalies Still PossibleCostPick When
Read Committed (default)Each statement sees data committed before it startedNon-repeatable reads, lost updates on read-modify-write, write skewCheapest90% 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 40001Write skew (two txns read overlapping sets, write disjoint rows)Retries on conflictReports needing a consistent view; multi-statement reads
Serializable (SSI)Equivalent to some serial orderNonePredicate-lock tracking; 1–10% abort rate under contention, must retryInvariants 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)#

SettingCommit Returns AfterLoss on Primary CrashLoss on FailoverLatency Added
offWAL in memory; flushed within ~3× wal_writer_delay (≈ 600 ms)Up to ~600 ms of commits (no corruption)Replication lagNone — often 2–5× write throughput on slow disks
localLocal WAL flushed0Replication lagLocal fsync (0.1–3 ms)
on (default) + no sync standbyLocal WAL flushed0Replication lagLocal fsync
remote_write + sync standbyStandby received WAL (OS buffer)0~0 unless both nodes crash together+1 RTT
on + sync standbyStandby flushed WAL00+1 RTT + standby fsync
remote_applyStandby applied WAL — visible on standby reads00Highest; 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 with synchronous_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-PatternDetection SignalBlast RadiusWho Pays
Blocking DDLpg_locks waiters > 50; lock wait graph rooted at DDLEvery query on the tableCustomer-facing services, on-call
Connection stormnumbackends > 3× cores; CPU in LWLock waitsWhole databaseAll tenants of the cluster
Long transactionsmax(now() - xact_start) > 5 min; n_dead_tup climbingDatabase-wide bloatFuture query latency, storage bill
OFFSET paginationpg_stat_statements mean time scales with pageOne endpoint, heavy I/OPrimary CPU
Analytics on primaryBuffer hit ratio drops < 95% during reportsOLTP latencyProduct SLOs
DELETE retentionWAL rate spikes nightly, replica lag > 30 sReplicas and vacuumRead path, backups
Unindexed FKSeq scans on child tables during parent deleteLock contentionWriters of both tables

The Technology Landscape / Head-to-Head Comparison#

DimensionPostgreSQL (self/RDS)Aurora PostgreSQLMySQL / InnoDBCitusCockroachDB / YugabyteDB / SpannerDynamoDB
Write scalingOne primary; ~10–50K TPS typical OLTPOne writer; storage scales to 128 TBOne primary; similarSharded across workersHorizontally scalable writesUnlimited with good keys
TransactionsFull ACID, SSISameACID, weaker default isolation semanticsSingle-shard fast; cross-shard 2PCDistributed serializableUp to 100 items per TransactWriteItems
Latency (point write)1–3 ms2–5 ms (6-way quorum storage)1–3 ms2–10 ms5–20 ms in-region (consensus)5–10 ms
FailoverPatroni / Multi-AZ: 30–120 s~30 s typicalSimilar to PostgresPer-node HAAutomatic, secondsTransparent
Replica lagms–sTypically < 20–100 ms (shared storage)ms–sPer workerConsensus — no lag for leaseholder readsEventually consistent reads < 1 s
EcosystemRichest: PostGIS, pgvector, extensionsMost extensionsHuge, simplerPostgres extensionPostgres wire-compatible, subsetAWS-native
Pick whenDefault system of recordWant managed + faster replicas, accept AWS lock-in and I/O pricingTeam expertise, Vitess sharding pathMulti-tenant SaaS outgrowing one nodeMulti-region writes with strong consistencyKnown 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#

ModelIsolationTenant Count CeilingOps CostPick When
Shared tables, tenant_id column + RLSLogicalMillionsLowestDefault SaaS
Schema per tenantNamespace~1–5K per cluster (catalog bloat, migration fan-out)Migrations × NRegulated mid-market, per-tenant customization
Database per tenantStrongHundreds per clusterHighestEnterprise contracts demanding isolation
Cell per group of tenantsBlast radiusUnlimited via more cellsFleet tooling neededLarge 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.

PatternContractCeiling SignalNext StepOwner
ReplicasBounded stalenessReplica lag p99 > budgetCache, or remove replica-heavy queryService team
PartitioningCheap retention> 500 partitions or cross-partition queries dominateTime-series DBService team
RLS multi-tenantLogical isolationOne tenant > 20% of loadMove tenant to own cellPlatform + account team
Outbox + CDCAt-least-once eventsSlot lag > 10 minScale consumers; cap slot WALData platform
App shardingHorizontal writesCross-shard queries > 5%Re-pick shard key (painful)Platform
Everything-storeGood enoughFeature's p99 or CPU share > 20%Dedicated systemService team

Scaling#

The Four Ceilings#

CeilingTypical Threshold (single primary, modern hardware)Early WarningFirst Response
Write TPS~10–50K simple OLTP write txns/s; less with many indexesPrimary CPU > 60% sustained; WAL generation > 100–200 MB/sBatch writes, drop unused indexes, move low-value writes out
Table sizeSingle tables > 1–2 TB get painful: vacuum takes hours, index builds take hoursAutovacuum runs never finish; n_dead_tup monotonicPartition; archive cold data
Connections~300–500 active backendsLWLock waits, CPU in kernelPooler; fewer, busier connections
GeographyOne writer regionRemote-region write p99 > 100 msRegional read replicas; then distributed SQL or cells

Back-of-Envelope Sizing#

InputRule of ThumbExample
Row size on diskPayload + ~24-byte tuple header + ~4-byte item pointer, rounded to 8 bytes200-byte payload ≈ 230 bytes
Index size~key width + 8–16 bytes per entry, ~70–90% page fill1B rows × BIGINT index ≈ 25–30 GB
Table growthrows/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 RAMHot indexes + hot rows should fit shared_buffers + OS cacheKeep cache hit ratio > 99%
Replica countRead QPS ÷ ~20–50K simple indexed reads/s per replica150K 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 CONCURRENTLY on a schedule.

Write Scaling Before Sharding#

  1. Batch: multi-row INSERT / COPY — 10–50× more rows/s than single-row commits.
  2. Remove indexes nobody uses (pg_stat_user_indexes.idx_scan = 0 over 30 days). Each index is a write.
  3. Move append-heavy, low-value data out (events, logs, metrics) to Kafka → warehouse or a time-series store.
  4. Split by domain (vertical partitioning): separate clusters for billing, catalog, messaging. Figma's public write-up describes doing this first before horizontally sharding.
  5. Then shard horizontally by the key most queries already filter on — usually tenant_id or user_id.
Diagram: Write Scaling Before Sharding

Multi-Region#

OptionWritesReadsRPO / RTOPick When
Single region + cross-region async replicaOne regionLocal replica in each regionRPO = lag (seconds); RTO 5–30 min with promotion runbookDefault DR
Aurora Global DatabaseOne regionReplicas in up to 5 secondary regions (typical lag < 1 s)RPO ~1 s; managed failover ~1 minAWS shops wanting managed DR
Regional cells (users homed to a region)Each region writes its own usersLocalPer-cellData residency, most traffic region-local
Distributed SQL (Spanner, CockroachDB)Any regionLocal or leaseholderRPO 0True 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#

FailureDetection SignalBlast RadiusMitigationOwner
XID wraparoundage(datfrozenxid) > 1BAll writes on the databaseKill xmin holder, freeze vacuumDBRE
Migration lock queuepg_locks waiters, blocking treeOne table → whole serviceCancel DDL; lock_timeoutService team
Replica lagreplay_lag > budgetStale readsRoute to primary, throttle batchService + batch owners
BloatDead-tuple ratio > 20%Latency, costpg_repack, per-table autovacuumDBRE
Failover loss / split-brainLSN gap at promotionRecent writesSync standby, fencingDB platform
Connection exhaustionnumbackends ≈ maxWhole clusterPooler, kill idleService team
Replication slot WAL pile-uppg_replication_slots retained WAL > 50 GBPrimary disk full → outageDrop/advance slot; cap WALData platform
Plan regressionpg_stat_statements mean time 10×One query → CPU for allANALYZE, extended stats, hint via indexService team
Diagram: Operational Reality Matrix

When to Use vs. Alternatives#

RequirementPickWhyNot Postgres Because
System of record with invariantsPostgreSQLConstraints, ACID, flexible queries—
Evolving access patterns, early productPostgreSQLPlanner answers new questions—
Multi-tenant SaaS < ~10 TBPostgreSQL (+ Citus later)RLS, tenant-leading keys—
> 100K sustained writes/s, known access patternsDynamoDB / CassandraHorizontal writes without sharding workSingle-writer ceiling
Multi-region active-active with strong consistencySpanner / CockroachDBConsensus replicationOne writer region
Full-text relevance, facets, fuzzy search at scaleElasticsearchInverted index, BM25, aggregationstsvector is fine to ~10M docs, weak relevance tooling
Event log with replayKafkaRetention + offsetsTable-as-log bloats and doesn't fan out
Sub-ms hot-key reads at 500K QPSRedis in frontRAM, no connection costConnection and CPU cost per query
Metrics / time-series at millions of points/sTime Series DBsCompression, downsamplingRow 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 use pg_upgrade.

The Dashboard the On-Call Actually Uses#

MetricHealthyPage
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 waiters0–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#

  1. 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.
  2. pg_blocking_pids() — find the root blocker; cancel before terminate.
  3. pg_stat_statements ordered by total time — find the query that changed.
  4. Check replication slots and replica lag before blaming the primary.
  5. 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 StudyHow Postgres Is UsedKey Pattern
Payment ProcessingLedger and payment state machineDouble-entry ledger, UNIQUE idempotency keys, sync standby
Reservation SystemsInventory holds and bookingsConditional UPDATE, exclusion constraints on time ranges
Flash SalesAuthoritative order record behind a Redis gateBucketed inventory rows, outbox
Database ShardingThe thing being shardedLogical shards over physical clusters, tenant keys
URL ShortenerMapping store at modest scaleSequence-based IDs, replicas + cache
Distributed Job SchedulerJob tableFOR UPDATE SKIP LOCKED, partial indexes
Database SelectionThe default to beatFour ceilings framework
Maps & GeospatialPostGIS for spatial queriesGiST 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#

ScenarioSenior (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#

  1. Name the intent and the ceiling: "System of record, 3K write TPS, ceiling ~2 years out on write throughput."
  2. Push invariants into the schema: "Unique constraint on the idempotency key, check constraint on balance, FK with an index."
  3. Design the hot path's index: "Composite (tenant_id, created_at DESC) covering the list endpoint — keyset pagination."
  4. State isolation and durability per write class: "Read Committed with conditional updates; sync standby for the ledger, synchronous_commit = off for clickstream."
  5. Plan retention and growth: "Monthly partitions, retention by DROP PARTITION, autovacuum tuned per table."
  6. Name the operational guardrails: "PgBouncer transaction pooling, statement_timeout, idle_in_transaction_session_timeout, lock_timeout on 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.

OptionWhat WorksWhat BreaksWho Pays
One shared monolith databaseJoins across domains, one thing to operateCoupled schema changes, one team's query hurts all, one failover takes everythingEvery team, on every incident
Database-per-service, self-runAutonomy, blast-radius isolation60 clusters, 60 vacuum configs, 60 untested backupsEach team's on-call; DBRE knowledge is diluted
Managed Postgres platform with per-service clustersIsolation + standard configs, backups, upgrades, migration toolingPlatform team is a dependency; needs a clear SLA3–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.

ScaleFootprintInfra $/monthPeopleOn-call LoadDominant Risk
Startup — 1 app, 200 GB, 1K TPS4 vCPU primary + standby, 1 replica~$1–1.5K~0.1 FTERareUnsafe migrations, no restore test
Growth — 10 services, 5 TB, 15K TPS peak3–4 clusters of 32 vCPU primary + standby + 2 replicas~$30–45K1–2 DBREs2–4 pages/monthVacuum debt, lock-queue outages, replica lag
Enterprise — 60 clusters + 1 sharded domain (64 logical → 16 physical)~200 nodes total~$250–400K4–6 platform/DBRE + sharding team of 3–410–20 pages/month fleet-wideCorrelated 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#

Diagram: The 3-Year Evolution Path

One-Way Doors vs Two-Way Doors#

DecisionDoorReversal CostWhy
Config knobs (autovacuum, work_mem, timeouts)Two-wayMinutesReload-only for most
Adding a read replica / cacheTwo-wayDaysRemove and re-route
Primary key type (UUIDv4 vs sequential vs UUIDv7)One-way at scaleTable rewrite + every FK: monthsEvery referencing table and client encodes it
Shard keyOne-wayFull data re-distribution: quartersAll queries and routing assume it
Tenancy model (shared vs schema vs DB per tenant)One-wayRe-platform all tenantsIsolation promises in contracts
Aurora-specific or vendor featuresMostly one-wayRe-test performance on vanilla; lose storage behaviorLock-in to storage engine semantics
Extensions (PostGIS, pgvector, Citus)One-way-ishRewrite feature on another storeData 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 and idle_in_transaction_session_timeout ≤ 60 s.
  • Ship schema changes through the migration framework: lock_timeout enforced, CONCURRENTLY for 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#

SignalWhat 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) and documents (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
  1. Loading the index…