Indexing Strategy and Design Questions
Choosing and designing indexes: B-tree, hash, composite, covering, partial, and full-text/inverted indexes, and the trade-offs between read acceleration and write/storage overhead. Covers selecting index columns from query patterns, cardinality and selectivity reasoning, diagnosing why an index is or is not used, and index maintenance: rebuilding or reorganizing a fragmented index, finding and dropping redundant or unused indexes, and rolling out a new index to production safely. Also covers indexing in analytical (bitmap, columnar), partitioned, and distributed/NoSQL systems. Central to database performance interviews.
Design an indexing strategy for a star-schema data warehouse: a 1 billion-row fact table (sales) joined to dimension tables (product, store, date). Typical queries compute time-range totals by product category and store. What would you index, and what alternatives to a traditional B-tree index would you consider at this scale?
Sample Answer
Direct answer
At 1 billion fact rows, don't reach for a single index type; layer three complementary techniques: partition the fact table (by date is the natural choice for time-range-by-category queries), use B-tree indexes on the foreign keys to small, low-cardinality dimensions only where the optimizer benefits from them (often not needed if the query always aggregates rather than point-looks-up), and use a BRIN (Block Range INdex, Postgres's built-in "zone map" style structure) index on any column that is naturally correlated with physical insert order, typically the load date or an ever-increasing surrogate key, because at this scale a BRIN index costs a tiny fraction of a B-tree's storage for range-style filtering. If the workload also needs sub-second dashboard response against a sliding recent window under real concurrency, add a second, narrower path: a materialized, pre-aggregated summary table refreshed on a schedule, since no index strategy alone reliably gets an ad hoc billion-row aggregation under a hard sub-second SLA (service-level objective) at real concurrency.
Structured elaboration
graph TD
F["fact_sales (1B rows)<br/>partitioned by date"] --> D1[dim_product]
F --> D2[dim_store]
F --> D3[dim_date]
Partitioning. Range-partitioning fact_sales by date (monthly or quarterly, tuned to keep individual partitions from getting unwieldy) lets a "time-range totals" query prune to only the relevant partitions before touching any index at all, often the single biggest win at this scale, since it shrinks the actual data volume any single query has to consider.
B-tree on dimension foreign keys. Whether a foreign key to product or store needs its own B-tree index depends on the query shape: a query that AGGREGATES (GROUP BY product_category) over a large date range benefits more from a good join order and partition pruning than from an index seek per row, since it's going to touch a large fraction of the partition anyway. A query that FILTERS down to one or a few specific products or stores before aggregating benefits from a real index on that foreign key, the same reasoning as any selective equality predicate.
BRIN for naturally-ordered columns, and why it wins at this scale specifically. A BRIN index doesn't store one entry per row; it stores a min/max (or similar) summary per block RANGE (a configurable group of physical pages), so its size barely grows with row count as long as the indexed column stays correlated with physical insert order, exactly the shape of a date or load-order column in a fact table that is loaded roughly chronologically.
Worked example
Extrapolated from a MEASURED 3-million-row sample (a naturally time-ordered append log on Postgres 16, real pg_relation_size figures), scaled linearly to 1 billion rows as an order-of-magnitude estimate, not a claim that behavior is perfectly linear at 333x the sample size:
measured_rows = 3_000_000
brin_bytes = 24 * 1024 # MEASURED
btree_bytes = 64 * 1024 * 1024 # MEASURED, same column, plain B-tree
table_bytes = 184 * 1024 * 1024 # MEASURED
target_rows = 1_000_000_000
scale = target_rows / measured_rows # 333.3x
def gib(b): return b / 1024**3
print(gib(table_bytes * scale)) # ~59.9 GiB fact table storage estimate
print(gib(brin_bytes * scale)) # ~0.0076 GiB (~7.8 MiB) BRIN index estimate
print(gib(btree_bytes * scale)) # ~20.8 GiB B-tree index estimate, same column
print(btree_bytes * scale / (brin_bytes * scale)) # ~2731x
The B-tree, on the SAME naturally-ordered column, would be an estimated 2,731 times larger than the BRIN index at this scale (about 20.8 GiB versus about 7.8 MiB), while the real, measured query behavior at the smaller scale showed BRIN answering a 1-day range filter (about 1/35th of the loaded date range) in single-digit milliseconds, competitive with or faster than the same query using the full B-tree, because a BRIN scan reads its whole tiny summary structure in a handful of pages and then only fetches the heap blocks the summary says could match, rather than walking a much deeper, much larger tree.
For the SLA-driven concurrency scenario, a concrete worked framing: 1,000 concurrent dashboard users, each running a sub-second aggregate query over a sliding 7-day window. Partition pruning gets each query down to roughly the last 2 to 3 monthly partitions instead of all 36+ months of history; a BRIN index on the load-date column inside each active partition gets the block-range filtering essentially free; but 1,000 CONCURRENT aggregate queries, even each individually fast, compete for the same buffer cache and I/O bandwidth, which is exactly the case where a scheduled, pre-aggregated rollup table (refreshed every few minutes, holding the last 7 days pre-summed by category and store) removes the aggregation cost from the hot path entirely and turns each dashboard hit into a cheap point lookup against a table sized in the thousands or low millions of rows, not billions.
Trade-offs & pitfalls
BRIN is not a free upgrade over B-tree; it only works well when the indexed column is genuinely correlated with physical row order. Indexing product_category with BRIN when categories are scattered randomly through the table (not the load order) would produce nearly useless summary ranges, since almost every block range would contain almost every category, making it deliver almost no pruning at all: for scattered, low-cardinality dimension columns, the earlier bitmap-vs-B-tree cardinality reasoning applies instead, not BRIN.
Denormalizing the star schema further, materializing pre-joined and pre-aggregated fact/dimension combinations, is a real, complementary escape hatch once indexing and partitioning alone can't hit a hard latency target under real concurrency, exactly the "materialized rollup" pattern above; the cost is staleness (the rollup lags the raw facts by its refresh interval) and a second thing to keep correctly in sync, not a free win.
At genuinely extreme scale, or if this workload is purely analytical with no OLTP (online transaction processing) concerns at all, a dedicated columnar/OLAP (online analytical processing) engine (Redshift, BigQuery, Snowflake, ClickHouse) automates a lot of what's being hand-rolled here: columnar storage, automatic zone-map pruning, vectorized aggregation. It is worth naming as the alternative to continuing to hand-tune Postgres indexes, particularly once the honest answer to whether a general-purpose row store can hit this SLA at this scale starts to look like: only with significant, ongoing manual tuning.
Discuss when you would use a clustered columnstore index versus a nonclustered columnstore index in SQL Server for a large analytical fact table that receives frequent micro-batch loads. Explain how the delta store and tuple-mover affect query and load performance, and how you would structure data loads to minimize fragmentation.
Sample Answer
Direct answer: For a large analytical fact table (a big, append-heavy table of business events or transactions, as opposed to a small reference "dimension" table like a customer or product list) that receives frequent micro-batch loads, use a CLUSTERED columnstore index (CCI) if this table's primary storage should just BE stored column-by-column, which is the normal choice for a warehouse fact table, since that's the natural fit for the scan-heavy queries fact tables usually see; use a NONCLUSTERED columnstore index (NCCI) only when you need to keep the table's primary storage in its traditional row-by-row layout, for fast single-row lookups and updates, and add a second, columnar copy alongside it purely to accelerate analytics run concurrently against the same data. Either way, micro-batch loads interact with the same underlying mechanism, the delta store, and the load pattern you choose has a direct, controllable effect on fragmentation.
Row-store vs. columnstore, and what "heap," "B-tree," and "clustered index" mean here
Two physical-storage concepts sit underneath this whole comparison, so it's worth naming them plainly before going further.
A row-store is the traditional table layout: every column belonging to one row is stored together on disk, one row after another. That's efficient for reading or changing ONE specific row (an OLTP query, short for online transaction processing: live application traffic doing many small reads and writes to individual rows, like one order being created or updated), since everything about that row sits in one place. It's inefficient for a query that scans millions of rows but only needs a handful of the table's columns, because the engine still has to read every row's full width off disk just to get those few columns.
A columnstore (columnar storage) flips that layout: it stores each COLUMN's values together instead, e.g. every order_date value in one place, every amount value in another. A large analytical scan that only touches a few columns now reads just those columns, not the full row width, and because values within one column tend to repeat or trend together, they also compress far better sitting next to each other than interleaved with unrelated columns. The cost shows up on the write side: touching one row now means touching several separate per-column structures instead of one row location, which is why a table under constant single-row transactional writes suits a row-store better than a columnstore.
Within a row-store, SQL Server organizes a table's rows one of two ways. A heap has no defined row order at all, a new row just goes wherever there's free space. A clustered index is not a separate structure sitting beside the table, it IS the table: the table's rows are physically stored on disk sorted in that index's key order, which is why a table can have at most one clustered index (or, if it defines none, it's stored as a heap instead). Once you add any index to speed up lookups, that index is typically implemented as a B-tree (specifically a B+ tree in SQL Server, a balanced, sorted tree structure that finds one row in a small, predictable number of steps instead of scanning every row).
CCI and NCCI, below, are both just ways of getting the columnar layout: either as the table's ONLY physical storage (CCI, replacing the heap/B-tree row-store entirely) or as a second, columnar copy layered next to a row-store original that keeps serving OLTP traffic (NCCI).
Clustered vs. nonclustered columnstore
- A clustered columnstore index (CCI) IS the table: it's the only, primary physical storage for every row and column, replacing what would otherwise be a heap or a B-tree clustered index. You create it directly (
CREATE TABLE ... WITH (INDEX ... CLUSTERED COLUMNSTORE)) or convert an existing rowstore table to one. - A nonclustered columnstore index (NCCI) is a secondary structure layered on top of a table whose primary storage stays row-based (a heap or a B-tree). It holds a compressed columnar copy of some or all columns, letting transactional (OLTP) queries keep using the row-store B-tree while analytical queries run against the columnstore copy concurrently, an SQL Server feature Microsoft calls real-time operational analytics. Since 2016, an NCCI is updatable and stays in sync automatically as the underlying rowstore table changes.
The delta store and the tuple-mover
- Rowgroup. A columnstore index physically slices the table into rowgroups, each holding up to 1,048,576 rows, compressed column-by-column. Bigger rowgroups compress better; too many small ones hurt query performance, so SQL Server also enforces a practical MINIMUM before it will compress a rowgroup: 102,400 rows.
- Delta store (a.k.a. the delta rowgroup(s)). Writing directly into compressed columnar format is expensive for a handful of rows at a time, so SQL Server stages small or incremental writes in the delta store instead: one or more ordinary, uncompressed, B-tree-indexed rowgroups (the docs are explicit that this is a real B+ tree, the same structure a normal clustered index uses, layered specifically for staging). Rows accumulate here until a rowgroup fills to the 1,048,576-row maximum, at which point it flips from OPEN to CLOSED.
- Tuple-mover. A background process that watches for CLOSED delta rowgroups and compresses them into the columnstore as new, permanent COMPRESSED rowgroups. Because it runs asynchronously in the background, a query issued right after a small load may still have to read from BOTH the compressed columnstore and the (uncompressed, B-tree) delta store and merge the results, which is measurably slower than reading from compressed columnstore alone. In SQL Server 2019 and later (and in Azure SQL Database/Managed Instance), the tuple-mover gets help from an additional background merge task that also compresses smaller long-lived OPEN delta rowgroups and merges COMPRESSED rowgroups that have had many rows deleted from them, improving index quality over time without manual intervention.
Structuring micro-batch loads to minimize fragmentation
The single highest-leverage lever is batch SIZE relative to the 102,400-row minimum:
- Batches at or above ~102,400 rows (and ideally close to a full 1,048,576-row rowgroup) can go DIRECTLY into compressed columnar rowgroups during a bulk load, bypassing the delta store entirely. This is the "big enough batch" regime and it produces the least fragmentation.
- Batches below that threshold always land in the delta store, because SQL Server won't compress a group that small on its own; frequent tiny micro-batches therefore generate many small delta rowgroups, all waiting on the tuple-mover, which is exactly the "too many small rowgroups" fragmentation state the docs warn degrades columnstore index quality.
- Practical structuring choices: batch application-level writes up to (or just above) 102,400 rows before committing them to the table when the source system allows it; use bulk-load paths with table-lock hints (
TABLOCK) where possible, since large parallel bulk loads write directly to compressed rowgroups more aggressively than row-by-row inserts; and scheduleALTER INDEX ... REORGANIZE(which forces the tuple-mover's work and merges small COMPRESSED rowgroups, always runs online, at lower resource cost than a full rebuild) during low-traffic windows if micro-batches can't be sized up at the source. ReserveALTER INDEX ... REBUILDfor periodic full defragmentation, since it's more resource-intensive though it can produce a cleaner result. - Tracing one load pattern through the two thresholds. Say a source system delivers a micro-batch of 40,000 rows every 10 minutes, well under the 102,400-row minimum, so every batch lands in the delta store instead of going straight to compressed columnar storage:
- Load 1 (40,000 rows): a new delta rowgroup opens, holding 40,000 rows, OPEN and uncompressed.
- Load 2 (another 40,000 rows, 10 minutes later): appended into the same OPEN delta rowgroup, now 80,000 rows, still OPEN; any query touching this data reads it from the slower, B-tree-indexed delta store, not compressed columnar storage.
- Load 3 (another 40,000 rows, 120,000 total): the rowgroup is now past the 102,400-row minimum, but that minimum only controls whether a bulk-load batch can skip the delta store entirely; it does NOT make an already-open delta rowgroup close. The rowgroup stays OPEN and keeps accepting rows.
- This keeps accumulating, roughly one delta rowgroup per 26 loads (26 x 40,000 = 1,040,000, just under the 1,048,576-row cap), until a load finally pushes it to the 1,048,576-row maximum. Only then does the rowgroup flip from OPEN to CLOSED, at which point the background tuple-mover picks it up, compresses it into a permanent COMPRESSED columnar rowgroup, and a fresh empty delta rowgroup opens for the next load.
- Contrast that with sizing the same source batches up to 102,400+ rows before committing (say, buffering three 40,000-row loads into one 120,000-row batch): each such batch can go straight into a compressed columnar rowgroup during a bulk load, never touching the delta store at all, so there is no lingering OPEN delta rowgroup for a query to fall back to between loads.
- On SQL Server 2019+, Azure SQL Database, Azure SQL Managed Instance (Microsoft's fully-managed, near-100%-compatible SQL Server instance running in Azure), and Synapse dedicated SQL pools (Azure Synapse Analytics' provisioned, large-scale data-warehouse compute tier), the automatic background merge task reduces (but doesn't eliminate) how much manual
REORGANIZE/REBUILDscheduling this needs.
Trade-offs and pitfalls
- Reading from the delta store is slower than reading from compressed columnstore alone, so a table under constant micro-batch load always has SOME fraction of recent data being served at row-store speed rather than columnstore speed; that's expected behavior, not a misconfiguration, but it means "how fresh is recent data" and "how fast are the most recent rows to query" are in tension.
- Batch mode execution, SQL Server's vectorized (multiple-rows-at-once) query processing mode used for columnstore scans, typically improves analytical query performance by roughly 2 to 4x over row-by-row execution; it only engages properly once data has left the delta store and lives in compressed rowgroups, another reason small, lingering delta rowgroups cost you more than their storage footprint suggests.
- Deleting from a compressed columnstore row only marks it logically deleted (physical space isn't reclaimed until a rebuild), while deleting from the delta store removes it immediately; a workload with heavy deletes against old, already-compressed data accumulates dead space that only
REBUILD(or, since 2019, the automatic merge task) reclaims. - A rowgroup that's had all its rows deleted transitions to TOMBSTONE and is cleaned up by the tuple-mover, so a table with a lot of delete churn benefits from the same background processes that manage delta-store compression, not a separate mechanism.
Discuss how secondary indexes are implemented in distributed NewSQL systems like CockroachDB or Spanner. Explain transaction implications, read/write amplification, how index-maintenance is coordinated across nodes, and the effect on latency for multi-region writes.
Sample Answer
Direct answer: In any horizontally-partitioned database, a secondary index is a second, independently-keyed data structure: its rows are ordered by the indexed column, so they almost always land on different partitions than the base row they point back to. Every write to an indexed column therefore has to touch at least two partitions instead of one, and if you want that index to return transactionally correct results, the database has to pay for cross-node coordination on every single write, not just occasionally. CockroachDB and Spanner both choose to pay this cost synchronously (the index write is part of the same distributed transaction as the base-row write), which is exactly what buys you a secondary index whose reads you can trust; systems that instead update indexes asynchronously trade that trust for lower write latency and more availability under partition.
The general cost, before any vendor specifics
This is true of secondary indexes in ANY distributed database, not just CockroachDB or Spanner: maintaining one is expensive because it is a second write path that has to stay correct relative to the first. That cost shows up as three concrete things:
- Write amplification. One logical write to the base row becomes N physical writes: one to the base table, one per secondary index that covers a changed column.
- Cross-node coordination. If the index entry and the base row live in different partitions (the normal case), keeping them consistent needs a distributed transaction, a locking/intent protocol, or an eventual-consistency compromise.
- Read amplification when the index isn't covering. A non-covering index gives you a matching key, not the full row, so a query often pays for an index lookup AND a follow-up fetch of the base row by primary key.
When you'd avoid a secondary index on a distributed table entirely
Given that cost, the honest "don't" cases are: (a) a column that's written far more often than it's queried by, where the write amplification dwarfs the read benefit, particularly under high write throughput per node; (b) a query pattern that can be served by denormalizing the needed columns directly onto the base row or into a separate purpose-built table instead (an async projection, a change-feed into a search system like Elasticsearch/OpenSearch, or a materialized summary table), decoupling the write-hot path from the index-maintenance cost entirely; and (c) a column whose selectivity is so low (e.g., a boolean status flag) that an index barely narrows the scan, so a full or partial scan is cheap enough not to be worth the write cost. In all three cases the fix is the same shape: stop trying to make the base table serve every access pattern, and let a second, purpose-built and separately-scaled structure absorb the query instead of a tightly-coupled secondary index.
CockroachDB and Spanner specifics
Physical layout. In both systems a secondary index is not a pointer structure bolted onto the table; it is itself a distinct, independently-partitioned keyspace, sharded and distributed across nodes exactly like a table is. An index row's key is built from the indexed column(s) (plus, if not unique, the primary key to disambiguate), and its value holds either nothing (a bare index, requiring a lookup back to the primary key for other columns) or the extra columns you name in a STORING clause (CockroachDB) / STORING(GoogleSQL)-INCLUDE(PostgreSQL-dialect) clause (Spanner). A STORING clause turns the index into a covering index for queries that only need the stored columns: the read is satisfied from the index alone, with no extra round trip to fetch the base row, which is the main lever you have for cutting the read-amplification cost above.
Transaction implications. Both systems fold index writes into the SAME distributed transaction as the base-row write. CockroachDB calls its partition unit a "range" and Spanner calls its "split"; both are just the vendor-specific name for the same idea this answer has been calling a partition, the physical chunk of sorted data that gets replicated and moved as one piece. In CockroachDB, each table and each index lives in one or more Raft-replicated ranges (Raft is a consensus protocol: a way for multiple replicas holding the same data to agree on the order of writes and elect a leader, even when some replicas are slow or unreachable), and a multi-range write goes through the transaction coordinator, which writes provisional "intents" across every affected range and only resolves them atomically at commit; the underlying per-range consensus is Raft. In Spanner, commits across the ranges (splits) touched by a transaction use two-phase commit (2PC: the coordinator first asks every participant "can you commit this?" and only tells them to actually commit once every participant has said yes, so a transaction spanning multiple partitions either commits everywhere or nowhere) coordinated with TrueTime, Google's globally-synchronized clock service, which lets Spanner assign a commit timestamp that every replica can agree on without extra communication just to establish ordering. Either way, correctness costs coordination: you cannot get a transactionally-consistent secondary index in a distributed table for free.
Coordination across nodes and write amplification. Because base row and index row are usually on different ranges/splits (unless deliberately colocated), a single-row UPDATE on an indexed column becomes at minimum two range writes, more if multiple indexes cover the changed column, each requiring intent-write-then-resolve (CockroachDB) or inclusion in the same 2PC commit (Spanner). Spanner has a specific mitigation for this: INTERLEAVE IN on the index lets you colocate the index's rows with an ancestor table's rows in the same split, so that when the index is interleaved in an ancestor of the base table, the index write becomes a LOCAL operation instead of a cross-split one, avoiding the two-phase commit overhead entirely for that write.
Multi-region write latency, and the eventual-consistency trade-off. When the ranges or splits touched by a transaction (base row plus every indexed column's index entry) span regions, the commit has to wait on cross-region round trips and quorum acknowledgment from replicas in those regions, regardless of how fast any single node is. Spanner's TrueTime bounds the UNCERTAINTY in that commit's timestamp, it does not remove the network latency of reaching a remote replica. CockroachDB can place replicas and leaseholders (the one replica of a range currently responsible for serving that range's reads and writes) to keep hot ranges local to where they're written most, but a write that must update an index range in another region still pays that region's round-trip time. This is the real lever for the "multi-region low-latency reads" trade-off: both systems let you trade strict, transactionally-fresh secondary-index reads for latency by reading with bounded/stale staleness (Spanner's stale reads, CockroachDB's AS OF SYSTEM TIME / follower reads) from a nearby replica instead of the current leaseholder. That gets you a fast local read of a secondary index at the cost of it possibly being a few hundred milliseconds to a few seconds out of date, an explicit eventual-consistency trade you opt into per query rather than a property of the index itself.
Trade-offs and pitfalls
- Minimize the number of secondary indexes on any table with a high write rate; each one is a proportional tax on every write to a covered column, not a one-time cost.
- Prefer covering (
STORING) indexes for read-heavy queries specifically to avoid the extra primary-key fetch, but remember every stored column also has to be kept in sync on every write, so storing everything defeats the purpose. - Use locality (Spanner interleaving, CockroachDB replica/leaseholder placement) to keep an index's writes local to where its base rows are written, before reaching for staleness as the fix.
- For extremely high write throughput where even synchronous, well-placed index maintenance is too expensive, the honest answer is often "don't index it in the primary OLTP store": stream changes out (CDC, a change feed) to an external, independently-scaled search or analytics system instead, and accept that its results are eventually consistent by design rather than accidentally so.
Explain the different index types relevant to analytical systems: B-tree, bitmap, inverted, and zone map indexes. For each index type, describe how it works at a high level, what query patterns it accelerates, and which storage engines (row or columnar) make best use of it.
Sample Answer
Direct answer
Four structures, four different jobs. B-tree: a balanced tree keyed by sorted value, general-purpose, equally at home in a row store or a columnar engine, best for high-cardinality equality and range lookups. Bitmap: one bit per row per distinct value, best for low-cardinality columns where multiple predicates get combined with bitwise AND/OR, the natural fit for a columnar warehouse's dimension columns. Inverted index: a mapping from each distinct TOKEN (a word, a tag, a substring) to the list of rows containing it, the structure behind full-text and log search, whether that's Postgres's GIN (Generalized Inverted iNdex) over tsvector, or Elasticsearch's core data structure. Zone map (in Postgres, a BRIN, Block Range INdex): a tiny summary, typically a min/max, kept per block RANGE rather than per row, cheap to store and cheap to consult, valuable only when the indexed column stays correlated with physical storage order.
Structured elaboration
| Index type | Mechanism | Query pattern accelerated | Row store or columnar |
|---|---|---|---|
| B-tree | sorted tree, one entry per row | equality, range, sort, on any cardinality | both, universal |
| Bitmap | one bit per row per distinct value | multi-predicate AND/OR on low-cardinality columns | both, but shines in columnar/OLAP (online analytical processing) |
| Inverted | token -> row-list mapping | full-text search, tag/substring search | both; the standard structure for dedicated search engines |
| Zone map / BRIN | min/max summary per block range | range filters on a naturally-ordered column | columnar engines build this in natively; Postgres offers it as an explicit index type on any row-store table |
Worked example
Three of the four demonstrated directly on Postgres 16, real numbers, real plans.
Bitmap-style execution (Postgres builds bitmaps dynamically at query time from ordinary B-tree indexes, rather than storing a persistent bitmap index type, but the execution pattern is the same bitwise mechanism): on a 1,000,000-row fact table with a 4-value region column and a 2-value gender column, EXPLAIN on region = 'eu' AND gender = 'F' showed a Bitmap Heap Scan using only the region index and applying gender as an in-memory filter, the planner correctly judging gender's 50 percent selectivity wasn't worth a second bitmap probe to combine. Sizing the two approaches: a measured B-tree on the 2-value gender column costs 6,792 KiB; a textbook bitmap (1 bit per row per value) would cost about 244 KiB for the same column, 28x smaller, exactly the storage argument for using bitmap-style structures on genuinely low-cardinality analytical dimension columns.
Zone map (BRIN), on a 3,000,000-row append-style log table with a timestamp column inserted in true chronological order:
CREATE INDEX idx_audit_brin ON audit_log USING BRIN (logged_at);
-- pg_relation_size: 24 kB (versus 64 MB for an ordinary B-tree on the same column)
EXPLAIN (ANALYZE, BUFFERS)
SELECT id FROM audit_log WHERE logged_at >= '2024-01-15' AND logged_at < '2024-01-16';
-- Bitmap Heap Scan using idx_audit_brin, Execution Time: 9.426 ms
A 24 KB index (about 2,700 times smaller than the equivalent B-tree) answering a one-day range filter out of a 270-day span in single-digit milliseconds, competitive with the much larger B-tree, because the column stayed correlated with insert order: BRIN doesn't need to know exactly which rows match, only which block RANGES could possibly contain a match, and it skips the rest.
Inverted index, via Postgres's GIN over a tsvector (a preprocessed, tokenized representation of text, used for full-text search): on a 500,000-row log-message table:
CREATE INDEX idx_log_lines_fts ON log_lines USING GIN (to_tsvector('english', message));
-- pg_relation_size: 2,776 kB (versus 40 MB for the table itself)
EXPLAIN (ANALYZE)
SELECT id FROM log_lines
WHERE to_tsvector('english', message) @@ to_tsquery('english', 'timeout & upstream');
-- unindexed: Parallel Seq Scan (re-tokenizing every row live), Execution Time: 505.366 ms
-- indexed: Bitmap Heap Scan using idx_log_lines_fts, Execution Time: 17.469 ms
About 29x faster once the token-to-row mapping is precomputed and indexed, rather than re-tokenizing every row's text on every query, and the index itself is under 7 percent of the table's size.
Trade-offs & pitfalls
Each structure's win condition is narrow, and using the wrong one is worse than using none, not just less optimal: a BRIN index on a column that ISN'T correlated with physical order (a randomly-scattered category column, for instance) provides almost no pruning at all, since nearly every block range ends up containing nearly every value, while still costing something to build and consult. A bitmap-style approach on a high-cardinality column inverts the storage argument entirely: run the arithmetic for a million-row table with a customer-identifier column carrying about 500,000 distinct values, and a textbook one-bit-per-row-per-value bitmap costs 1,000,000 x 500,000 / 8 bytes, on the order of 58 GiB, against a plain B-tree on the same column that typically costs a few tens of megabytes. That inversion, cheap at low cardinality, ruinous at high cardinality, is why bitmap-style structures stay scoped to dimension-shaped columns. An inverted index built over a column that's rarely searched by substring or token match is pure write and storage overhead for no read benefit, since a B-tree already answers exact-match and prefix queries more cheaply.
Storage engine matters here too: a genuine columnar engine (Redshift, BigQuery, Snowflake, ClickHouse) builds zone maps and often bitmap-style encodings into its storage format automatically, as a property of how columns are physically laid out and compressed, not as a separate index you opt into; a row store like Postgres makes you choose and build each of these structures explicitly, which is more manual work but also more precise control over exactly which columns pay which costs. Knowing which of these four structures fits a given column, rather than defaulting to "just add a B-tree," is the senior-level judgment this question is really testing: it comes down to the column's cardinality, whether it's naturally correlated with physical order, and whether the query pattern is exact-match, range, or token search.
You're building on a NoSQL key-value store that has no native secondary-index feature, and the application needs to look up records efficiently by an attribute other than the primary key. Design an indexing approach that supports this at scale, and explain the consistency, write-amplification, and hot-key risks your design introduces as write volume grows.
Sample Answer
Direct answer: Build the index as its own key-value (KV) collection: a second set of items keyed by the attribute you want to look up by (or a hash of it), whose value holds the matching primary key(s). Every write to the base item becomes two writes, one to the base item and one to its index entry, and the central design decisions are (a) whether those two writes happen synchronously in the request path or asynchronously via a change log, and (b) how you prevent a popular attribute value from concentrating too many writes onto one physical index partition.
The pattern: an inverted index as a second collection
For a base item stored at key item:<id> and an attribute email you need to look up by:
- Index collection key:
idx:email:<email_value> - Index collection value: the set of matching primary keys (usually just one, if
emailis meant to be unique; a list/set if not)
A lookup by email becomes: read idx:email:<email_value> to get the primary key(s), then read each item:<id> by primary key, the store's one native fast operation. This is exactly the shape DynamoDB's own Global Secondary Indexes and Cassandra's native 2i implement internally (a separate keyed structure that maps indexed-value to primary key); the difference here is you're building and maintaining that structure yourself because the underlying store doesn't offer it.
Keeping the two writes consistent
- Synchronous dual write (in the request path). The application writes the index entry and the base item as part of handling the request, ideally through whatever the store offers closest to a transaction (a conditional/compare-and-swap write, or a native multi-key transaction if the store has one). This gives read-your-own-write consistency immediately, at the cost of coupling the base write's availability and latency to the index write's: if the index write fails after the base write succeeds (or a process crashes between the two), the index silently drifts out of sync with no automatic recovery.
- Asynchronous propagation (via a change log/stream). The application writes only the base item; a separate consumer reads a change log (an append-only stream of "item X changed" events, the same idea as change-data-capture, CDC) and applies the corresponding index write afterward, idempotently (safe to re-apply the same change twice without corrupting state, usually by making the index write a last-write-wins upsert keyed by a version or timestamp). This decouples the base write's latency and availability from index maintenance, and survives partial failures by simply retrying the log entry, at the cost of a real consistency lag: a read immediately after a write may not see it reflected in the index yet.
- Reconciliation. Whichever approach you pick, run a periodic consistency-check job that compares a sample of base items against their expected index entries and repairs drift; the synchronous approach needs this for the "crashed mid-write" case, and the asynchronous approach needs it as a standing safety net against a stuck or lossy consumer.
Consistency, write amplification, and hot-key risk as write volume grows
- Consistency: synchronous dual-write gives strong-ish consistency (bounded by whatever atomicity the store's conditional write actually guarantees, which for a pure KV store is usually "one key at a time," so a true two-key atomic write may not exist and you're really doing best-effort ordering plus reconciliation). Asynchronous propagation is explicitly eventually consistent, and the lag grows under load exactly when you most need the index to be fresh (a traffic spike).
- Write amplification: every base write becomes at least 2 physical writes (1 base + 1 index), growing linearly with the number of indexed attributes; a base item indexed on 3 attributes is 4 writes per logical update. This tax is unavoidable, since it's the same fundamental cost every managed secondary-index feature pays; the only choice is whether YOUR system or the DATABASE's internals absorb it.
- Hot-key risk: if the indexed attribute has low cardinality (a
statusfield with 3 values, a country code with a handful of common values), a large share of index writes and reads collapse onto a small number of index keys, each living on one physical partition. Concrete point of reference: AWS documents DynamoDB's own per-partition ceiling at 3,000 read-capacity-units/second and 1,000 write-capacity-units/second; the same physical constraint (one partition, one set of disks/CPU, a hard ceiling) exists in any partition-per-key infrastructure, including a hand-rolled one, even if the exact numbers differ. Push enough writes at one index key and you hit that ceiling long before the rest of the cluster is under any real load.
Mitigation as write volume grows
- Write (key) sharding. Append a bucket suffix, random or computed (e.g., a hash of the primary key modulo N), to hot index keys:
idx:status:pending#7instead ofidx:status:pending. Writes spread across N physical index partitions instead of one; reads for "all pending items" now have to fan out across N buckets and merge results, trading read simplicity for write scalability. Choose N based on the ceiling you're targeting divided by expected peak write rate for that value, not an arbitrary round number. - Bounded fan-out per index entry. If an index key's value is a growing list (many primary keys per indexed value), cap and paginate it rather than letting one index item grow unbounded; an unbounded list under one key is itself a form of the same hot-key problem, now expressed as item size instead of throughput.
- Skip indexing what doesn't need it. If an attribute is low-cardinality enough that the mitigation above becomes the dominant complexity, a full or filtered scan may genuinely be cheaper than a hand-built index that needs its own sharding scheme to survive; don't build the index reflexively just because the requirement says "look up by X."
Trade-offs and pitfalls
- This pattern is exactly what a managed secondary index gives you for free, plus the operational burden of writing and running the propagation and reconciliation logic yourself; before building it, confirm the store really has no native option (some "pure KV" stores add limited secondary-index or search features over time).
- A synchronous dual write without a real cross-key transaction is not actually atomic; be explicit with your team about what failure mode you're accepting (a crash between the two writes leaves a stale or missing index entry) and build the reconciliation job as a first-class part of the design, not an afterthought.
- Async propagation needs idempotent application logic; without it, a replayed log entry after a retry can apply an update twice or apply events out of order, corrupting the index in a way that's hard to notice until a lookup returns a stale or wrong primary key.
Unlock Full Question Bank
Get access to all 6 Indexing Strategy and Design interview questions and detailed answers.
Sign in to ContinueJoin thousands of developers preparing for their dream job.