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.
For recurring analytics that compute ROW_NUMBER() OVER (PARTITION BY customer_id ORDER BY event_ts DESC) on a 2B row events table, propose index and partitioning strategies across OLAP systems to speed queries, and discuss trade-offs such as insert throughput vs query latency and maintenance costs.
Sample Answer
Direct answer: For "give me the latest N events per customer" run recurrently over a 2-billion-row table, the fastest and cheapest long-term fix is usually to stop computing it from raw events at query time at all: maintain a small, incrementally-updated "latest per customer" summary table, and keep the raw events table partitioned by time with a composite (customer_id, event_ts DESC) clustering key purely to make ad hoc/audit queries and the summary's own refresh job cheap. Where a system can't or shouldn't maintain that summary table (irregular queries, exploratory analytics), the fallback is engine-appropriate physical ordering: partition by date for prune-and-isolate, and sort/cluster by (customer_id, event_ts DESC) so each customer's most recent rows are physically contiguous, at the cost of slower, more resource-intensive loads.
Why the ROW_NUMBER() shape forces this discussion
ROW_NUMBER() OVER (PARTITION BY customer_id ORDER BY event_ts DESC) computed over the whole table means the engine has to, for every distinct customer_id, sort that customer's rows by event_ts before it can number them, even though the query usually only wants the top 1-few rows per customer. Without physical layout support, this is a full sort of up to 2 billion rows partitioned into however many customer_id groups exist. The goal of indexing/partitioning here is to avoid ever materializing that full sort.
Strategy, by layer
- Time partitioning (all systems). Partition the raw events table by
event_ts(day or month, depending on retention and query granularity). This isolates each load to its own partition (so ingest doesn't contend with older partitions) and lets any query with a time bound prune partitions before scanning. - Physical ordering within partitions. Cluster/sort by
(customer_id, event_ts DESC):- Postgres:
CREATE INDEX ON events_p2026_09 (customer_id, event_ts DESC);on each partition; combined with partition pruning, an index-only or index scan can find a given customer's most recent rows without a table-wide sort. - Redshift:
SORTKEY (customer_id, event_ts)withDISTKEYon a high-cardinality column (oftencustomer_iditself, or a hash of it) to spread load rather than concentrate one customer's writes on one slice. Redshift'sSORTKEYclause takes only column names, not a per-column ASC/DESC direction, so this sort key does not physically store a customer's rows in descending time order; it only keeps one customer's rows contiguous, and the query's ownORDER BY ... DESCstill does the reversing at read time. - BigQuery:
PARTITION BY DATE(event_ts) CLUSTER BY customer_id. No native secondary sort within a cluster block, but clustering colocates one customer's rows within the same storage blocks. - Snowflake: time-ordered micro-partitions plus an explicit clustering key
(customer_id, event_ts); Snowflake'sCLUSTER BYclause likewise takes only column expressions, no per-column sort direction, so the same caveat as Redshift's sort key above applies: clustering co-locates a customer's rows within the same micro-partitions, it does not pre-sort them descending. Watch the automatic-reclustering credit cost on a table this large and this write-heavy. - A dedicated row-per-key engine like ClickHouse (mentioned for completeness, not assumed to be in scope):
ORDER BY (customer_id, event_ts DESC)in a MergeTree table makes this exact "latest-N per key" query close to free, since the physical sort order IS the query's required order.
- Postgres:
- The materialized "latest per customer" table. For the recurring case (this query runs on a schedule, not ad hoc), maintain a small table
(customer_id PK, event_ts, ...payload)updated incrementally via change-data-capture (CDC, streaming the events table's inserts) or a periodicMERGE/upsert job. Queries against it are simple primary-key lookups, no window function, no per-customer sort, at query time.
Worked sizing (assumptions stated; this is an ESTIMATE, not a measurement, and is meant to show the shape of the trade-off, not a guaranteed number for any specific schema)
rows = 2_000_000_000
avg_row_bytes = 150 # customer_id bigint(8) + event_ts timestamptz(8) + a small enum(2)
# + ~16B of numeric/enum payload + ~24B row/page overhead + a short string(~40B) + padding
raw_bytes = rows * avg_row_bytes
print(f"raw table size estimate: {raw_bytes/1e9:.1f} GB")
distinct_customers = 50_000_000
idx_entry_bytes = 24 # customer_id(8) + event_ts(8) + row pointer/overhead(~8)
idx_bytes = rows * idx_entry_bytes
print(f"(customer_id, event_ts DESC) index estimate: {idx_bytes/1e9:.1f} GB")
latest_table_bytes = distinct_customers * avg_row_bytes
print(f"latest-per-customer summary table estimate: {latest_table_bytes/1e9:.2f} GB "
f"({latest_table_bytes/raw_bytes*100:.3f}% of raw table size)")
Output:
raw table size estimate: 300.0 GB
(customer_id, event_ts DESC) index estimate: 48.0 GB
latest-per-customer summary table estimate: 7.50 GB (2.500% of raw table size)
Even under generous per-row assumptions, a maintained summary table lands about 40x smaller than the raw events table (7.5 GB versus 300 GB) and about 6x smaller than the composite index (7.5 GB versus 48 GB); that is a real, substantial gap, but it is not the two full orders of magnitude (100x) that the 2-billion-events-versus-50-million-customers row-count ratio alone might suggest, since each summary row is the same width as a raw event row rather than a smaller aggregate. That gap is still the entire argument for maintaining it: the summary table's storage, backup, and query cost stay flat as the raw table grows, while a per-query full window-function pass over the raw table grows linearly with total row count regardless of how many customers actually queried it.
Trade-offs: insert throughput vs. query latency vs. maintenance
| Approach | Insert throughput | Query latency | Maintenance cost |
|---|---|---|---|
No physical ordering, plain ROW_NUMBER() scan | Best: pure append | Worst: full sort per partition-by group, scales with total row count | None |
| Time partition only | Good: isolated per-partition writes | Better for time-bounded queries, still a customer-level sort within scanned partitions | Low: partition creation/retention only |
| Time partition + (customer_id, event_ts DESC) sort/cluster/index | Worse: every insert pays to maintain physical order or a composite index | Best of the "query raw data" options: engine can walk pre-ordered data instead of sorting | Medium/high: index maintenance (Postgres), reclustering credits (Snowflake), or interleaved-write cost (Redshift SORTKEY) |
| Materialized latest-per-customer table, refreshed via CDC/upsert | Adds a second write path (the refresh job), but decoupled from the main ingest path's throughput | Best possible: primary-key lookup, no scan or sort at all | Requires a reliable incremental refresh pipeline (CDC lag, upsert correctness) as new maintenance surface |
Trade-offs and pitfalls
- Sorting/clustering at write time is a real tax on insert throughput; systems optimized for fast bulk loads (BigQuery, ClickHouse) absorb it more gracefully than Postgres, where many indexes plus frequent updates also mean more autovacuum work.
- A materialized "latest" table's incremental refresh (CDC or scheduled
MERGE) becomes new maintenance surface: correctness now depends on the refresh job never silently falling behind, so it needs its own lag/freshness monitoring, not just the base table's health checks. - Replacing the window function with a join to a
GROUP BY customer_idmax-timestamp subquery (... JOIN (SELECT customer_id, MAX(event_ts) FROM events GROUP BY customer_id) USING (customer_id, event_ts)) can be cheaper thanROW_NUMBER()for the "top 1" special case specifically, since it lets the optimizer use the same(customer_id, event_ts DESC)structure as a straightforward aggregate rather than a window frame; it stops being a shortcut once you need top-N for N > 1. - Watch re-clustering/rebuild costs as ongoing OPERATING expense, not a one-time migration cost: Snowflake's automatic clustering keeps consuming credits as new data arrives out of clustering order, and Postgres index maintenance on a 2-billion-row, constantly-inserted table needs a real vacuum/reindex strategy, not defaults.
For columnar analytical databases (e.g., Redshift, BigQuery, Snowflake), indexing concepts differ from row stores. Explain how sortkeys, clustering, and partitioning map to traditional index concepts and describe a migration plan from a row-store with many indexes to a columnar store while preserving dashboard performance.
Sample Answer
Direct answer: Sort keys, clustering, and micro-partition pruning do the same two jobs a row-store index does, ordering and skipping, but they do them at the level of storage blocks (megabytes of rows at a time), not individual rows, and none of the three engines maintain a real B-tree for point lookups. Migrating off a row-store with many indexes means deliberately giving up cheap single-row lookups in exchange for scan-pruning efficiency on the big aggregate queries dashboards actually run, and then compensating any point-lookup-shaped access pattern with a narrow side table or a materialized rollup instead of an index.
How the primitives map (and where the mapping breaks)
| Row-store concept | Redshift | BigQuery | Snowflake | What it's actually doing |
|---|---|---|---|---|
| Ordered index, for range scans | Sort key (COMPOUND by default, or INTERLEAVED, up to 8 columns) | Clustering columns (up to 4, ordered) | Clustering key | Places rows so a range predicate on the leading column(s) touches contiguous storage instead of scattered blocks |
| Coarse block-skipping index | Zone maps: automatic min/max per 1 MB disk block | Automatic block-level metadata and pruning on clustered columns | Micro-partition metadata: automatic min/max plus distinct-value counts per 50-500 MB micro-partition | Lets the scanner discard whole blocks whose min/max can't satisfy the predicate, with no B-tree traversal at all |
| Table-level partition pruning | Not a first-class DDL concept; usually folded into the sort key plus zone maps, or emulated with date-suffixed tables | Native PARTITION BY on a date/timestamp or integer range; a partition is physically separate storage | Not a DDL concept; time-ordered loads plus a clustering key achieve the same skip behavior | Prunes entire chunks of storage before scanning even begins |
| Point-lookup index (B-tree) | None | None as a general mechanism (BigQuery's separate search indexes target exact/substring match on string fields, a different feature) | None by default; the optional search optimization service adds a purpose-built point-lookup structure on specific columns, at extra credit cost | Not present in any of the three out of the box; you either pay for an add-on feature or build a narrow side table |
An AWS-documented illustration of what zone-map pruning buys you: a table holding five years of data sorted by date, queried with a one-month predicate, can skip up to 98% of its disk blocks, because contiguous storage plus per-block min/max means the scanner never opens blocks outside the date range. That is Amazon's own worked example for Redshift zone maps, not a number I measured; it's cited to show the mechanism's shape, not to be treated as a guaranteed ratio for every workload.
Migration plan: row-store with many indexes to columnar, preserving dashboard performance
- Inventory before touching anything. Capture the actual queries dashboards run: filter predicates,
GROUP BYkeys, join patterns, and each dashboard's latency budget (its service-level objective, SLO). Tag each filtered column by cardinality (few distinct values vs. many) and by whether it's typically used as a range or an equality predicate. This inventory, not intuition, drives every choice below. - Assign each access pattern to the right primitive.
- Heavy time-range dashboards: partition by day/month (BigQuery) or lean on sort-key/clustering ordering by timestamp (Redshift/Snowflake) so the date filter prunes blocks before anything else runs.
- Frequent equality filters on low-cardinality dimensions (country, plan tier): these benefit from block-level min/max pruning almost for free once the table is loaded in roughly sorted order; putting a low-cardinality column first in a compound sort key or clustering list wastes its ordering power, since equality on that column alone doesn't narrow much.
- High-cardinality point lookups (a specific order ID, a specific user's single row): none of the three engines make this cheap through indexing primitives. Route it to a narrow, purpose-built lookup table (a small table keyed by that ID, refreshed incrementally) or a real OLTP database kept alongside the warehouse; do not try to force the warehouse to do a row-store's job.
- Load historical data with the target physical layout already in place. Partition/cluster/sort at load time rather than loading flat and re-clustering after the fact; a full re-cluster of billions of already-loaded rows is exactly the "maintenance cost" this migration is trying to control.
- Pre-aggregate the expensive rollups dashboards recompute on every refresh. Hourly/daily materialized aggregates (a small summary table maintained on a schedule or via
MERGE) turn a scan-heavy dashboard query into a lookup against a table that's orders of magnitude smaller than the raw fact table. - Validate against the SAME captured queries from step 1, comparing
EXPLAIN/query-profile output and bytes-scanned, not just wall-clock time (which is environment-dependent); confirm pruning is actually engaging (e.g., scanned-bytes should shrink roughly in proportion to how selective the date/cluster predicate is) before declaring the migration done. - Dual-write and compare before cutover. Run both systems for a defined window, compare dashboard outputs and latency against the row-store baseline, and keep a rollback path (reverting dashboards to the row-store) until the columnar system has cleared its SLOs for a full business cycle (e.g., a month, to catch periodic reporting jobs).
Trade-offs and pitfalls
- Columnar storage is a poor fit for many small random single-row lookups; all three engines compress data well and scan large ranges fast, but each pays a much higher relative cost per point lookup than a row-store B-tree does. Denormalize or pre-aggregate rather than fighting this.
- Sort/clustering key order matters: put the most selective and most frequently range-filtered column first. A compound sort key or clustering list ordered wrong (e.g., leading with a rarely-filtered column) gets almost none of the pruning benefit even though the DDL "has" a sort key.
- Over-partitioning (too many small partitions, common in BigQuery
PARTITION BYmisuse on high-cardinality columns) and over-clustering both create their own overhead: BigQuery must manage more partition metadata, Snowflake's automatic reclustering consumes credits proportional to how much it has to reshuffle, Redshift'sVACUUM REINDEXon an interleaved sort key takes an extra analysis pass and can run materially longer than a plainVACUUM FULL. - Instrument the migration with the engines' own pruning telemetry (bytes scanned, query profile block/partition counts) rather than trusting the DDL alone; a sort key or clustering key that theoretically matches the query pattern can still fail to prune if load order drifted from it over time (Redshift specifically recommends
SORTKEY AUTOtoday precisely so the optimizer, not a one-time DDL choice, keeps the physical layout matched to the actual query mix).
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.
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.
Compare bitmap indexes and B-tree indexes for high-cardinality columns in a data warehouse. Explain storage footprints, query execution patterns (bitwise operations vs B-tree seeks), and update/maintenance characteristics.
Sample Answer
Direct answer
For a genuinely high-cardinality column, a B-tree index is the right choice: it stays compact (roughly proportional to row count, not to distinct-value count) and supports equality, range and sort operations. A bitmap index goes the other way: it is extremely compact and fast for low-cardinality columns (few distinct values, each common), but its storage cost grows with the number of distinct values, so on a high-cardinality column it can end up far bigger than a B-tree and its bitwise merge advantage disappears. Bitmap indexes are a data-warehouse (OLAP, online analytical processing) tool for dimension columns like region, gender or status, not for keys or high-cardinality attributes like customer_id or email.
Structured elaboration
Storage footprint. A B-tree stores one entry per row (key value plus a pointer to the row), so its size scales with row count and key width, largely independent of how many distinct values exist. A bitmap index instead keeps one bitmap per distinct value, one bit per row, marking whether that row has that value. Its size scales with distinct_values x row_count / 8 bytes. That formula is why the two structures trade places as cardinality rises: at low cardinality the bitmap is tiny; at high cardinality it explodes.
Query execution pattern. A B-tree seek walks down the tree (root -> internal nodes -> leaf) in O(log n) comparisons (notation for "grows very slowly as the table gets bigger": doubling the row count adds roughly one more comparison, not twice the work) to find a value or a range boundary, then follows leaf pointers, or the heap, to fetch rows. Combining two B-tree-indexed conditions with AND/OR either narrows to one B-tree scan and filters in memory, or, when the database supports it, builds a temporary bitmap of matching row positions from each index and combines them with fast bitwise AND/OR/NOT before touching the table. A dedicated bitmap index skips the "build a bitmap on the fly" step because the bitmap already IS the index: combining predicates on several bitmap-indexed columns is a pure bitwise operation over the stored bitmaps, no tree traversal at all, which is why bitmap indexes shine on multi-predicate analytical filters (WHERE region = 'eu' AND channel = 'web' AND status = 'active').
Update and maintenance characteristics. A B-tree update touches one leaf page (and occasionally splits it); cost is O(log n) and localized. A bitmap index update has to touch the bitmap for the row's OLD value (clear a bit) and the bitmap for its NEW value (set a bit), and because bitmaps are typically stored compressed (run-length encoded, since a mostly-0/mostly-1 bitmap compresses very well), a single-row update can force re-encoding a run of the bitmap, which is comparatively expensive and often involves row-level or page-level locking that serializes concurrent writers more than a B-tree does. This is the core reason bitmap indexes are steered toward read-heavy, batch-loaded warehouse tables and avoided on frequently-updated OLTP (online transaction processing) tables.
| B-tree | Bitmap index | |
|---|---|---|
| Storage scales with | row count | distinct_values x row_count |
| Best cardinality | any, especially high | low (few, common, distinct values) |
| Multi-predicate AND/OR | tree scan + in-memory filter, or a dynamically built bitmap | native bitwise AND/OR/NOT on stored bitmaps |
| Point update cost | O(log n), localized to one leaf | re-encode a compressed run, more contention |
| Typical home | OLTP and any selective lookup | OLAP dimension columns |
Worked example
Consider a 1,000,000-row fact table with a gender column (2 distinct values) and a region column (4 distinct values), next to a customer_id column with about 500,000 distinct values. Measured on Postgres 16 (which builds its indexes as B-trees and only forms bitmaps dynamically at query time, so these are B-tree sizes, but the arithmetic for a true stored bitmap generalizes across engines):
CREATE TABLE fact_sales (
id bigserial PRIMARY KEY,
region text NOT NULL,
gender text NOT NULL,
customer_id bigint NOT NULL,
amount_cents bigint NOT NULL
);
-- 1,000,000 rows loaded, region uniform over 4 values, gender uniform over 2
CREATE INDEX idx_fact_sales_gender ON fact_sales (gender);
CREATE INDEX idx_fact_sales_region ON fact_sales (region);
CREATE INDEX idx_fact_sales_customer ON fact_sales (customer_id);
Measured B-tree sizes (pg_relation_size): gender index 6,792 KiB, region index 6,792 KiB, customer_id index 16 MiB, table itself 61 MiB.
A textbook bitmap index would instead cost rows x distinct_values / 8 bytes:
rows = 1_000_000
gender_bitmap_kib = rows * 2 / 8 / 1024 # 2 distinct values
region_bitmap_kib = rows * 4 / 8 / 1024 # 4 distinct values
customer_bitmap_gib = rows * 500_000 / 8 / 1024**3 # 500k distinct values
print(gender_bitmap_kib, region_bitmap_kib, customer_bitmap_gib)
# 244.14... 488.28... 58.2...
That gives a bitmap estimate of about 244 KiB for gender (versus a measured 6,792 KiB B-tree: the bitmap is about 28x smaller) and about 488 KiB for region (about 14x smaller than its B-tree). For customer_id the picture inverts completely: a bitmap would cost roughly 58 GiB against a measured 16 MiB B-tree, about 3,700x LARGER. That is the whole cardinality argument in one number.
On execution pattern: running EXPLAIN (ANALYZE, BUFFERS) on region = 'eu' AND gender = 'F' against the B-tree indexes above, Postgres chose a Bitmap Heap Scan using only the region index (cost 2786.43..14439.93, actual time 4.171..68.761 ms, 125,209 rows), applying gender = 'F' as an in-memory filter on the fetched rows rather than also probing the gender index and combining bitmaps. That is the planner correctly judging that gender's 50 percent selectivity makes a second index probe not worth it, a live illustration that even a dynamic bitmap-combination engine only pays the bitwise-AND cost when it is worth paying, exactly the same selectivity judgment that decides whether a genuine bitmap index earns its keep on a given column.
Trade-offs & pitfalls
The in-memory-bitset-inside-a-BI-engine angle is the same idea one layer up: many BI (business intelligence) and OLAP engines keep a compact in-memory bitset per low-cardinality column purely to accelerate ad hoc filtering during a session, without persisting it as a database index at all. It has the same cardinality ceiling (a bitset over a million-row column with a handful of distinct values is cheap; over a high-cardinality column it is not) and the same update cost profile (fine for a batch-refreshed cube, bad for a live-updated one).
A classic worked low-cardinality case is a gender or region column: 2 to 10 distinct values, each covering a large, roughly even fraction of rows, which is exactly the shape bitmap indexes are built for and exactly the shape a B-tree wastes space and selectivity on (a B-tree seek on a 50 percent-selective predicate barely beats a sequential scan). The common mistake is reaching for a bitmap index (or, in Postgres, assuming a plain B-tree will behave like one) on a column like customer_id, order_id or email: the storage blows up, and worse, if the table is being written to concurrently, bitmap-index maintenance contention can serialize writers in ways a B-tree never would.
Unlock Full Question Bank
Get access to all 10 Indexing Strategy and Design interview questions and detailed answers.
Sign in to ContinueJoin thousands of developers preparing for their dream job.