Hiring BarSupport

Database Indexing

Foundation31 min read4 diagrams

Why This Matters#

Indexing is not a query-optimization question. It is a question about which workload you have decided to make expensive. Every index makes some reads cheaper and every write more expensive — in CPU, in I/O, in storage, and in replication lag. An index is a standing tax paid on every insert so that some query can skip a scan. The Staff question is always: who pays that tax, and is the query worth it?

Most candidates treat indexes as free: "add an index on user_id" and move on. That answer is usually correct and always incomplete. At 50K writes/sec, each additional secondary index adds a random write, a WAL record, and replica apply work. Seven indexes on a hot table can turn a 1 KB logical insert into 10–20 KB of physical I/O. Under the wrong storage engine, that becomes the bottleneck that forces you to shard a year early.

The Staff-level understanding has three layers: the data structure (B-tree vs LSM, and why that choice decides read vs write amplification), the index shape (composite column order, covering, partial — the difference between a 2ms query and a 2s one), and the distribution problem (a secondary index in a sharded database is either local and scatter-gathered or global and asynchronously maintained — there is no third option that is free).

The 60-Second Version#

  • An index is a sorted, redundant copy of some columns plus a pointer. It turns an O(n) scan into an O(log n) seek, and turns every write into 1 + (number of indexes) writes.
  • B-trees optimize reads; LSM trees optimize writes. B-tree: ~3–4 page reads per lookup, in-place updates, write amplification ~10–30× under random writes. LSM: sequential appends and compaction, write amplification ~10–30× from compaction, read amplification across levels mitigated by bloom filters (~1% false positive at 10 bits/key).
  • Composite index column order is the whole game. Equality columns first, then the range/sort column. (tenant_id, status, created_at) serves WHERE tenant_id=? AND status=? ORDER BY created_at; (created_at, tenant_id) does not.
  • Covering indexes eliminate the table lookup. A query that reads only indexed columns skips the heap fetch — often 5–50× faster for range scans returning thousands of rows.
  • Partial indexes index only the rows you query. WHERE status = 'PENDING' on a table that is 99% COMPLETED gives an index 100× smaller.
  • In a sharded database, secondary indexes are local (scatter-gather reads) or global (async, eventually consistent writes). DynamoDB GSIs, Cassandra's local secondary indexes, and Spanner's distributed indexes are each a different point on that tradeoff.
  • Unused indexes are pure cost. Audit them. In most mature schemas, 20–30% of indexes are never read.

How Database Indexing Works#

The Basic Idea#

Without an index, SELECT * FROM orders WHERE customer_id = 42 reads every row. On a 100M-row, 50 GB table at ~500 MB/s sequential throughput, that is ~100 seconds — and it evicts your hot working set from the buffer pool while it does it.

An index on customer_id is a separate structure, sorted by customer_id, where each entry points to the row's location (a physical tuple ID in Postgres, the primary key in InnoDB). The lookup descends a tree of ~3–4 levels to the first matching entry, reads the matching entries sequentially, then fetches each row. For 20 matching rows: ~4 index page reads + up to 20 heap page reads ≈ 1–5ms with a warm cache.

The cost: every INSERT, every DELETE, and every UPDATE that touches customer_id must also modify the index — and do it in the same transaction.

Key Terms#

TermMeaningWhy It Matters
Clustered indexTable rows stored physically in index order (InnoDB primary key, SQL Server clustered)Range scans on the key are sequential; secondary indexes point to the PK, adding a second lookup
Heap + secondary indexRows stored unordered; all indexes point to physical location (Postgres)Updates that move a row can require updating every index (mitigated by HOT updates)
SelectivityFraction of rows matching a predicateIndex on is_active (90% true) is useless; the planner will scan
CardinalityNumber of distinct valuesHigh cardinality → selective → index-friendly
Covering indexIndex containing every column the query needsIndex-only scan, no heap fetch
Write amplificationPhysical bytes written ÷ logical bytes writtenDetermines disk wear, IOPS budget, and max write throughput
Read amplificationPhysical reads per logical readLSM point reads may touch several SSTables without bloom filters
Space amplificationPhysical size ÷ logical data sizeLSM with stale versions pending compaction; B-tree page fill factor

Where Indexes Live#

Engine FamilyIndex StructureExamplesSweet Spot
Relational OLTPB+tree (clustered or heap)PostgreSQL, MySQL/InnoDB, SQL Server, OracleRead-heavy, mixed queries, secondary indexes, transactions
LSM key-value / wide-columnLSM tree + bloom filtersCassandra, RocksDB, ScyllaDB, HBase, (DynamoDB's internals are not public — treat as managed)Write-heavy, key-based access, time series
NewSQLLSM (RocksDB/Pebble) under distributed SQLCockroachDB, TiDB, YugabyteDBSQL + horizontal scale; secondary indexes are distributed
SearchInverted indexElasticsearch, OpenSearch, LuceneFull-text, faceted, "any field" queries
SpecializedGiST/R-tree, GIN, BRIN, HNSWPostGIS, Postgres JSONB, vector indexesGeospatial, containment, append-ordered, similarity

Core Strategies#

Strategy 1: B-Tree (B+Tree) Indexes#

A balanced tree of fixed-size pages (8 KB in Postgres, 16 KB in InnoDB). Internal pages hold separator keys; leaf pages hold entries and are linked for range scans. With ~200–500 keys per page, a 4-level tree addresses billions of rows.

lookup(key):
    page = root                               # almost always cached
    while page is internal:
        page = page.child_for(key)            # binary search within page
    return page.find(key)                     # leaf; follow sibling links for ranges

insert(key, ptr):
    leaf = descend(key)
    if leaf.has_space: leaf.insert(key, ptr)  # rewrite whole 8-16 KB page for one entry
    else: split(leaf)                          # propagates up; random I/O

When to use: read-heavy or mixed workloads needing point lookups, range scans, and ordering on many different columns. The default for OLTP.

Failure mode: random-write amplification. Inserting a 100-byte entry dirties an entire 8–16 KB page; with random keys (UUIDv4), every insert touches a different leaf, the working set of dirty pages exceeds memory, and throughput collapses to disk IOPS. Page splits also fragment the index (fill factor drops toward ~50–70%). Mitigation: time-ordered keys (UUIDv7, ULID, snowflake IDs) so inserts hit the rightmost leaf.

Strategy 2: LSM-Tree Indexes#

Writes go to an in-memory sorted buffer (memtable) and a sequential commit log. When the memtable fills (~64–256 MB), it is flushed as an immutable sorted file (SSTable). Background compaction merges SSTables to discard overwritten and deleted entries.

write(key, value):
    commit_log.append(key, value)             # sequential, fsync-batched
    memtable.put(key, value)                  # O(log n) in memory
    if memtable.size > threshold: flush_to_sstable()

read(key):
    if key in memtable: return memtable[key]
    for sstable in newest_to_oldest:          # L0, L1, ... Ln
        if sstable.bloom.might_contain(key):  # ~1% false positive at 10 bits/key
            v = sstable.lookup(key)
            if v: return v                    # may be a tombstone
    return null

When to use: write-heavy workloads (> ~10K writes/s per node), append-mostly data (events, metrics, messages), and key-based access. Sequential writes make LSM friendly to both SSDs and spinning disks.

Failure mode: compaction debt. If writes outpace compaction, SSTable count grows, reads touch more files, p99 latency climbs, and disk usage can temporarily double during a major compaction. Tombstone accumulation: deletes are writes; a range scan over 100K tombstones to find 10 live rows is the classic Cassandra latency incident.

B-Tree vs LSM: The Core Tradeoff#

DimensionB-TreeLSM
Point read~3–4 page reads, predictableMemtable + bloom checks per level; 1–2 disk reads typical, worse under compaction debt
Range scanExcellent — linked leavesGood, but merges across levels; tombstones hurt
Write throughputBounded by random IOPS for random keysBounded by sequential bandwidth and compaction
Write amplification~10–30× under random writes (page rewrites + WAL + full-page writes)~10–30× leveled compaction; ~4–10× size-tiered
Space amplificationFragmentation, fill factor ~70%Size-tiered can need ~2× during compaction; leveled ~1.1×
Latency predictabilityHighCompaction causes p99 spikes unless throttled
ChoiceWhat WorksWhat BreaksWho Pays
B-tree (Postgres/MySQL)Rich secondary indexes, predictable reads, transactionsRandom-write ceiling; bloat and vacuum on churn-heavy tablesDBAs/on-call pay in vacuum tuning; product pays when write growth forces early sharding
LSM (Cassandra/RocksDB)5–10× higher sustained write throughput per nodeSecondary indexes weak; tombstones; compaction p99 spikesApplication teams pay in query-first data modeling; on-call pays in compaction and repair tuning

🎯 Staff Move: "The write-to-read ratio picks the engine. At 2K writes/sec with ad-hoc queries on six columns, Postgres with B-tree indexes. At 200K writes/sec of append-only events read by key and time range, an LSM store where the primary key is the index. I'd rather design around one access path than bolt secondary indexes onto an LSM store."

Index Shape: Composite, Covering, and Partial#

Choosing the structure is the storage engine's job. Choosing the shape is yours — and it is where 90% of real indexing wins and failures come from.

Composite Index Column Order#

A composite index on (a, b, c) is sorted by a, then b within a, then c within b. It can serve predicates on a leftmost prefix: a, (a, b), (a, b, c). It cannot efficiently serve b alone.

The rule: Equality → Sort → Range (ESR). Put columns filtered by = first, then the column you ORDER BY, then range predicates.

-- Query: a tenant's open tickets, newest first
SELECT id, subject, created_at FROM tickets
WHERE tenant_id = $1 AND status = 'OPEN'
ORDER BY created_at DESC LIMIT 50;

-- Good: seeks to (tenant, OPEN), reads 50 entries in order, stops.  ~1-3 ms
CREATE INDEX idx_t_status_created ON tickets (tenant_id, status, created_at DESC);

-- Bad: range column first. Must read all of the tenant's tickets in time order and filter status.
CREATE INDEX idx_created_t ON tickets (created_at, tenant_id);

-- Also bad: (tenant_id, created_at, status) -- works for ORDER BY, but scans CLOSED rows
-- interleaved in time; if 98% are CLOSED, it reads ~2,500 entries to return 50.

Covering Indexes (Index-Only Scans)#

If every column the query reads is in the index, the engine never touches the table.

-- Postgres: INCLUDE adds payload columns to leaf pages without affecting sort order
CREATE INDEX idx_orders_cust_cover ON orders (customer_id, created_at DESC)
    INCLUDE (status, total_cents);

SELECT created_at, status, total_cents FROM orders
WHERE customer_id = $1 ORDER BY created_at DESC LIMIT 100;
-- Index Only Scan: 100 rows from ~2-3 leaf pages instead of ~100 random heap pages

When: hot, high-frequency queries that return many rows from a narrow column set. Cost: a wider index — more memory, more write bytes. In Postgres, index-only scans also depend on the visibility map being current (i.e., vacuum keeping up); on a churn-heavy table, "index only" can quietly revert to heap fetches.

Partial (Filtered) Indexes#

Index only the rows your queries target.

-- 500M jobs, 99.8% COMPLETED; workers poll the ~1M PENDING ones
CREATE INDEX idx_jobs_pending ON jobs (priority DESC, created_at)
    WHERE status = 'PENDING';
-- ~1M entries instead of 500M: fits in memory, ~250x smaller, cheap to maintain
-- Rows transitioning to COMPLETED are removed from the index

Also for: uniqueness with conditions (UNIQUE (email) WHERE deleted_at IS NULL — soft-delete-friendly), and sparse columns (index external_ref only where it is not null).

Other Index Types You Should Name#

TypeUse CaseExample
HashPure equality, no rangesPostgres hash index; in-memory hash in Redis
GIN (inverted)Containment: JSONB keys, arrays, full-textWHERE tags @> '{urgent}'
GiST / R-tree / SP-GiSTGeospatial, ranges, nearest-neighborPostGIS ST_DWithin — see Maps & Geospatial
BRINHuge append-ordered tables; stores min/max per block range10 TB time-series table; index is ~MBs instead of ~100s of GB
Expression indexQuery on a computed valueCREATE INDEX ON users (lower(email))
Vector (HNSW/IVF)Approximate nearest neighbor on embeddingspgvector, dedicated vector stores

Write Amplification and Secondary Indexes in Distributed Databases#

This is the hard sub-problem: indexes stop being free the moment writes are hot or data is distributed.

Write Amplification: The Real Cost of "Just Add an Index"#

For a Postgres table with a primary key and N secondary indexes, one INSERT does roughly:

heap tuple write                          1
index entry writes                        1 (PK) + N
WAL records                               1 + 1 + N   (plus full-page images after each checkpoint)
replica apply                             all of the above, on every replica
vacuum later                              visits heap + every index to clean dead entries
Secondary IndexesRelative Insert Cost (approx.)Typical Effect at 20K inserts/s
01×Comfortable on one primary
3~2–3×Fine; WAL volume noticeable
7~4–6×Replication lag during peaks; checkpoint I/O spikes
12+~8×+Write ceiling hit; the team starts discussing sharding

UPDATE is worse in Postgres: MVCC writes a new tuple version, and unless the update is HOT (heap-only tuple — no indexed column changed and the page has free space), every index gets a new entry, even for columns that didn't change. Indexing a frequently updated column (last_seen_at) disables HOT updates and can multiply write volume several-fold. Uber's widely read 2016 engineering post on moving from Postgres to MySQL cited exactly this write amplification across indexes and replicas as a primary motivation.

🎯 Staff Insight: "Before I add the eighth index to a hot table, I ask what it costs at peak write rate — in WAL bytes and replica lag — and whether the query could be served by the existing composite index, a partial index, or by moving it to a read replica or search index instead."

Secondary Indexes When Data Is Sharded#

When a table is partitioned by user_id, a query by user_id goes to one shard. A query by email does not know which shard to ask. There are exactly two designs.

DesignWrite PathRead PathConsistencyExamples
Local (document-partitioned) indexIndex lives on the same shard as the row; updated in the same local transactionScatter-gather to all N shards, merge resultsStrong (per shard)Cassandra secondary indexes, Elasticsearch shards, MongoDB non-shard-key indexes, DynamoDB LSIs
Global (term-partitioned) indexIndex partitioned by the indexed value; a write updates the row's shard and the index's shardSingle-shard lookup by indexed valueAsync (eventually consistent) or distributed transaction (slower writes)DynamoDB GSIs (async), Spanner / CockroachDB secondary indexes (transactional)
# Local index: read fans out
find_by_email(email):
    results = parallel_for shard in all_shards: shard.index_lookup("email", email)
    return merge(results)             # latency = slowest shard; cost = N queries

# Global index (async): write fans out
create_user(user):
    shard_for(user.id).insert(user)                     # source of truth, committed
    index_stream.publish(email=user.email, id=user.id)  # applied to index shard later (ms to s)
    # a read-by-email immediately after create may miss the user

Who pays:

  • Local index: readers pay. At 64 shards, one query becomes 64 queries; p99 is the max of 64 p99s. Tail latency amplification is the silent killer: if each shard's p99 is 10ms, the fan-out query's median is near that 10ms.
  • Global async index: correctness pays. Read-your-writes breaks. A uniqueness check against a GSI can race. The index can lag by seconds under write bursts.
  • Global synchronous index: writers pay. Each write becomes a cross-shard (often cross-node) transaction — Spanner/CockroachDB make it correct at the price of commit latency and coordination.

🎯 Staff Move: "Email lookup is low-QPS and must be unique, so I'd keep a separate email → user_id table partitioned by email and write both in a transaction — or, if the store can't do that, claim the email first with a conditional put and then create the user. For 'search users by name' I'd not use a DB index at all — I'd stream changes into Elasticsearch and accept a few seconds of lag."

Visual Guide#

LSM Write and Read Path#

Diagram: LSM Write and Read Path

Choosing an Index Strategy#

Diagram: Choosing an Index Strategy

Local vs Global Secondary Index Under Sharding#

Diagram: Local vs Global Secondary Index Under Sharding

Implementation Patterns#

Build Indexes Online, Never Blocking#

A plain CREATE INDEX in Postgres takes a lock that blocks writes for the entire build — on a 500 GB table, that can be an hour of outage. Use CREATE INDEX CONCURRENTLY (slower, two table passes, can leave an INVALID index on failure that you must drop and retry). In MySQL, use online DDL (ALGORITHM=INPLACE, LOCK=NONE) or tools like gh-ost / pt-online-schema-change. In any engine: build off-peak, watch replication lag, and have an abort plan.

Read the Plan, Not the Schema#

EXPLAIN (ANALYZE, BUFFERS)
SELECT ... FROM orders WHERE customer_id = 42 ORDER BY created_at DESC LIMIT 20;

Limit (actual rows=20)
  -> Index Scan using idx_orders_cust_created on orders  (actual time=0.04..0.21 ms)
       Buffers: shared hit=24

What to look for: Seq Scan on large tables, Rows Removed by Filter ≫ rows returned (wrong column order or missing predicate column), Sort nodes on paginated queries (index doesn't match ORDER BY), Heap Fetches on an expected index-only scan (vacuum behind), and estimate-vs-actual row mismatches > 10× (stale statistics).

Audit and Drop Unused Indexes#

-- Postgres: indexes never scanned since stats reset, largest first
SELECT relname, indexrelname, idx_scan, pg_size_pretty(pg_relation_size(indexrelid))
FROM pg_stat_user_indexes
WHERE idx_scan = 0 AND indexrelid NOT IN (SELECT conindid FROM pg_constraint)
ORDER BY pg_relation_size(indexrelid) DESC;

Check across all replicas (a replica-only reporting query may be the one user), across a full business cycle (month-end jobs), then drop. Some engines let you hide an index first (MySQL INVISIBLE) — a two-way door before the one-way drop.

Time-Ordered Keys for Insert-Heavy B-Trees#

Random UUIDv4 primary keys scatter inserts across every leaf page. At 100M+ rows, the index no longer fits in memory and each insert is a random read-modify-write. UUIDv7 / ULID / snowflake IDs keep inserts on the right edge of the tree, cutting buffer churn dramatically. The trade: a sequential key concentrates writes on one range — ideal for a single B-tree, a hot spot in a range-partitioned distributed database (see Sharding & Partitioning).

Offload Queries That Don't Belong on the Primary#

Not every query deserves an OLTP index. Analytical queries → columnar warehouse. Full-text and multi-field filtering → Elasticsearch fed by CDC. Rare admin queries → read replica with its own indexes (possible with logical replication). Every index moved off the primary is write capacity returned to the product.

The Numbers in Context#

NumberValueWhat It Means for Your Design
B-tree depth3–4 levels for 100M–10B rowsUpper levels are always cached; a point lookup ≈ 1 disk read when cold, ~0.05ms when warm
Page size8 KB Postgres, 16 KB InnoDBA 100-byte index insert dirties a whole page — the root of B-tree write amplification
Point lookup (warm)~0.05–0.5 ms in-engineNetwork + driver overhead (0.3–1 ms) usually dominates
Full scan throughput~0.5–2 GB/s per node from SSDA 50 GB table scan is ~25–100s — never on a request path
Bloom filter~10 bits/key → ~1% false positivesLSM point reads for absent keys skip nearly every SSTable
LSM write amp (leveled)~10–30×100 MB/s of logical writes can mean 1–3 GB/s of compaction I/O — size disks for it
Secondary index insert overhead~0.5–1× base insert cost each7 indexes ≈ 4–6× the write cost of the unindexed table
Selectivity thresholdPlanner prefers scan above ~5–20% of rowsIndexing a boolean or low-cardinality status column alone is usually useless
Scatter-gather fan-outQuery cost × N shardsAt 64 shards, a 1,000 QPS local-index query becomes 64,000 shard queries/s
Global index lag (async)Typically ms, seconds under burstsNever use an async index for uniqueness or read-after-write
Unused indexesCommonly 20–30% of indexes in mature schemasFree write capacity waiting to be reclaimed
Online index buildHours for 100s of GBSchedule, monitor replica lag, plan for failure and retry

How This Shows Up in Interviews#

Scenario 1: "How would you make this query fast?"#

The L5 answer: "Add an index on the WHERE column." The Staff answer derives the index from the full query — equality, sort, range — and states the cost: "Composite on (tenant_id, status, created_at DESC) with INCLUDE (subject), so it's an index-only seek of 50 entries. This table takes 3K writes/sec and already has four indexes; I'd check whether the existing (tenant_id, created_at) index can be replaced rather than adding a fifth."

Scenario 2: "Writes got slow after we launched a feature" (Full Walkthrough)#

Step 1 — Measure before guessing. "I'd look at write latency p99, WAL generation rate (bytes/sec), replication lag, and checkpoint frequency before and after the launch. If WAL bytes per insert doubled, something changed the physical write path."

Step 2 — Find the index change. "The feature added two indexes: one on last_active_at and a GIN index on a JSONB preferences column. last_active_at is updated on every request — indexing it disabled HOT updates, so every heartbeat update now writes a new entry into all six indexes on the table. That's the multiplier."

Step 3 — Fix the access pattern, not just the index. "Why does the feature query last_active_at? To find users active in the last 5 minutes. That belongs in a Redis sorted set or a small separate presence table, not an index on the hot users table. I'd drop the index and move the write off the main table. For the GIN index: GIN updates are expensive; enable fastupdate or, if the query is rare, move it to the search index."

Step 4 — Guard against recurrence. "Index DDL on tables above ~10K writes/sec goes through a checklist: estimated write overhead, expected QPS of the query it serves, and a shadow-period measurement of WAL bytes per write. Metrics: db.wal_bytes_per_sec, db.replication_lag_seconds, db.hot_update_ratio. The feature team owns the query; the database platform team owns the checklist and gets paged on lag."

Why this is a Staff answer: it quantifies the physical write path, identifies the mechanism (HOT updates disabled), fixes the access pattern rather than tuning the index, and installs a process guard with an owner.

Scenario 3: "Users need to be looked up by email, but we shard by user_id"#

Name the two options (scatter-gather over local indexes vs a global index), then pick based on semantics: login requires uniqueness and read-after-write, so a synchronously maintained email → user_id mapping partitioned by email — with the email claim as the first, conditional step of signup. Mention that this is exactly the DynamoDB pattern of a separate uniqueness item written in a transaction, because a GSI is eventually consistent and cannot enforce uniqueness.

Scenario 4: "Our Cassandra reads got 20× slower on one table"#

Tombstones. A queue-like table with deletes accumulates tombstones inside partitions; reads scan through them before gc_grace_seconds (default 10 days) allows compaction to purge them. Look for tombstone_warn_threshold warnings in logs and TombstoneScannedHistogram. The fix is a data-model change — time-bucketed partitions that are dropped whole, or TTLs with a time-window compaction strategy — not a tuning knob. See Cassandra.

What Interviewers Probe#

After You Say...They Will Ask...(What They're Evaluating)
"Add an index on user_id""This table takes 20K writes/sec. What does that index cost?"Whether you see indexes as a write-side tax
"Composite index on (a, b, c)""Why that order? What about a query on b alone?"Leftmost-prefix and ESR reasoning
"Use Cassandra for the write rate""Now product wants to filter by status. How?"Whether you know LSM secondary indexes are weak and model a second table
"Global secondary index for email""Two users sign up with the same email at the same instant."Async index vs uniqueness; conditional writes
"We'll build the index in production""The table is 800 GB. Walk me through the rollout."Online DDL, replica lag, abort plan
"Index-only scan will make it fast""It's still doing heap fetches. Why?"Visibility map / vacuum awareness — engine depth, not trivia

Advanced Patterns#

PatternHow It WorksWhen to Use
Covering index with INCLUDENon-key payload columns in leaf pagesHot range queries over a narrow column set
Partial indexIndex only rows matching a predicateSkewed status distributions, soft deletes, sparse columns
Index-organized / clustered tableRows stored in PK orderRange scans on PK dominate (InnoDB default)
BRINMin/max summaries per block rangeMulti-TB append-ordered tables
Materialized index tableApp-maintained lookup table in a different partition keyGlobal secondary access in stores without transactional global indexes
CDC to search indexStream row changes into Elasticsearch/OpenSearchMulti-field filtering, full-text, "search by anything"
Invisible / hidden indexIndex maintained but ignored by plannerTest drop impact safely before the one-way door
Time-window compactionLSM SSTables grouped by time window, dropped whole when expiredTTL'd time-series in Cassandra/Scylla

Failure Modes & Operational Reality#

FailureDetection SignalBlast RadiusMitigationOwner
Index build locks writesdb.lock_wait_seconds spike; write errors during DDLEntire table; often the whole product surfaceCONCURRENTLY / online DDL; DDL review; off-peak windowDatabase platform + migrating team
Planner flips to a bad planQuery p99 jumps 100×; seq_scan count rises; no deployEvery caller of that queryFresh statistics (ANALYZE), extended stats, plan pinning as last resortOwning service team
Write amplification from index sprawldb.wal_bytes_per_sec ↑, db.replication_lag_seconds ↑, hot_update_ratio ↓Primary and all replicas; stale reads from replicasDrop unused indexes, avoid indexing hot-updated columnsDatabase platform
LSM compaction debtpending_compactions ↑, sstables_per_read > 5, p99 ↑All reads on the node; disk-full riskThrottle ingest, add nodes, switch compaction strategyStorage on-call
Tombstone scantombstones_scanned_per_read ↑; read timeouts on one tableOne table, often one hot partitionRemodel partitions; TTL + TWCSOwning service team
Global index laggsi.replication_lag_ms ↑; users "missing" right after createReads by secondary attributeRead from base table for read-your-writes; alert on lagOwning service team

Failure Scenario: The Midnight Plan Flip#

t=0       Nightly batch deletes 40M expired sessions from a 60M-row table.
t=+20min  Autovacuum/analyze runs; statistics now show 20M rows with a different status distribution.
t=+21min  Planner switches login lookup from idx_sessions_user to a Seq Scan (bad estimate after churn).
t=+22min  Login p99: 3ms → 1.8s. DB CPU 95%. Connection pool saturates; auth timeouts cascade.
t=+30min  On-call sees CPU but no deploy. pg_stat_statements shows one query at 98% of total time.
t=+38min  Manual ANALYZE with higher statistics target; plan reverts. Recovery by t=+40min.

Detection: alert on per-query p99 regressions from pg_stat_statements (not just aggregate CPU), and on seq_scan increases for tables above 1M rows. Prevention: avoid mass deletes on hot tables (partition by time and DROP PARTITION instead), raise statistics targets on skewed columns, and run ANALYZE as part of the batch job. Owner: the session service owns the batch design; the database platform owns query-regression alerting.

The Principal Lens#

Why L7 Sees This Problem Differently#

At Staff level, indexing is a per-query, per-table decision. At Principal level, indexes are the largest uncontrolled consumer of database capacity in the company. Every team adds indexes; almost no team removes them; nobody's dashboard shows the write cost of an index they didn't create. The Principal question is how the organization governs a shared resource where the benefit accrues to one team's query and the cost is paid by everyone who writes to the table and every replica that applies the WAL — and when to stop indexing in the OLTP store altogether and move query patterns to purpose-built systems.

The Org-Level Fault Line#

The shared OLTP database as "query anything" vs the OLTP database as a narrow system of record. If every team can add indexes to shared tables, the primary becomes a search engine, a reporting database, and a queue simultaneously — and its write ceiling drops every quarter. If the primary is locked down to its core access paths, teams must stand up CDC pipelines, search clusters, and warehouses for anything else, which costs platform headcount.

OptionWhat WorksWhat BreaksWho Pays
Any team may index shared tablesFast feature deliveryWrite ceiling erodes; incidents blamed on "the database"Database platform on-call; every writer
DBA approval for every indexControlled capacity1–2 week queue; teams add caches and shadow tables insteadProduct velocity
Core access paths only + paved-road CDC to search/warehousePrimary stays fast; secondary patterns get fit-for-purpose storesCDC and search infra need owners; eventual consistency for secondary queriesPlatform team (2–4 FTE for CDC + search)

Cost Model#

Assumptions: managed Postgres-class primary with 2 replicas; storage ~$0.10–0.25/GB-month; engineer ≈ $25K/month fully loaded.

ScaleData / WritesIndex FootprintMonthly Infra Attributable to IndexesPeople / On-Call
Small200 GB, 500 writes/s~100 GB across 3 nodes~$500–1K (storage + slightly larger instances)Part of a generalist's time
Medium5 TB, 10K writes/s~3–5 TB × 3 nodes; ~40% of IOPS spent on index maintenance~$10–25K (larger instances and provisioned IOPS largely driven by index writes)1 FTE DB reliability; index-related pages ~1–2/month
Large (sharded, 64 shards)200 TB, 500K writes/sGlobal indexes + search cluster + CDC~$150–400K across OLTP indexes, search cluster, CDC pipeline4–8 FTE storage/search platform; indexing reviews part of design review

The lever: at medium scale, removing 25% unused indexes and moving two query patterns to a search index can defer sharding by 12–18 months — worth far more than the index storage itself, because sharding typically costs 2–4 engineers for 2–3 quarters.

The 3-Year Evolution Path#

Diagram: The 3-Year Evolution Path

One-Way Doors vs Two-Way Doors#

DecisionReversibilityCost to Reverse
Adding a secondary indexTwo-wayDrop it (hide first) — but callers may silently depend on it
Dropping an indexTwo-way, slowlyRebuild takes hours on large tables; outage risk if a critical query depended on it
Primary key type (auto-increment vs UUIDv4 vs UUIDv7)One-wayRewrite of every table and every foreign key; client-visible IDs
Storage engine family (B-tree RDBMS vs LSM wide-column)One-wayFull data model redesign and migration — quarters
Global secondary index semantics (async vs transactional)One-way-ishApplication logic built on the consistency assumption
Moving a query pattern to search via CDCTwo-wayThe pipeline can be kept or retired

The Standard I'd Write#

RFC-DB-004: Index Governance for Shared OLTP Databases

Scope: Tables with > 1K writes/sec or > 100 GB on any shared primary.

Mandatory (MUST):

  1. New indexes MUST be built online (CONCURRENTLY / online DDL) and MUST cite the query they serve and its expected QPS.
  2. Tables above 10K writes/sec MUST NOT index columns updated more than once per minute per row.
  3. Each index MUST have an owning team recorded in the schema catalog.
  4. Indexes with zero scans over 60 days on all replicas MUST be hidden, then dropped after 30 more days, unless the owner files an exception.

Recommended (SHOULD): use composite indexes following equality → sort → range; prefer partial indexes for skewed predicates; route full-text and ad-hoc multi-field queries through the CDC → search paved road.

Exceptions: 1-page request to database platform; decided in 3 business days.

Success metrics: WAL bytes per logical write (−20% in 2 quarters); % of indexes with owners (100%); index-related Sev2+ incidents (target 0 per quarter).

What I'd Tell the VP#

"Our main database is slowing down not because we have too much data, but because every team has added lookups to it and nobody removes them — about a quarter of them are unused, and each one makes every write more expensive. If we do nothing, we hit the write limit in roughly a year and have to split the database, which is a two-to-three-quarter project for several engineers. I'm proposing a lightweight index policy with ownership and automated cleanup, and moving search-style queries to our search system. That should buy us 12 to 18 months before any sharding work and reduce database-related incidents."

Principal Interview Signals#

SignalWhat It Sounds Like
Treats indexes as shared capacity"This index helps one team's query and costs every writer. I'd make that cost visible per owner."
Prices deferral of sharding"Dropping 30 unused indexes buys us ~40% write headroom — that's a year of not sharding."
Knows when to leave the OLTP store"Five teams want 'filter by anything.' That's a search product, not five more B-trees."
Identifies one-way doors in key choice"The primary key type is the decision we can't undo. Index choices we can."
Designs governance through tooling"Unused-index detection and DDL linting in CI, not a DBA approval queue."

Staff answers that L7 interviewers find insufficient:

  • "Add a composite index with the right column order" — correct for this query, silent on the table's cumulative write budget.
  • "Use a GSI for the email lookup" — misses that the org needs a standard pattern for uniqueness across sharded stores.
  • "We'll audit indexes when there's a problem" — reactive; no ownership, no automation, no metric.

🧭 Principal Move: "I'd publish WAL bytes per write per table on the same dashboard as each team's error rate. The moment index cost is attributed to an owner, most of the sprawl stops on its own."

In the Wild#

Uber: Write Amplification and the Postgres-to-MySQL Post#

Uber's 2016 engineering blog post on migrating from Postgres to MySQL described how Postgres's MVCC and physical-pointer indexes caused updates to write new entries to every index and replicate large volumes of WAL, amplifying writes on tables with many indexes and complicating replication. Whatever one thinks of the conclusion — the Postgres community responded with detailed counterpoints — the post is a canonical public case of index-driven write amplification shaping a storage decision.

Staff insight: index cost is an engine-specific physical cost. Name the mechanism (HOT updates, secondary indexes pointing to PKs vs tuple IDs) instead of saying "Postgres doesn't scale."

Amazon DynamoDB: Global Secondary Indexes Are Eventually Consistent#

DynamoDB documents that global secondary indexes are updated asynchronously and support only eventually consistent reads, while local secondary indexes share the partition key with the base table and support strongly consistent reads. AWS's own guidance for uniqueness on non-key attributes is to write a separate item keyed by that attribute inside a transaction.

Staff insight: this is the local-vs-global tradeoff made explicit in a product's API. In an interview, citing it shows you know a GSI is not a unique constraint.

Google Bigtable, Cassandra, and the LSM Lineage#

Google's Bigtable paper (2006) popularized the memtable/SSTable/compaction design with bloom filters; Cassandra, HBase, LevelDB and RocksDB followed. RocksDB, developed at Facebook, became the embedded LSM engine inside many databases, including CockroachDB's early storage layer and MyRocks, which Facebook used to reduce MySQL storage footprint for its user database.

Staff insight: LSM trades read and space predictability for write throughput and compression. Choosing it is choosing to design the data model around the primary key.


Staff Calibration#

What Staff Engineers Say (That Seniors Don't)#

ConceptSenior (L5)Staff (L6)Principal (L7)
Adding an index"Index the WHERE column""Composite in equality → sort → range order, covering if hot; here's the write cost at peak""Index cost is attributed to an owner; unused indexes are auto-hidden after 60 days"
Engine choice"Postgres is reliable""Write-to-read ratio and access paths pick B-tree vs LSM""Engine family is a one-way door; I'd standardize two paved roads, not five engines"
Sharded lookup"Add a secondary index""Local = scatter-gather, global = async or 2PC; uniqueness needs a synchronous mapping""One org pattern for secondary access in sharded stores so 20 teams don't reinvent it"
Slow writes"Bigger instance""WAL bytes/write and HOT ratio point to index sprawl on a hot column""Write headroom is a capacity budget; index governance defers a 3-quarter sharding project"
Search queries"Add a LIKE index""Full-text belongs in a search index via CDC""CDC-to-search is a paved road with an owner and a freshness SLO"
Why "Adding an index" separates levels

Everyone knows indexes speed reads. The L6 gap is that indexes are a write-side decision: column order determines whether it helps at all, and every index is paid for at peak write rate on primary and replicas. The L7 gap is that nobody sees this cost unless someone makes it visible per owner.

Why "Sharded lookup" separates levels

L5 answers as if the database were a single node. L6 knows that a secondary index across shards forces a choice between read fan-out and write fan-out, and that async global indexes cannot enforce uniqueness. L7 recognizes that every team with a sharded store will face this and standardizes the pattern.

Common Interview Traps#

  • Indexing every column in the WHERE clause separately. Five single-column indexes are not a composite index; the planner usually uses one.
  • Wrong column order. Range column first kills the index for equality-plus-sort queries.
  • Indexing low-cardinality columns alone. status with 4 values — use a partial or composite index instead.
  • Forgetting writes. Never propose an index on a hot table without saying what it costs at peak write rate.
  • Assuming a GSI is a unique constraint. Async global indexes cannot enforce uniqueness.
  • Random UUIDv4 primary keys on a large B-tree. Insert locality disappears; use time-ordered IDs.
  • Blocking DDL. Proposing CREATE INDEX on a 500 GB production table without the online variant.
  • Ignoring tombstones in LSM stores. Delete-heavy patterns in Cassandra are a data-model bug.

Practice Drill#

Prompt: "A multi-tenant SaaS events table has 2 billion rows, 30K inserts/sec, and 9 secondary indexes. Replication lag hits 30 seconds at peak and the team wants to shard. What do you do first?"

Staff Answer

Before sharding, I'd find out how much of the write cost is indexes. I'd pull pg_stat_user_indexes across the primary and all replicas: my prior is 2–3 of the 9 indexes have near-zero scans — hide, then drop them. Next I'd check which indexed columns are updated after insert; events should be append-only, so any updated indexed column is a modeling bug. Then I'd map the remaining indexes to queries: tenant dashboards probably need (tenant_id, created_at); free-form filters on event_type, user_agent, country belong in a search or columnar store fed by CDC, not the OLTP primary. I'd time-partition the table by day so retention is DROP PARTITION, not 40M-row deletes, and so each partition's indexes stay small and hot. Metrics to confirm: WAL bytes/insert (expect −40–60%), replication lag p99 at peak, and index-only scan ratio. If after that we're still within ~12 months of the write ceiling, then shard by tenant_id with a plan for hot tenants — but now with 4 indexes instead of 9.

Why this is L6:

  • Treats index count as the first lever on write throughput, with a measurement plan.
  • Moves query patterns that don't belong in OLTP to fit-for-purpose stores.
  • Uses time partitioning to fix retention deletes and keep indexes small.

What L7 adds:

  • Prices the alternative: sharding ≈ 3 engineers × 2–3 quarters vs ~1 engineer × 1 quarter for index cleanup and CDC.
  • Turns the one-time cleanup into governance: index ownership, auto-hide after 60 days, WAL-per-write dashboards per team.
  • Decides whether the org needs a shared events/analytics platform so the next team doesn't put 2B rows into OLTP.

Where This Appears#

Related Technologies: PostgreSQL · Cassandra · DynamoDB · Elasticsearch

  1. Loading the index…