skip to content

Two queries each match about 50,000 rows through secondary indexes on the same table and each needs columns the index does not hold. One finishes in milliseconds, the other takes many seconds in row lookups. What property of the data explains the difference, and how would you confirm it?

level: seniorimportance: should knowfreq 36%

answer

  1. same rows, different page spread
  2. correlation = insert behaviour, not index definition
  3. 500 pages vs 50,000 pages
  4. clustering factor statistic
  5. cold vs warm run reveals scatter

basics

~20 s

Correlation between the index key order and the physical order of rows. Well-correlated keys make consecutive lookups land on the same few pages, so they are cache hits or sequential reads; uncorrelated keys scatter 50,000 lookups across 50,000 different pages, each a random read. Confirm with the clustering statistic and buffer hit counts.

solid answer

~60 s

The index search is the same in both cases; the difference is entirely in the row lookups. If the indexed column correlates with physical row order - a creation timestamp on an append-only table, or a column that mirrors the clustered key - then successive index entries point into the same or adjacent pages. Fifty thousand rows at a hundred rows per page become about five hundred page visits, mostly cached and sequential. If the column is uncorrelated - a status flag, a random identifier, a foreign key to a hot parent - those fifty thousand entries point at up to fifty thousand distinct pages, each a random read, and the working set may exceed the cache so repeats do not help either. I confirm by reading the plan with actual timings plus buffer hit and read counts for the lookup node, and by checking the engine's clustering or correlation statistic for the index. Remedies are physically reordering the table by that key, partitioning so the relevant rows are colocated, or removing the lookup by putting the needed columns in the index - each with its own cost.

code

text · 7 lines
text
-- well-correlated key (rows inserted in this order)
Index Scan using idx_orders_created_at   rows=50,000  time=41 ms
  buffers: shared hit=612 read=88

-- uncorrelated key (values interleaved throughout the table)
Index Scan using idx_orders_status       rows=50,000  time=7,340 ms
  buffers: shared hit=9,110 read=41,255

go deeper

for a junior

Say that if the matching rows happen to sit close together in the table the lookups read few pages, whereas scattered rows mean one read each.

for a middle

Do the arithmetic with rows per page, name clustering or correlation as the statistic, and connect it to how rows were inserted.

for a senior

Diagnose with per-node buffer hit and read counts plus cold-versus-warm comparison, then weigh reordering, partitioning, and removing the lookup with their real costs.

for a principal

Treat physical layout as a design lever: which access path deserves colocation, whether partitioning is justified system-wide, and how correlation decay is monitored as data ages.

## Same row count, different cost Row lookup cost is not counted in rows, it is counted in **distinct pages touched and how they are reached**. Two queries returning the same number of rows can differ by two orders of magnitude because of how those rows are spread across pages. - **Well-correlated index.** The indexed column's order tracks physical row order. Consecutive index entries resolve to rows on the same page. With a hundred rows per page, 50,000 matches touch roughly 500 pages, and they are visited in increasing order, so read-ahead works and most visits after the first per page are cache hits. - **Uncorrelated index.** Matching rows are sprinkled across the whole table. 50,000 matches touch up to 50,000 distinct pages, in unpredictable order, none of which prefetch helps. Each is an independent random access, and if the touched set exceeds available cache, revisits also miss. That is a difference of about a hundred times in page visits and, once you account for random versus sequential access, potentially much more in wall time. ## Where correlation comes from Correlation is a property of *insert behaviour*, not of the index definition: - Append-only tables have rows physically ordered by insertion time, so any monotonically increasing column - creation timestamp, sequence-based identifier - is naturally well correlated. - In clustered storage, an index on a prefix of the clustered key, or on a column functionally tied to it, is well correlated by construction. - Low-cardinality flags such as status are usually badly correlated, because rows of every status are interleaved throughout the table - and worse, status values *change*, so rows drift. - Foreign keys are correlated only if children were inserted grouped by parent; if they arrive interleaved by time, they are scattered. - Update-heavy workloads erode correlation over time as rows move. Optimizers track this explicitly - as a correlation figure or a clustering factor - precisely because it changes the estimated cost of the lookup step by an order of magnitude. A stale value here is a common cause of a plan that looks reasonable and behaves terribly. ## Confirming it 1. Get the plan with **actual** timings and per-node buffer statistics: pages hit in cache versus pages read from storage, at the lookup node. Massive read counts relative to matched rows is the signature of scatter. 2. Read the engine's clustering or correlation statistic for the two indexes and compare. 3. Sanity-check with a query that computes how many distinct pages the matching rows occupy, where the engine exposes physical row locations. 4. Compare cold and warm runs. A well-correlated query is fast even cold; a scattered one is fast only when its whole working set happens to be cached, which makes it look intermittent in production. ## Remedies and their costs - **Physically reorder the table** by the important index's key. This is dramatic and cheap to reason about, but it is a one-time operation that usually rewrites the table and takes locks, it only helps *one* access order, and correlation decays again as new rows arrive in a different order. - **Partition** so that the interesting rows are colocated in their own segment. This keeps the benefit as data grows and helps retention operations too, but partitioning is a schema-wide commitment that affects every query, not just this one. - **Eliminate the lookup** by including the needed columns in the index, so the row is never visited. Effective, at the cost of a wider index and more write work. - **Sort-by-locator batching.** Several engines can collect locators, sort them into physical order, and sweep the table once. Where the engine chooses this automatically, scatter hurts less; where it must be enabled or is blocked by needing index order for a sort, you may be able to influence it by removing the ordering requirement. - **Reduce the match count.** Often the honest fix: 50,000 rows fetched to display twenty is a query-design problem, and a more selective predicate or aggregation in the database beats any physical tuning. ## The interview signal Weak answers say "the second index is less selective" - but the premise fixed the row count, so selectivity is equal. The expected move is to separate *how many rows match* from *how those rows are laid out*, name correlation or clustering factor, and then propose a diagnosis based on measured page reads rather than a guess.

  • Why does a scattered-lookup query often look fast in testing and slow in production?
    In a small test dataset the whole table fits in cache, so scattered lookups are memory accesses and the scatter is invisible. In production the touched pages exceed available cache, so each lookup becomes a physical random read and latency jumps by orders of magnitude. The same effect makes production behaviour intermittent - the query is quick while its working set happens to be resident and slow after cache pressure evicts it.
  • How does sorting the row locators before fetching help, and what does it cost?
    Sorting turns scattered random reads into one increasing sweep over the table, so the storage layer can prefetch and each page is visited once instead of repeatedly. The costs are that all qualifying locators must be materialized and sorted first, needing memory proportional to the match count and delaying the first row, and that the output no longer arrives in index key order, so any query relying on that ordering needs an explicit sort afterwards.

saying these in an interview costs you the question

  • Explaining the difference by selectivity when both queries match the same number of rows
  • Assuming lookup cost scales only with row count and not with pages touched
  • Believing physically reordering a table keeps it clustered permanently
  • Judging lookup performance only from warm-cache runs
  • Treating clustering as a property of the index definition rather than of insert and update behaviour

context