skip to content

When is sorting or clustering during a rewrite worth its cost compared with plain bin-packing?

level: seniorimportance: should knowfreq 56%

answer

  1. one fixes how many files, the other fixes which files
  2. overlapping min/max means nothing can be skipped
  3. concentrate matching rows into few files
  4. the price of ordering is a shuffle
  5. hygiene on a schedule versus a deliberate campaign

basics

~20 s

Bin-packing fixes file count; sorting fixes pruning. Pay for a sort-based rewrite only when queries filter on a column whose values are scattered across every file, so per-file statistics overlap and nothing can be skipped. It costs a full shuffle.

solid answer

~50 s

Both operations rewrite data files, but they solve different symptoms. If planning is slow and scans issue huge numbers of requests, the problem is **file count** and bin-packing is enough — no shuffle, cheap, run it often. If planning is fast and file count is fine but queries still read most of the table despite a selective predicate, the problem is **data layout**: rows matching the filter are smeared across every file, so each file's min/max spans the whole range and pruning excludes nothing. Sorting or clustering by that column during the rewrite concentrates matching rows into few files, tightening statistics so the engine skips the rest. The price is a global shuffle, which is far more expensive and harder to run continuously. So: bin-pack on a schedule, cluster deliberately and less often, keyed to the columns your real query filters use.

go deeper

for a junior

Know that there are two different rewrites: one that makes files bigger and one that reorders rows so queries can skip files, and that the second is much more expensive.

for a middle

Explain why unsorted data makes per-file min/max ranges overlap so nothing prunes, and why fixing that requires a shuffle while plain repacking does not.

for a senior

An interviewer expects a diagnosis from metrics — planning time and file count versus bytes-scanned against a selective predicate — and a defensible choice of sort key tied to observed query patterns, plus a plan for the degradation as new data lands.

for a principal

Own the layout strategy as a cost decision: which tables justify recurring shuffle spend, how clustering choices are revisited when query patterns shift, and where partitioning ends and clustering begins across the platform.

## Two different diseases with one shared cure shape Both bin-packing and sort-based rewriting take existing data files and produce new ones, then atomically swap them into the table. That structural similarity makes candidates treat them as the same operation with a knob. They are not. They target different bottlenecks, cost different amounts, and are scheduled differently. **Bin-packing** answers: *the table has too many files.* **Sorting/clustering** answers: *the table has the right number of files, but the rows I want are in all of them.* ## The pruning mechanism you are trying to fix A table's metadata carries per-file (and often per-row-group) minimum and maximum values for columns. When a query filters `WHERE customer_id = 'C-8842'`, the engine compares that constant against each file's min/max for `customer_id` and skips any file whose range cannot contain it. This is the primary mechanism by which a lakehouse query reads a fraction of a huge table. It only works if the values are **clustered**. Consider data arriving as it happens: every batch contains a random mix of customer IDs. Every file therefore has a minimum near the lowest possible ID and a maximum near the highest. Every file's range contains 'C-8842'. Nothing is skipped, and a selective query reads 100% of the table while the engine reports that pruning was attempted. Sorting the table by `customer_id` during a rewrite makes each output file cover a narrow, disjoint slice of the ID space. Now exactly one file (or a few) can possibly contain that ID, and the rest are excluded at planning time. ## How to tell which disease you have Diagnose from measurements, not intuition: - **Planning time is long before any data is read; file count is enormous** → small-files problem → bin-pack. - **Bytes scanned ≈ full table size, despite a highly selective predicate; file count is reasonable** → layout problem → cluster on the predicate's column. - **Both** → bin-pack first (it is cheap and reduces the input size of everything else), then decide whether clustering is still needed. The give-away metric is the ratio of *files pruned* to *files considered*. If that ratio is near zero for a query that logically touches a tiny fraction of the data, sorting is the lever. ## Why sorting is expensive Bin-packing can be done locally: pick a set of files in one partition, read them on whatever node is convenient, write one file. No data crosses the cluster in a coordinated way. A global sort cannot. To make output files cover disjoint ranges, rows must be redistributed by the sort key across the whole rewritten scope — a shuffle, with network transfer, sort buffers and spill to disk. On a large table that is a serious job, sometimes hours. That cost is why you cluster **on a cadence measured in days**, not minutes, and why you cluster only the columns that actually appear in filters. ## Choosing the ordering Sorting is not free of tradeoffs even when justified: - **A single sort column** gives excellent pruning for that column and essentially none for others. Rows sorted by `customer_id` are random with respect to `event_time`. - **Multi-column sort** prioritizes strictly: the first column prunes well, the second only within ties of the first, and so on. A two-column sort is rarely symmetric. - **Space-filling-curve orderings** (the general idea behind multi-dimensional clustering) trade some per-column selectivity to give *several* columns moderate pruning at once. Choose these when queries filter on different columns at different times rather than always the same one. - **Partitioning already gives you coarse clustering for free** on the partition column. Do not sort on a column you already partition by; sort on the *next* most common filter. A practical rule: cluster on the column your dashboards and point lookups filter by, which is very often a high-cardinality identifier — precisely the column you should *not* partition by, because partitioning on it would shred the table into tiny directories. Clustering is how you get selectivity on high-cardinality columns without the small-file explosion that partitioning on them causes. ## Interaction with continuous ingestion A clustered table degrades as new unsorted data arrives, exactly as a compacted table degrades. Newly written files are unclustered, so queries prune the old well-sorted region and scan all of the recent one. Where that matters, the usual pattern is to re-cluster recent partitions on a schedule while leaving cold, already-clustered history alone — the layout equivalent of only compacting the hot edge. An incremental strategy also exists conceptually: rewrite only the files that are demonstrably poorly clustered rather than the whole table. This bounds cost but yields weaker global ordering. ## The honest summary Run bin-packing routinely and cheaply; treat it as hygiene. Run sorting deliberately, justified by a measured pruning failure on a named query pattern, on a slower cadence, and re-evaluate it when query patterns change — a table clustered for last year's dashboard is paying shuffle cost for nothing today.

  • How do you decide which column to cluster on when queries filter several different ways?
    Rank filters by frequency and selectivity from real query history. A single sort key serves one pattern well; when several columns matter roughly equally, a multi-dimensional clustering that gives all of them moderate pruning beats a linear sort that gives one excellent pruning and the others none. Do not cluster on a column already used for partitioning.
  • Why does clustering help high-cardinality columns where partitioning would hurt?
    Partitioning on a high-cardinality column creates one directory per value, shredding data into tiny files and bloating metadata. Clustering gets the same skipping benefit through per-file value ranges without creating any directories, so file sizes stay at target and the file count stays bounded.
  • A table was clustered six months ago and pruning is poor again. What happened?
    New writes landed unsorted. Recent files span the full value range, so any query touching recent data scans all of it, even though the older clustered region still prunes well. The fix is re-clustering the recent, degraded portion rather than the whole table.
  • Can you get clustering benefits without a global shuffle?
    Partially. Sorting within each partition or within each rewrite group avoids a table-wide exchange and still tightens per-file statistics inside that scope. It gives weaker global ordering — files from different groups still overlap — but costs a fraction of a full sort and can run far more often.

Bin-packing is consolidating a hundred half-empty filing boxes into ten full ones; sorting is alphabetizing the contents so you can pull one box instead of opening all ten.

saying these in an interview costs you the question

  • Treating sorting and bin-packing as the same operation with a flag
  • Clustering by the same column the table is already partitioned by
  • Assuming a multi-column sort prunes equally well on every column
  • Recommending a global sort without measuring whether pruning actually fails
  • Believing a clustered table stays clustered as new data arrives

context