Physical ordering should be chosen from observed predicates and maintained only while its read savings exceed rewrite and write costs.
Workload-driven clustering and layout drift
Pick from real query shapes
Inspect a representative week of filters, joins and scanned bytes. If most selective requests constrain account ID within a day, daily partitioning plus account clustering may outperform adding a separate partition for every account. Too many tiny partitions multiply metadata and writer work. Do not choose a clustering key merely because it has many distinct values; verify that the query engine stores usable bounds and actually applies the predicate. Safe skipping is the mechanism being purchased.
Compare layouts on the same snapshot
Build a candidate layout from one pinned generation and run the same cold and warm queries against both layouts. Capture planning time, listed files, bytes read, CPU and end-to-end latency. Check selective equality filters, broad scans, writes and maintenance jobs. A single fast dashboard query does not justify a rewrite if ingest latency and storage amplification harm the rest of the workload.
Detect decay from new writes
Fresh append files may arrive unsorted even when older partitions are tightly grouped. As their share grows, min/max ranges overlap and skip efficiency falls. Track candidate files per predicate over time, along with file sizes and overlap, and schedule maintenance when projected read waste exceeds rewrite cost. Compaction manages file count; sorting addresses locality, so measure both independently.
Budget rewrite side effects
A sort rewrite consumes shuffle, I/O and temporary storage. It may increase write amplification or conflict with concurrent appends, and it must preserve row values, deletes and snapshot visibility. Choose a bounded subset such as recent high-traffic days, publish it atomically, and keep the prior generation available until verification. Avoid a daily full-table rewrite if only a small tail is responsible for query degradation.
Set an exit rule
Define the minimum saved bytes or latency per rewrite hour, the largest allowed writer lag, and the maximum query regression for broad scans. Revisit key choice when the workload changes. A former account-filter workload may become a time-range workload, making the old clustering expensive but ineffective. Store the benchmark data with the release so a later operator can reverse the choice rather than defending it by habit.
Implementation
query_weights = {
"account_lookup": 47,
"regional_rollup": 29,
"full_scan": 8,
}
baseline_bytes = {
"account_lookup": 920,
"regional_rollup": 680,
"full_scan": 1200,
}
candidate_bytes = {
"account_lookup": 145,
"regional_rollup": 510,
"full_scan": 1200,
}
def weighted_scan_cost(weights, bytes_by_query):
return sum(count * bytes_by_query[query_name]
for query_name, count in weights.items())
saved_mb = (weighted_scan_cost(query_weights, baseline_bytes)
- weighted_scan_cost(query_weights, candidate_bytes))
rewrite_mb = 19_000
assert saved_mb == 47 * 775 + 29 * 170
assert saved_mb > rewrite_mbPerformance and operating cost
For Q query classes, weighted cost comparison is O(Q) time and O(1) extra space, but the measured workload is the hard part. Sort maintenance may read and write every selected file, often with shuffle and temporary duplicate storage. Its benefit depends on predicate selectivity, statistics availability and how quickly new writes degrade locality.
Common Mistakes
- Do not add high-cardinality partitions just to imitate clustering.
- Do not compare layouts using different table snapshots or cache states.
- Do not ignore unsorted append files when reporting a one-time benchmark.
