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)servesWHERE 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%COMPLETEDgives 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#
| Term | Meaning | Why It Matters |
|---|---|---|
| Clustered index | Table 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 index | Rows stored unordered; all indexes point to physical location (Postgres) | Updates that move a row can require updating every index (mitigated by HOT updates) |
| Selectivity | Fraction of rows matching a predicate | Index on is_active (90% true) is useless; the planner will scan |
| Cardinality | Number of distinct values | High cardinality → selective → index-friendly |
| Covering index | Index containing every column the query needs | Index-only scan, no heap fetch |
| Write amplification | Physical bytes written ÷ logical bytes written | Determines disk wear, IOPS budget, and max write throughput |
| Read amplification | Physical reads per logical read | LSM point reads may touch several SSTables without bloom filters |
| Space amplification | Physical size ÷ logical data size | LSM with stale versions pending compaction; B-tree page fill factor |
Where Indexes Live#
| Engine Family | Index Structure | Examples | Sweet Spot |
|---|---|---|---|
| Relational OLTP | B+tree (clustered or heap) | PostgreSQL, MySQL/InnoDB, SQL Server, Oracle | Read-heavy, mixed queries, secondary indexes, transactions |
| LSM key-value / wide-column | LSM tree + bloom filters | Cassandra, RocksDB, ScyllaDB, HBase, (DynamoDB's internals are not public — treat as managed) | Write-heavy, key-based access, time series |
| NewSQL | LSM (RocksDB/Pebble) under distributed SQL | CockroachDB, TiDB, YugabyteDB | SQL + horizontal scale; secondary indexes are distributed |
| Search | Inverted index | Elasticsearch, OpenSearch, Lucene | Full-text, faceted, "any field" queries |
| Specialized | GiST/R-tree, GIN, BRIN, HNSW | PostGIS, Postgres JSONB, vector indexes | Geospatial, 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#
| Dimension | B-Tree | LSM |
|---|---|---|
| Point read | ~3–4 page reads, predictable | Memtable + bloom checks per level; 1–2 disk reads typical, worse under compaction debt |
| Range scan | Excellent — linked leaves | Good, but merges across levels; tombstones hurt |
| Write throughput | Bounded by random IOPS for random keys | Bounded 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 amplification | Fragmentation, fill factor ~70% | Size-tiered can need ~2× during compaction; leveled ~1.1× |
| Latency predictability | High | Compaction causes p99 spikes unless throttled |
| Choice | What Works | What Breaks | Who Pays |
|---|---|---|---|
| B-tree (Postgres/MySQL) | Rich secondary indexes, predictable reads, transactions | Random-write ceiling; bloat and vacuum on churn-heavy tables | DBAs/on-call pay in vacuum tuning; product pays when write growth forces early sharding |
| LSM (Cassandra/RocksDB) | 5–10× higher sustained write throughput per node | Secondary indexes weak; tombstones; compaction p99 spikes | Application 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#
| Type | Use Case | Example |
|---|---|---|
| Hash | Pure equality, no ranges | Postgres hash index; in-memory hash in Redis |
| GIN (inverted) | Containment: JSONB keys, arrays, full-text | WHERE tags @> '{urgent}' |
| GiST / R-tree / SP-GiST | Geospatial, ranges, nearest-neighbor | PostGIS ST_DWithin — see Maps & Geospatial |
| BRIN | Huge append-ordered tables; stores min/max per block range | 10 TB time-series table; index is ~MBs instead of ~100s of GB |
| Expression index | Query on a computed value | CREATE INDEX ON users (lower(email)) |
| Vector (HNSW/IVF) | Approximate nearest neighbor on embeddings | pgvector, 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 Indexes | Relative Insert Cost (approx.) | Typical Effect at 20K inserts/s |
|---|---|---|
| 0 | 1× | 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.
| Design | Write Path | Read Path | Consistency | Examples |
|---|---|---|---|---|
| Local (document-partitioned) index | Index lives on the same shard as the row; updated in the same local transaction | Scatter-gather to all N shards, merge results | Strong (per shard) | Cassandra secondary indexes, Elasticsearch shards, MongoDB non-shard-key indexes, DynamoDB LSIs |
| Global (term-partitioned) index | Index partitioned by the indexed value; a write updates the row's shard and the index's shard | Single-shard lookup by indexed value | Async (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_idtable 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#
Choosing an Index Strategy#
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#
| Number | Value | What It Means for Your Design |
|---|---|---|
| B-tree depth | 3–4 levels for 100M–10B rows | Upper levels are always cached; a point lookup ≈ 1 disk read when cold, ~0.05ms when warm |
| Page size | 8 KB Postgres, 16 KB InnoDB | A 100-byte index insert dirties a whole page — the root of B-tree write amplification |
| Point lookup (warm) | ~0.05–0.5 ms in-engine | Network + driver overhead (0.3–1 ms) usually dominates |
| Full scan throughput | ~0.5–2 GB/s per node from SSD | A 50 GB table scan is ~25–100s — never on a request path |
| Bloom filter | ~10 bits/key → ~1% false positives | LSM 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 each | 7 indexes ≈ 4–6× the write cost of the unindexed table |
| Selectivity threshold | Planner prefers scan above ~5–20% of rows | Indexing a boolean or low-cardinality status column alone is usually useless |
| Scatter-gather fan-out | Query cost × N shards | At 64 shards, a 1,000 QPS local-index query becomes 64,000 shard queries/s |
| Global index lag (async) | Typically ms, seconds under bursts | Never use an async index for uniqueness or read-after-write |
| Unused indexes | Commonly 20–30% of indexes in mature schemas | Free write capacity waiting to be reclaimed |
| Online index build | Hours for 100s of GB | Schedule, 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#
| Pattern | How It Works | When to Use |
|---|---|---|
| Covering index with INCLUDE | Non-key payload columns in leaf pages | Hot range queries over a narrow column set |
| Partial index | Index only rows matching a predicate | Skewed status distributions, soft deletes, sparse columns |
| Index-organized / clustered table | Rows stored in PK order | Range scans on PK dominate (InnoDB default) |
| BRIN | Min/max summaries per block range | Multi-TB append-ordered tables |
| Materialized index table | App-maintained lookup table in a different partition key | Global secondary access in stores without transactional global indexes |
| CDC to search index | Stream row changes into Elasticsearch/OpenSearch | Multi-field filtering, full-text, "search by anything" |
| Invisible / hidden index | Index maintained but ignored by planner | Test drop impact safely before the one-way door |
| Time-window compaction | LSM SSTables grouped by time window, dropped whole when expired | TTL'd time-series in Cassandra/Scylla |
Failure Modes & Operational Reality#
| Failure | Detection Signal | Blast Radius | Mitigation | Owner |
|---|---|---|---|---|
| Index build locks writes | db.lock_wait_seconds spike; write errors during DDL | Entire table; often the whole product surface | CONCURRENTLY / online DDL; DDL review; off-peak window | Database platform + migrating team |
| Planner flips to a bad plan | Query p99 jumps 100×; seq_scan count rises; no deploy | Every caller of that query | Fresh statistics (ANALYZE), extended stats, plan pinning as last resort | Owning service team |
| Write amplification from index sprawl | db.wal_bytes_per_sec ↑, db.replication_lag_seconds ↑, hot_update_ratio ↓ | Primary and all replicas; stale reads from replicas | Drop unused indexes, avoid indexing hot-updated columns | Database platform |
| LSM compaction debt | pending_compactions ↑, sstables_per_read > 5, p99 ↑ | All reads on the node; disk-full risk | Throttle ingest, add nodes, switch compaction strategy | Storage on-call |
| Tombstone scan | tombstones_scanned_per_read ↑; read timeouts on one table | One table, often one hot partition | Remodel partitions; TTL + TWCS | Owning service team |
| Global index lag | gsi.replication_lag_ms ↑; users "missing" right after create | Reads by secondary attribute | Read from base table for read-your-writes; alert on lag | Owning 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.
| Option | What Works | What Breaks | Who Pays |
|---|---|---|---|
| Any team may index shared tables | Fast feature delivery | Write ceiling erodes; incidents blamed on "the database" | Database platform on-call; every writer |
| DBA approval for every index | Controlled capacity | 1–2 week queue; teams add caches and shadow tables instead | Product velocity |
| Core access paths only + paved-road CDC to search/warehouse | Primary stays fast; secondary patterns get fit-for-purpose stores | CDC and search infra need owners; eventual consistency for secondary queries | Platform 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.
| Scale | Data / Writes | Index Footprint | Monthly Infra Attributable to Indexes | People / On-Call |
|---|---|---|---|---|
| Small | 200 GB, 500 writes/s | ~100 GB across 3 nodes | ~$500–1K (storage + slightly larger instances) | Part of a generalist's time |
| Medium | 5 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/s | Global indexes + search cluster + CDC | ~$150–400K across OLTP indexes, search cluster, CDC pipeline | 4–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#
One-Way Doors vs Two-Way Doors#
| Decision | Reversibility | Cost to Reverse |
|---|---|---|
| Adding a secondary index | Two-way | Drop it (hide first) — but callers may silently depend on it |
| Dropping an index | Two-way, slowly | Rebuild takes hours on large tables; outage risk if a critical query depended on it |
| Primary key type (auto-increment vs UUIDv4 vs UUIDv7) | One-way | Rewrite of every table and every foreign key; client-visible IDs |
| Storage engine family (B-tree RDBMS vs LSM wide-column) | One-way | Full data model redesign and migration — quarters |
| Global secondary index semantics (async vs transactional) | One-way-ish | Application logic built on the consistency assumption |
| Moving a query pattern to search via CDC | Two-way | The 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):
- New indexes MUST be built online (
CONCURRENTLY/ online DDL) and MUST cite the query they serve and its expected QPS.- Tables above 10K writes/sec MUST NOT index columns updated more than once per minute per row.
- Each index MUST have an owning team recorded in the schema catalog.
- 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#
| Signal | What 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)#
| Concept | Senior (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.
statuswith 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 INDEXon 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
eventstable 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#
- Database Sharding — secondary indexes across shards, global vs local, and resharding
- Search Indexing — inverted indexes and CDC-fed search as the alternative to OLTP indexes
- Maps & Geospatial — spatial indexes (geohash, R-tree, S2) for proximity queries
- Database Selection — B-tree vs LSM as a selection criterion
- Metrics & Monitoring — LSM and time-partitioned storage for time series
- Scaling Writes — index write amplification as a write-scaling limit
- Data Modeling & Schema Design — designing tables around access paths first
- Sharding & Partitioning — partition keys and distributed secondary indexes
Related Technologies: PostgreSQL · Cassandra · DynamoDB · Elasticsearch