Skip to content

[Spark] Lightweight compaction for clustered tables - #19

Draft
sezruby wants to merge 1 commit into
masterfrom
stats-aware-compaction-binning
Draft

sezruby wants to merge 1 commit into
masterfrom
stats-aware-compaction-binning

Conversation

@sezruby

@sezruby sezruby commented Oct 6, 2026 •

Copy link
Copy Markdown
Owner

Which Delta project/connector is this regarding?

  • Spark

Description

Adds an opt-in lightweight compaction tier to OPTIMIZE and auto compaction on clustered (liquid) tables.

Motivation

Today, incremental OPTIMIZE on a clustered table:

  • takes all unclustered files plus every file in a partial Z-cube (smaller than minCubeSize, 100 GB by default);
  • range-partitions and sorts all of them along the clustering curve.

Any unclustered file makes the partial cubes eligible again. So on tables that get small batches between frequent OPTIMIZE or auto compaction runs, a little new data can rewrite up to ~100 GB of already-clustered data on every run. That is high write amplification, plus a full shuffle each time.

Change

When spark.databricks.delta.optimize.clustering.lightweight.enabled is true, and the table has some unclustered data but less than ...lightweight.maxUnclusteredBytes, OPTIMIZE uses a new LightweightClusteringStrategy instead of ClusteringStrategy:

  • Which files: only small unclustered files (or files with many deleted rows), using the regular compaction filter. Clustered files and partial Z-cubes are not touched.
  • How they are grouped: candidate files are ordered by the per-file (min, max) stats of a single clustering column before bin packing, so each output file covers a narrow range of that column. Rows are not sorted; each bin is coalesced, as in plain compaction.
  • Which column: set by ...lightweight.column. Otherwise it is auto-picked as the clustering column whose files overlap least. The overlap is measured on ranks of the distinct min/max values, so it works for any orderable type. If no column's average relative file range is at most ...lightweight.maxRelativeFileRange, files are grouped by size. One column is used because grouping whole files cannot keep several columns narrow at once; full clustering handles the rest.
  • Output: files stay unclustered, with no clusteringProvider or Z-cube tags. They are clustered once enough unclustered data has accumulated. The commit records clusterBy=[], so lightweight runs can be told apart in history.
  • Clustering still runs with OPTIMIZE FULL, during REORG, at or above the threshold, and when there is no unclustered data (so partial cubes can still merge).
Config (spark.databricks.delta. prefix) Default
optimize.clustering.lightweight.enabled false
optimize.clustering.lightweight.maxUnclusteredBytes 10 GB
optimize.clustering.lightweight.column (auto)
optimize.clustering.lightweight.maxRelativeFileRange 0.5

Defaults leave current behavior unchanged.

Prior art

  • Hudi clustering can order file slices by commit time (INSTANT_TIME).
  • Cassandra TWCS groups SSTables by their min/max timestamp windows.

This change generalizes that idea to a stats-chosen clustering column, as a cheap tier in front of full clustering.

How was this patch tested?

New LightweightClusteringSuite (15 tests):

  • default behavior still clusters;
  • auto-pick of the narrowest column, with and without column mapping;
  • fallback to size grouping;
  • configured column (case-insensitive, nested, non-narrow);
  • invalid column error;
  • clustered files are left untouched;
  • clustering at or above the threshold;
  • FULL clusters;
  • lightweight output is clustered later;
  • auto compaction;
  • unit tests for the overlap measure.

Also ran IncrementalZCubeClusteringSuite, ClusteringProviderSuite, OptimizeCompaction{SQL,Scala}Suite, OptimizeMetricsSuite, AutoCompact{Execution,Configuration}Suite, DeltaErrorsSuite, DeltaThrowableSuite, and scalastyle.

Does this PR introduce any user-facing changes?

Yes, new opt-in SQL configs (documented in delta-clustering.mdx). No behavior change by default.

@sezruby
sezruby force-pushed the stats-aware-compaction-binning branch from 8de554e to 4596f79 Compare October 6, 2026 20:51
@sezruby sezruby changed the title [Spark] Group files by column statistics in compaction [Spark] Lightweight compaction for clustered tables Oct 6, 2026
Add an opt-in lightweight tier to OPTIMIZE and auto compaction on clustered
tables. While the table has some unclustered data, but less than
spark.databricks.delta.optimize.clustering.lightweight.maxUnclusteredBytes,
OPTIMIZE compacts only the small unclustered files instead of clustering:
clustered files and partial Z-cubes are not rewritten, and rows are not sorted.
Candidate files are ordered by the per-file min/max statistics of a single
clustering column before they are packed into bins, so each output file covers
a narrow range of that column. The column is configurable, or picked as the
clustering column whose files overlap least. Output files stay unclustered and
are clustered once enough unclustered data has accumulated. OPTIMIZE FULL
always clusters.

Configs:
- optimize.clustering.lightweight.enabled (default false)
- optimize.clustering.lightweight.maxUnclusteredBytes (default 10GB)
- optimize.clustering.lightweight.column (default: auto-pick)
- optimize.clustering.lightweight.maxRelativeFileRange (default 0.5)

Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com>
@sezruby
sezruby force-pushed the stats-aware-compaction-binning branch from 4596f79 to 4fcdc8d Compare October 7, 2026 05:34
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant