Skip to content

Compactor

The ingestor writes many small files. The compactor merges them into large files, builds the index sidecars for each merged file, and removes snapshots and files that the table does not use any more. It is the same program as the ingestor (ingestor compact), but it runs as a separate process, so that a large merge cannot take CPU from ingestion.

Compactor flow Compactor flow

Open the interactive diagram for pan and zoom, search, and export to PNG or SVG.

When a pass runs

Each backlog tick starts a pass. A pass examines each dataset in turn. After each time in the schedule, the next pass uses the rules of the schedule, which select more:

Pass Configuration key Default Selects
Backlog tick compaction.backlogCheckIntervalMinutes 15 minutes Partitions with more than compaction.backlogFileThreshold (default 100) small files.
Schedule compaction.schedule (cron, UTC) 0 2 * * * Partitions with two or more small files, and files that have no text index.

A small file is a file below compaction.minFileSizeMb (default 32). The first tick occurs when the compactor starts. The compactor examines the schedule only at a tick, so a scheduled pass starts at the first tick after the scheduled time. If a long pass continues through more than one scheduled time, the compactor does one scheduled pass, not one for each. The compactor counts scheduled times from its start: a scheduled time that passes while the compactor is stopped is not done later.

The sizes and thresholds on this page can have a different value for each dataset (ingestion.datasets[].compaction). The tick interval and the schedule apply to all datasets.

A partition that the schedule selects is not always merged. If its small files are fewer than backlogFileThreshold and their total size is below minFileSizeMb, the merged file would also be small. The compactor then leaves the partition until more data arrives, unless one of its files needs a text index (see Text index backfill).

The pass, for each selected partition

  1. Skip if done. The name of each output file comes from a hash of the input file names. If the table already contains that output, a previous run merged this partition, and the compactor skips it. This makes it safe to run a pass again after a crash.
  2. Merge. The compactor reads one input file at a time. It keeps the merged files of one partition in memory until it uploads them, so the memory that a pass needs grows with the total size of the input files of the largest selected partition. It writes files of compaction.targetSizeMb (default 128). If the dataset has sortColumns, the rows in each output file are in that order, and the compactor also writes sorted runs to the temporary disk of the pod. The compactor reads all the input files with the current table schema. At the end it makes sure that the number of rows out is equal to the number of rows in.
  3. Build the sidecars. In the same read, the compactor builds a full-text index for the dataset's tantivy.indexingColumns, a bloom filter for its bloomFilterColumns, and distinct-value sketches for its ndvColumns. The name of each sidecar contains the UUID of its data file.
  4. Upload. The compactor writes each merged file and its sidecars to Azure Blob.
  5. Commit. One Iceberg commit replaces the small files with the merged files. Queries see the old files or the new files, not a mixture.

Merge sequence

This is the order of the calls for one partition. An ingestor can add files to the table at the same time. If it commits first, the compactor sends its commit again.

Compactor merge sequence Compactor merge sequence

Open the interactive diagram.

Rules that protect your data

Upload first, then commit. The table does not refer to a file until that file is in your storage. If the upload of a data file fails, there is no commit.

A sidecar that is missing is not an error. If the upload of a sidecar fails, the commit continues. The query gateway scans a file that has no sidecar. It does not skip it.

An ingestor append does not undo a merge. The commit replaces only the files that the compactor merged. If another writer commits first, the catalog refuses the commit. The compactor then reads the table again and sends the commit again, up to compaction.commitRetryMax times (default 5). It does not merge the files again. Each refusal increments tracelake_catalog_commit_conflicts_total{component="compactor"}. If the retries run out, that partition fails, and the pass stops for that dataset. A later pass merges it again. If two compactors merge the same files, the second one finds that its input files are gone, and that partition fails in that pass.

The commit does not delete the old files. Queries that started before the commit can continue to read them. The two maintenance loops remove them later.

Text index backfill

The full-text index is built only by the compactor. The ingestor does not build it, because that would take CPU from ingestion. A file that the ingestor wrote at its maximum size is not a small file, so the size rule does not select it. On a scheduled pass, the compactor therefore also selects a file when all of these are true:

  • the dataset has tantivy.indexingColumns;
  • the file has no text index;
  • TraceLake wrote the file;
  • the file is not larger than targetSizeMb.

Maintenance loops

Two loops run on their own timers, independently of the compaction pass:

Loop Interval key Default Does
Expire snapshots maintenance.expireSnapshotsIntervalMinutes 60 minutes Removes table snapshots older than compaction.snapshotRetentionHours (default 24). It deletes no files.
Orphan cleanup maintenance.orphanCleanupIntervalMinutes 60 minutes Deletes data files that no snapshot refers to and that are older than maintenance.orphanMinAgeHours (default 24). It deletes the sidecars of each file with it.

The age limit protects files that an ingestor has uploaded but not committed yet. Orphan cleanup has no switch: it starts with the compactor. Before you start the compactor for the first time in a deployment, run ingestor compact --config <file> --orphan-dry-run. It lists the files that a cleanup would delete, and deletes nothing.

States of a data file

Data file states Data file states

Open the interactive diagram.

A file goes through the numbered states in order. The diagram draws no arrows between them.

  • Live. The current snapshot of the table refers to the file.
  • Replaced. A merge commit replaced the file. Only older snapshots refer to it, and queries that use those snapshots can continue to read it.
  • Unreferenced. Snapshot expiry removed the last snapshot that referred to the file.
  • Deleted. Orphan cleanup deleted the file and its sidecars, after the file was older than maintenance.orphanMinAgeHours.

A file that a writer has uploaded but not committed yet has no snapshot that refers to it. To orphan cleanup, it looks the same as an unreferenced file. The age limit prevents its deletion before the commit. If the writer stopped and the commit never occurs, the file becomes an unreferenced file, and orphan cleanup deletes it later.

When a pass fails

The compactor counts the failure, writes the cause to its log, and continues with the next dataset. The process does not stop, so the schedule continues to operate for the other datasets. An alert on the failure counter is the signal to the operator.

If the files of one partition have column types that are not compatible, the compactor does not merge that partition. The small files stay, and the pass reports a failure. You must migrate that partition with an external engine.

compaction.runtimeDeadlineHours (default 6) does not stop a pass. A pass that runs longer increments a counter for an alert.

Manual runs

You can run one pass by hand. A run is a manual run only if it has one or more of the five flags below --config. With none of them, ingestor compact --config <file> starts the long-running compactor, which includes orphan cleanup.

A manual run does one pass and exits. It does not expire snapshots or delete files, and its exit code is not zero if the pass failed. Like a scheduled pass, it also selects files that have no text index.

Flag Effect
--config <file> Required. The configuration file, the same file that the ingestor reads.
--dataset <name> Compact one dataset only.
--date-range <from>:<to> Compact the partitions in an inclusive YYYY-MM-DD date range. The dataset must have a partition column that holds a date.
--partition <column>=<value> Compact the partitions with this value. You can repeat the flag; a partition must then match all the values.
--dry-run List the files that a run would merge. Merge and commit nothing.
--orphan-dry-run List the files that orphan cleanup would delete. Delete nothing, and do no compaction.

You can run several compactors at the same time on different datasets or date ranges.