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
| Requirement | Unaware-bucket append (bucket = -1) | Bucketed append (bucket > 0) |
|---|---|---|
| Primary key | Must not be defined. | Must not be defined. |
| Enable clustering | clustering.incremental = true and nonempty clustering.columns. | Same. |
| Append ordering | No bucket-order guarantee. | Must set bucket-append-ordered = false. |
| Deletion vectors | Supported. | Must remain disabled. |
| Compaction execution | Schedule 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 mode | Configurable for batch clustering jobs. | Clustering is performed within each partition and bucket; the global/local option does not select this path. |
| Historical-partition auto-clustering | Supported. | 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:
- Flink
- Spark
ALTER TABLE my_table SET (
'clustering.incremental' = 'true',
'clustering.columns' = 'product_id,price'
);
ALTER TABLE my_table SET TBLPROPERTIES (
'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.
- Flink
- Spark
ALTER TABLE bucketed_table SET (
'bucket-append-ordered' = 'false',
'clustering.incremental' = 'true',
'clustering.columns' = 'product_id,price'
);
ALTER TABLE bucketed_table SET TBLPROPERTIES (
'bucket-append-ordered' = 'false',
'clustering.incremental' = 'true',
'clustering.columns' = 'product_id,price'
);
Clustering Options
| Option | Default | How to use it |
|---|---|---|
clustering.incremental | false | Set to true to enable incremental clustering. |
clustering.columns | Not set | Comma-separated columns, such as product_id,price. Prefer frequently filtered data columns over partition columns. |
clustering.strategy | auto | order, zorder, or hilbert. Automatic selection uses order for one column, zorder for two to four, and hilbert for five or more. |
clustering.incremental.mode | global-sort | Sort 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:
| Mode | Execution | Tradeoff |
|---|---|---|
global-sort | Range-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-sort | Sorts 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.
- Spark SQL
- Flink Action
-- 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'
);
<FLINK_HOME>/bin/flink run \
-Dexecution.runtime-mode=batch \
/path/to/paimon-flink-action-2.2-SNAPSHOT.jar \
compact \
--warehouse <warehouse-path> \
--database <database-name> \
--table <table-name> \
--compact_strategy minor \
--table_conf sink.parallelism=2
Replace the placeholders with your installation and table details. Use --compact_strategy full for full clustering.
For the partitioned my_table example, add --partition dt=2026-09-10 to select that partition. This limits the whole
job when historical-partition auto-clustering is disabled; when enabled,
the job can also fully cluster eligible historical partitions outside the requested partition.
For an unaware-bucket table, add --table_conf clustering.incremental.mode=local-sort to override the default sort mode.
Add --catalog_conf key=value arguments as required by your catalog and storage.
sink.parallelism controls the Flink action's sink parallelism. Size it for the amount of selected data rather than
using a large value for every run.
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.
| Option | Default | Purpose |
|---|---|---|
clustering.history-partition.idle-to-full-sort | Not set (disabled) | How long a partition must have no new updates before it is considered historical. |
clustering.history-partition.limit | 5 | Maximum 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.
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.