Notes · Published 25 September 2026
Iceberg compaction on Spark: why it is expensive and what to measure
ETL Compare staff · Published 25 September 2026 · Figures checked September 2026
- Workloads
- Benchmarks
Compaction reads, sorts, shuffles and rewrites table data on a schedule, which makes it a good test of an accelerator's shuffle and spill claims.
In brief
Iceberg compaction rewrites many small data files into fewer large ones, and sort-based compaction also re-sorts the data, so one run reads, decodes, sorts, shuffles and re-encodes a large share of a table. That makes it shuffle- and spill-heavy and a useful workload for testing accelerators. Measure it as cost per TB rewritten and by its stage profile, not by a headline speedup.
What does compaction do in Apache Iceberg?
Frequent appends and streaming writes leave tables with many small files, and every reader then pays to open and plan them (see the small files problem). Iceberg ships Spark procedures for table maintenance. rewrite_data_files combines data files using a binpack, sort or zorder strategy; rewrite_manifests rewrites manifest files "for better query performance and metadata efficiency"; and expire_snapshots removes expired snapshot metadata. Compaction in this note means rewrite_data_files.
Binpack groups small files into larger ones. Sort and zorder also reorder rows by chosen columns, which helps later reads skip data but requires Spark to sort the data it rewrites.
Why is compaction so expensive?
In Spark terms a sort-based compaction is a full ETL job: a scan stage that reads and decodes Parquet, an exchange that shuffles rows to sort them, a sort that can spill to disk when a partition is larger than the memory its task gets, and a write stage that compresses and commits new files. DualBird's benchmark post describes sort-based compaction as one of the most compute-intensive jobs in a lakehouse because it "reads, decodes, sorts, shuffles, and re-encodes potentially hundreds of gigabytes of Parquet data". Because it runs on a schedule across many tables, its cost recurs.
The cost also grows with how far maintenance falls behind. A table that is compacted rarely has more small files per partition, so each run reads more files and rewrites more data, and a sort over a large partition is more likely to run out of task memory. Smaller, more frequent rewrites spread the work, but every run still pays the scan and write cost, which is why cost per TB rewritten is the number to track.
Which rewrite_data_files options shape the cost?
The Iceberg Spark procedures page lists the options and their defaults. These are the ones that most affect how much work a run does and how much memory it needs:
| Option | Default | What it controls |
|---|---|---|
| target-file-size-bytes | 536870912 (512 MB) | The size of output files the rewrite aims for |
| min-input-files | 4 | The minimum number of files in a group before that group is rewritten |
| max-concurrent-file-group-rewrites | 10 | How many file groups are rewritten at the same time, which drives parallelism and memory demand |
Two more options change the scale and risk of a run: partial-progress.enabled lets a long rewrite commit in parts, so a failure loses less work, and rewrite-all forces every matching file to be rewritten, which is the most expensive choice.
Test the strategy you actually run. A binpack rewrite and a sort rewrite of the same table stress different stages, and a result for one does not carry over to the other.
What do vendors and projects state about compaction and Iceberg?
Of the seven products we track, DualBird is the one whose pages we reviewed include a compaction-specific benchmark. Its June 2026 post covers sort-based compaction of a 100 GB dataset of about 370 million rows and 64 mixed-type columns with ZSTD level 3 compression, against vanilla Spark and "state-of-the-art C++ accelerated Spark".
Vendor states
“12-20x faster Spark task execution” and “55-85% lower EC2 cost”, about $2/TB against $6/TB to more than $13/TB for the CPU clusters tested
Source: DualBird: Sort-based Iceberg compaction benchmark · Fetched 27 Sep 2026
Vendor figure. Not measured by ETL Compare.
The post is written by DualBird staff. As noted in our DualBird review, cluster sizes and Spark versions are not stated and we found no third-party results.
Other products publish Iceberg results that are not about compaction. Comet 1.0.0 (7 August 2026) adds Iceberg table format V3 support with table encryption, and the AWS engineers' June 2026 post reports Comet 37% faster overall on Iceberg tables for their TPC-DS queries. The RAPIDS Accelerator's 26.08.0 release adds Iceberg 1.11 support with accelerated incremental append scans. We found no compaction-specific figures from Gluten, Auron, Photon or Flarion on the pages we reviewed.
How should you measure a compaction trial?
- Pick one real table and a fixed set of partitions to compact, and copy it so baseline and trial start from the same files.
- Record files and bytes before and after, and the TB rewritten.
- Compute cost per TB: cost per run divided by TB rewritten, using cost per job.
- Record the stage profile: shuffle write, Shuffle spill (disk) and the slowest task.
- Run sort and binpack separately if you use both, at two cluster sizes.
- Check the result: row counts match and the new snapshot commits cleanly.
For the full trial method, see how to run a Spark accelerator proof of concept.
Sources
- DualBird Iceberg compaction benchmark · Fetched 27 Sep 2026
- Iceberg Spark procedures · Fetched 27 Sep 2026
- Comet 1.0.0 release notes · Fetched 27 Sep 2026
- Comet on Amazon EKS · Fetched 27 Sep 2026
- NVIDIA cuDF for Apache Spark release notes · Fetched 27 Sep 2026
- Web UI - Spark documentation · Fetched 27 Sep 2026