Clustering Decay
also called Sort Order Drift, Layout Degradation
The gradual loss of data-skipping as newly written files overlap the sort order of existing ones - so query cost rises month after month with no change to the query or the data volume.
A table was clustered by customer_id in March and queries were fast. By September the same dashboard reads four
times the bytes and costs four times as much. The query is unchanged, daily volume is unchanged, nothing was
deployed. Every file written since March arrived in event order, so each spans the full range of customer ids,
and a filter on one customer now matches the statistics of every recent file.
Clustering decay is that drift. A sort order is a property of the files that existed when the sort ran, not of the table, and every later write erodes it until a maintenance job restores it.
Why it matters
Data skipping is the dominant lever in analytical cost - far larger than cluster size - and it degrades silently. A decayed sort order throws no error; the only symptom is a cost curve bending gently upward over months, indistinguishable from growth and usually explained away as growth.
The quantities are not marginal. A single-customer query against a well-clustered table may open tens of files; against a decayed one it opens thousands and reads hundreds of times the bytes for the same answer. Teams respond by adding compute, which makes the query fast again and the bill worse.
Implementation patterns
- Measure overlap, not elapsed time. The signal is how many files a typical filter value appears in - the overlap depth - or the bytes scanned by a fixed benchmark query. Alert on the trend.
- Trigger re-clustering on a threshold, not a calendar. A table taking one late batch a week decays slowly and a streaming table decays daily, so one schedule for all tables either wastes compute or starves the worst tables.
- Re-cluster incrementally by level. Merge new files into the existing sorted runs, as a log-structured merge does, so cost scales with new data rather than table size.
- Sort by the filter column and partition by the time column, and scope each pass to the files written since the last one, so the job's duration stays roughly constant as the table grows.
- Give maintenance reserved capacity, because clustering has no deadline and loses to jobs that do.
Industry example
Datadog has described two compaction strategies in Husky, its event store over object storage. Size-tiered compaction merges writer fragments of a few thousand rows through exponentially larger size classes up to roughly a million rows. Locality compaction then re-sorts fragments by a lexical sort schema - service, status, environment, timestamp - through log-structured levels, so that queries can prune. Datadog reports that rolling it out produced about a 30% reduction in query worker replicas, which it calls its most expensive component (Datadog engineering blog, 2025).
Read as a statement about decay: the sort order has to be re-established continuously as fragments land, and the payoff is counted in serving capacity rather than storage.
Failure scenarios
- Silent cost drift over two or three quarters, attributed to growth, with the layout never examined.
- Maintenance starvation: clustering submitted as low-priority work that never gets a slot in the busy season, so decay compounds when query volume is highest.
- Re-clustering fighting a streaming writer for the commit, so the job retries and completes nothing.
- A sort order chosen for yesterday's queries: the dashboards changed, the clustering key did not, and perfect maintenance prunes nothing.
- Storage doubling after a full re-sort, because snapshots hold every rewritten file until expiry runs.
Trade-offs
Re-clustering costs compute and write amplification - each pass rewrites what it touches - and competes with the workloads it accelerates. Leaving it undone costs bytes on every query, forever. The arithmetic almost always favours maintenance, because a rewrite is paid once per batch of new data and a scan is paid once per query. The real trade-off is not whether but how far back.
When not to use it
When queries do not filter on a high-cardinality column at all. A table only ever read in full - a small dimension, a daily aggregate, a feature table consumed whole - cannot decay, so maintaining a sort order is pure cost. The same applies where partitioning alone prunes selectively, and to write-once tables: a single sort at write time never decays.
Interview question
Q: "A dashboard's cost has risen 4x in six months. Daily volume is flat, the query text has not changed, the engine version is the same. Walk me through the diagnosis, and what you would put in place so the next occurrence is detected rather than discovered."
What a strong answer covers: bytes scanned as the first metric rather than duration; overlap depth as the mechanism; separating decay from a changed query pattern and from partition pruning that stopped working; incremental re-clustering on a threshold; reserved maintenance capacity; and a benchmark query tracked over time as the detection mechanism.
Quick check
Quiz: Clustering was applied once and queries have slowly become more expensive. What changed? Nothing in the query - new files arrived in event order, each spanning the full key range, so per-file statistics no longer exclude them and skipping has degraded.
Flashcard: Which signal says a clustered table needs re-clustering, and which misleads? — Bytes scanned or files opened by a fixed benchmark query, trending upward; wall-clock duration misleads, because adding compute hides the decay while the bill keeps rising.