Skip to main content

Incremental Clustering

Incremental clustering improves the data layout of append tables by sorting selected files on frequently filtered columns. Compared with repeatedly sorting an entire partition, it can reduce the amount of data rewritten while improving file-statistics pruning. A run may select no files when its compaction criteria are not met. Full mode considers all runs in the selected scope, but can skip work that is already clustered; see file selection.

Clustering also merges small files, respecting target-file-size. It changes the physical layout, not the rows returned by a query, and does not replace SQL ORDER BY.

Requirements​

RequirementUnaware-bucket append (bucket = -1)Bucketed append (bucket > 0)
Primary keyMust not be defined.Must not be defined.
Enable clusteringclustering.incremental = true and nonempty clustering.columns.Same.
Append orderingNo bucket-order guarantee.Must set bucket-append-ordered = false.
Deletion vectorsSupported.Must remain disabled.
Compaction executionSchedule explicit clustering jobs; the Flink sink's normal background compaction is disabled.Writer compaction and dedicated compact jobs use the bucket clustering path.
Global/local sort modeConfigurable for batch clustering jobs.Clustering is performed within each partition and bucket; the global/local option does not select this path.
Historical-partition auto-clusteringSupported.clustering.history-partition.* does not apply.

Data Evolution tables cannot enable incremental clustering. If streaming consumers require ordered append reads from a bucketed table, keep that ordering and do not enable clustering.

Enable Incremental Clustering​

Set the clustering keys on the table using the DDL for its layout. Choose the tab for your engine.

Unaware-Bucket Table​

For my_table from the overview, enable clustering and specify the columns:

ALTER TABLE my_table SET (
'clustering.incremental' = 'true',
'clustering.columns' = 'product_id,price'
);

Bucketed Table​

For bucketed_table from Bucketed append, disable append ordering in the same statement that enables clustering. Keep deletion vectors disabled. This explicitly opts out of ordered streaming reads.

ALTER TABLE bucketed_table SET (
'bucket-append-ordered' = 'false',
'clustering.incremental' = 'true',
'clustering.columns' = 'product_id,price'
);

Clustering Options​

OptionDefaultHow to use it
clustering.incrementalfalseSet to true to enable incremental clustering.
clustering.columnsNot setComma-separated columns, such as product_id,price. Prefer frequently filtered data columns over partition columns.
clustering.strategyautoorder, zorder, or hilbert. Automatic selection uses order for one column, zorder for two to four, and hilbert for five or more.
clustering.incremental.modeglobal-sortSort execution mode for unaware-bucket batch clustering; see below.

Choose a Sort Mode​

For unaware-bucket tables, the mode controls how the selected files in each partition are sorted:

ModeExecutionTradeoff
global-sortRange-shuffles rows across tasks, then sorts within tasks using the configured clustering strategy.Coordinates the layout across the selected output files, at the cost of a network shuffle.
local-sortSorts rows independently within each compaction task, without the global range shuffle.Less shuffle work; ranges in files produced by different tasks can overlap. Useful when ordering within files is sufficient, such as for Parquet lookup optimizations.

Here, “global” refers to the selected clustering work within a partition. It does not imply that all existing files or all table partitions become globally ordered after an incremental run.

Run Incremental Clustering​

Run explicit compact jobs in batch mode. Table options supply the clustering columns and strategy. The examples below show routine incremental selection (minor) and full clustering (full) of the selected table or partition scope.

-- Choose parallelism for the workload; too many tasks can produce small files.
SET spark.sql.shuffle.partitions = 10;

-- Select files using the incremental compaction strategy.
CALL sys.compact(table => 'my_table', compact_strategy => 'minor');

-- Alternatively, request full clustering; already-clustered runs can be skipped.
CALL sys.compact(table => 'my_table', compact_strategy => 'full');

With historical-partition auto-clustering disabled (the default), use partitions to limit the work to one partition of my_table:

CALL sys.compact(
table => 'my_table',
partitions => 'dt=2026-09-10',
compact_strategy => 'full'
);

Alternatively, where accepts a predicate on partition columns. Do not combine partitions and where in one call.

On unaware-bucket tables, historical-partition auto-clustering can add full clustering of partitions outside either filter, so these arguments are not a hard job boundary when that feature is enabled.

For the unpartitioned bucketed_table example, change the table name and omit the partition filter.

For an unaware-bucket table, you can override the sort mode for one invocation:

CALL sys.compact(
table => 'my_table',
compact_strategy => 'minor',
options => 'clustering.incremental.mode=local-sort'
);

The compact entry points route to the appropriate clustering implementation when clustering is enabled. On unaware-bucket tables, use these jobs for recurring small-file merging and layout maintenance. On bucketed tables, writer compaction can also perform incremental clustering; set write-only = true on ingestion if that work should be performed only by dedicated jobs.

Verify the Result​

Inspect the files before and after a clustering job. For example, in Spark SQL:

SELECT `partition`, bucket, level,
COUNT(*) AS file_count,
SUM(file_size_in_bytes) AS total_bytes
FROM `my_table$files`
GROUP BY `partition`, bucket, level
ORDER BY `partition`, bucket, level;

SELECT snapshot_id, commit_kind, commit_time
FROM `my_table$snapshots`
ORDER BY snapshot_id DESC
LIMIT 10;

The files system table shows the current physical layout. Compare file counts, sizes, and levels in the partitions and buckets selected by the job. The snapshots table helps identify commits around the job's execution time; concurrent ingestion can also create snapshots. A successful job may leave the files unchanged when the planner selects no work, including the full-mode skip cases.

To inspect pruning potential, query file_path, min_value_stats, and max_value_stats from my_table$files and compare the ranges for your filter columns. Then run a representative query and compare scan metrics such as files or bytes read. Higher levels or fewer files alone do not establish that a query reads less data. Keep the data and query comparable when measuring the effect, especially if ingestion continues during maintenance.

Change Clustering Keys​

Update clustering.columns and, if needed, clustering.strategy using the same ALTER TABLE syntax as above. Changing the options does not immediately rewrite existing data. Subsequent clustering uses the new settings for selected files. After changing the clustering columns, use compact_strategy = full to apply them across the selected scope. Changing only clustering.strategy does not force an already-clustered run to be rewritten: the full-mode skip check compares the clustering columns, not the sorting strategy.

Auto-Clustering for Historical Partitions​

For partitioned unaware-bucket tables, a clustering run can also select inactive historical partitions outside the requested partition predicate for full clustering. This additional selection requires both a configured idle duration and an explicit partition predicate: Spark's partitions or where, or Flink's --partition.

Without an explicit partition predicate, the historical auto-full path is inactive. For example, setting the idle duration does not activate it for the unscoped minor call above; that call still uses normal incremental selection across the table. Auto-clustering happens within a submitted job; configuring the options does not start a scheduler.

OptionDefaultPurpose
clustering.history-partition.idle-to-full-sortNot set (disabled)How long a partition must have no new updates before it is considered historical.
clustering.history-partition.limit5Maximum number of additional historical partitions selected outside the requested partition predicate in a run.

For example, in Flink SQL:

ALTER TABLE my_table SET (
'clustering.history-partition.idle-to-full-sort' = '3 d',
'clustering.history-partition.limit' = '5'
);

With these options, a job requesting dt=2026-09-10 can also fully cluster up to five eligible historical partitions outside that date. The additional partitions use full clustering even if the requested partitions use compact_strategy = minor. The planner evaluates their eligibility and whether a full rewrite is needed.

To keep the entire job within the requested partitions, leave clustering.history-partition.idle-to-full-sort unset, or remove it from the table options before submitting the job. The partition limit controls the additional historical work, not the number of explicitly requested partitions. These options do not apply to bucketed append tables.

How File Selection Works​

Incremental clustering organizes files into levels and uses a universal compaction strategy to select sorted runs. Normally, newly appended files enter level 0. At level 0, each file is treated as its own run; at a higher level, the files produced by a clustering set are treated as one run.

Incremental clustering selects new files and an eligible existing run, rewrites that set, and keeps an unselected higher-level run.

The planner balances the number and relative sizes of runs against write amplification. A run rewrites only the selected set, so higher-level files can remain untouched. The levels in the diagram are illustrative; the selected files and output level depend on the planner.

Full mode selects all runs in a compaction unit unless there is no data, or the unit already contains a single run at the highest level with the same clustering columns. In those cases, no rewrite is needed by the planner. This check applies per partition for unaware-bucket tables and per partition and bucket for bucketed tables, so a full job can rewrite some units while leaving others unchanged.