What is clustering depth in a columnar table, and how does it show that pruning has degraded?
answer
- overlap is the enemy of skipping
- how many blocks cover one key value
- depth near one means disjoint ranges
- late-arriving data widens block ranges
- re-sorting means rewriting unchanged data
basics
~20 sClustering depth is the average number of blocks whose key ranges overlap at a given point of the key domain. Depth near 1 means disjoint ranges and near-perfect skipping; depth that climbs over time means new or rewritten blocks now span wide ranges and every query reads more.
solid answer
~50 sPruning quality is not a property of the sort key you declared — it is a property of how much the blocks' key ranges **overlap** today. Clustering depth measures that: pick a value in the key domain, count how many blocks whose `[min, max]` covers it, and average over the domain. Depth 1 means the blocks partition the key space, so a point filter reads one block. Depth 40 means forty blocks must be opened for the same filter. Depth rises whenever data is written out of key order: continuous ingest of late-arriving events, small frequent loads, updates and deletes that rewrite rows. The fixes are re-sorting or compacting the offending blocks — automatic on some platforms, a manual rewrite on others — and the real decision is whether the pruning you regain is worth the write amplification of continuously re-sorting a table that keeps receiving unordered data.
code
text · 5 lines-- computed from sampled per-block min/max on the cluster key
blocks total : 12,400
blocks with a constant key value : 310
average overlapping blocks per value : 41.2 (clustering depth)
blocks spanning > 90% of key domain : 3,100go deeper
Know that declaring a sort or cluster key does not guarantee the data stays in that order, and that overlapping block ranges mean the engine has to read more.
Define depth as the average number of blocks covering one key value, and explain why continuous small loads, late-arriving rows and updates push it upward over time.
Demonstrate the production loop: watch bytes scanned per query trend upward, confirm with the clustering metric on the column queries actually filter on, then re-sort or compact and quantify the improvement.
Own the economics — reclustering is ongoing write amplification, so decide from measured scan savings versus maintenance spend, and shape ingest batch size, key granularity and partition boundaries so the maintenance stays bounded.
## The metric and why it exists Declaring a sort or cluster key tells the engine what order you *want*. It does not guarantee the order you *have*. Every table receiving writes drifts, and the interesting question in production is: how well is this table clustered right now, on the key that matters? **Clustering depth** answers it. Consider the key column's domain, and the interval `[min, max]` each block covers for that column. For a given point in the domain, the depth is the number of block intervals covering that point; the table's depth is that count averaged (or sampled) across the domain. - **Depth 1** — the intervals are disjoint. A point filter must open exactly one block; a range filter opens exactly the blocks the range crosses. This is the theoretical best. - **Depth 5** — five blocks must be opened for a single key value; four of them will contribute nothing. - **Depth in the hundreds** — the key is effectively decorative: nearly every block covers nearly every value and you scan the table. Depth is more useful than "is the table sorted?" because it is continuous, comparable over time, and directly proportional to the bytes a selective query will read. Related numbers engines report alongside it are the count of blocks with *constant* key values (perfectly clustered blocks) and the histogram of depth across the key domain, which shows whether the problem is the whole table or only the recent tail. ## Why depth degrades Columnar blocks are immutable, so "maintaining order" means rewriting. Depth rises whenever new blocks land whose ranges straddle existing ones: 1. **Continuous ingest with late or out-of-order data.** Each small load writes blocks spanning whatever arrived in that window. If the cluster key is not the arrival order, every new block spans a wide range and overlaps everything. 2. **Many small loads.** Frequent micro-batches produce many blocks, each covering the full key spread of its batch. Depth is roughly "how many batches contain this key". 3. **Updates and deletes.** A row-level change rewrites the block (or writes a delete marker plus a new block), and the rewritten data usually lands out of order at the tail of the table. 4. **Backfills.** Reloading two years of history into a table clustered by customer produces blocks that each cover the full customer range, in one shot. A table can be beautifully clustered for its first month and useless after six, with no schema change and no query change — which is why an interviewer asks about this: it is a monitoring problem, not a design-time one. ## Detecting it in production Three signals, in order of directness: - **The clustering metric itself.** Mature platforms expose depth (and an overlap histogram) for a table's key; some do it as a function you call, others through system tables of block metadata. Where the engine exposes only per-block min/max, you can compute a serviceable approximation yourself by sampling that metadata. - **Bytes scanned per query, tracked over time.** The same dashboard query scanning 400 GB this month and 40 GB last month, with unchanged SQL and roughly unchanged data volume, is degraded clustering until proven otherwise. - **Scan-side row counts.** Rows read by the scan far exceeding rows that survive the filter, and that gap widening month over month. Check depth on the column queries actually filter on. A table can be perfectly clustered on a key nobody filters by, which is a cost with no benefit. ## Fixing it, and what the fix costs Re-establishing order means rewriting blocks — sorting the affected data and writing new, disjoint blocks. Platforms differ in packaging: some run it as an automatic background service that continuously reclusters the worst-overlapping blocks and bills you for the work (Snowflake's automatic clustering is the well-known example); others require an explicit maintenance operation or a full rewrite of the table or partition. Either way the costs are the same and must be stated: - **Write amplification.** Data that has not changed is read, sorted and written again. A table under continuous unordered ingest can be re-sorted forever without ever reaching depth 1. - **Compute spend**, continuous rather than one-off, if the service runs automatically. - **Churn on the storage layer**, including new versions of files that time-travel or snapshot retention keeps alive, so storage grows during the process. ## The judgment call Good engineering here is not "keep depth at 1". It is: 1. Cluster on the column with the highest *weighted* benefit — bytes scanned per query times query frequency — not on the primary key by reflex. 2. Reduce cardinality of the key where possible: cluster on a truncated timestamp or a coarse bucket rather than a millisecond value, so blocks can actually reach disjoint ranges. 3. Prefer fewer, larger loads over many tiny ones, since batch shape drives depth directly. 4. Bound the problem with a partition axis so re-sorting only ever touches the recent partition rather than the whole table. 5. Stop reclustering a table whose maintenance cost exceeds the scan savings — and measure both sides before deciding.
- Why can a table under continuous out-of-order ingest never reach depth 1?Because every load writes blocks spanning whatever keys arrived in that window, so overlap is recreated as fast as reclustering removes it. Maintenance becomes a steady-state cost proportional to ingest rate, not a one-off. The usual mitigations are batching writes larger, coarsening the key so ranges can align, or partitioning so only the live partition churns.
- Would you cluster on a millisecond timestamp or on a truncated date?Usually the coarser value. A millisecond key gives nearly every row a distinct value, so blocks can never share a constant key and overlap is easy to reintroduce. Truncating to hour or day lets many rows share a key, blocks reach tight or constant ranges, and typical range filters still prune to the same handful of blocks.
- How do you decide whether continuous reclustering is worth its cost?Compare two measured numbers over a representative window: the compute and storage spent on reclustering, and the bytes-scanned reduction it delivers across actual query traffic multiplied by that traffic's frequency and price. If the savings do not exceed the maintenance spend, turn it off, coarsen the key, or bound it to the recent partition.
Well-clustered blocks are like a bookshelf where each box is labelled with a tight, non-overlapping range; depth 40 is forty boxes all labelled 'A–Z', so finding one book means opening forty.
saying these in an interview costs you the question
- Assumes clustering is set once and stays good
- Confuses depth with the number of key columns
- Judges clustering by table size rather than block overlap
- Reclusters continuously without measuring the scan savings
- Clusters on a key that no query filters on