firas bouzazi

2026-05-26

Apache Iceberg compaction in practice

A practical guide: architecture, compaction strategies, hands-on examples (AWS and Spark) and full table maintenance.

This article assumes you've heard of Apache Iceberg and maybe played with it a little, but compaction still feels like a black box. These are my notes from going deep on it: how it works, how to run it, and how to not shoot yourself in the foot.

1. Reminder of Apache Iceberg architecture

Before talking about compaction, let's do a quick recap of how Iceberg organizes things under the hood. It will help make sense of why compaction even matters.

Iceberg architecture: the catalog points to the current metadata file. Metadata files point to manifest lists, one per snapshot. Manifest lists point to manifest files, and manifest files point to the data files. METADATA LAYER DATA LAYER catalogpoints to current metadata metadata file[s0] metadata file[s0, s1] manifest listsnapshot s0 manifest listsnapshot s1 manifest file manifest file manifest file data files data files data files
How a query finds its data: catalog, metadata, manifests, data files.

Data layer

This is where the actual data lives. It's made of datafiles that store the rows themselves (Parquet, ORC, Avro…). Those files are the leaves of the tree.

There's also another type of file here: delete files. They don't store data, they track which records of existing datafiles have been deleted. Iceberg uses them to handle deletes efficiently without immediately rewriting the full datafile.

Metadata layer

This is the brain of the table. It's a three-level structure:

  • Manifest files: track a subset of the datafiles in the data layer, plus some extra stats per file (row counts, min/max values, etc.)
  • Manifest lists: a snapshot of the full table at a given point in time. Each snapshot maintains all the manifest files needed to reconstruct the table at that moment.
  • Metadata files: sit at the top. They track all the manifest lists (= all snapshots), and also store table-level info: table name, partition spec, schema history, etc.

The catalog

The catalog is basically the entry point. It holds a pointer to the current metadata file for each table. When you run a query, Iceberg goes to the catalog first, follows the pointer to the metadata file, walks down through the manifest list and manifest files, and finally reads the relevant datafiles.

You can use several catalog implementations: Hive Metastore, AWS Glue, REST catalog, Nessie, and others.

2. Compaction

Compaction is one of the most important maintenance operations for keeping an Iceberg table healthy. In short: it reads a bunch of small files and merges them into fewer, larger files. You can also sort the data in the process, which we'll get to later.

Before compaction: many small files between 1 and 32 MB. After rewrite_data_files: a few files of about 256 MB each. BEFORE AFTER 32MB 8MB 12MB 16MB 20MB 24MB 28MB 8MB 32MB 256MB 256MB 256MB 256MB rewrite_data_files() many small files few files, same size
Compaction rewrites many small files into a few files near the target size.

Why do we need it?

Imagine a table that has accumulated thousands of small datafiles. This is called the small files problem, and it causes real pain:

  • Queries slow down significantly. Each time you run a query, Iceberg must open each file, read metadata, plan execution, then close it. 10,000 files = 10,000 open/read operations.
  • More files means a bigger metadata footprint. More manifest entries = slower query planning.
  • On managed services (AWS, GCP, Azure…), every file open is a storage request. More requests = higher cost.

This happens more often than you'd think. Typical causes:

  • Streaming writes with low throughput (Flink, Kafka consumers)
  • Spark micro-batches writing frequently
  • CDC event logging

The fix: periodically run compaction to reshape your data into fewer, larger files.

3. Hands-on compaction with Spark

3.1 Simple way: the SQL command

The easiest way to trigger compaction:

CALL system.rewrite_data_files(
  table => 'my_db.my_table'
);

This scans for small files and merges them. By default it targets the file size set in the table property write.target-file-size-bytes, which defaults to 512 MB.

You can also pass parameters explicitly:

CALL system.rewrite_data_files(
  table => 'my_db.my_table',
  options => map(
    'min-input-files', '5',
    'target-file-size-bytes', '536870912'  -- 512MB
  )
);

Or compact only a specific partition, very useful for daily incremental pipelines:

CALL system.rewrite_data_files(
  table => 'my_db.my_table',
  where => 'event_date = DATE "2026-05-20"'
);

Compaction is heuristic-based, not perfect packing. You may still see files with slightly different sizes at the end, that's normal.

3.2 Using the DataFrame API

If you want to trigger compaction programmatically (useful inside an Airflow DAG or a cron maintenance script), you can just wrap the SQL call:

spark.sql("""
  CALL system.rewrite_data_files(
    table => 'my_db.my_table',
    options => map(
      'target-file-size-bytes', '268435456'
    )
  )
""")

3.3 Sort and compact

You can also sort and compact at the same time, which gives you better query performance down the line (more on strategies in section 4):

CALL system.rewrite_data_files(
  table => 'my_db.my_table',
  strategy => 'sort',
  sort_order => 'event_time'
);

A typical Airflow setup might look like this, running on a 1-day or 1-week interval:

def compact():
    spark.sql("""
        CALL system.rewrite_data_files(
            table => 'prod.events',
            where => 'event_date >= current_date - INTERVAL 1 DAY'
        )
    """)

3.4 Advanced parameters

A few options worth knowing about:

  • max-concurrent-file-group-rewrites: ceiling on how many file groups to rewrite in parallel.
  • max-file-group-size-bytes: limits how much data is processed in one rewrite task. Useful to avoid OOM issues on large partitions. Iceberg splits the work into manageable groups.
  • partial-progress-enabled: allows Iceberg to commit results incrementally during compaction instead of waiting for the whole job to finish. Already-compacted groups become visible sooner, and if the job fails, you don't lose everything.
  • partial-progress-max-commits: controls how many intermediate commits Iceberg is allowed to make when partial progress is enabled.

Without partial progress, compaction works like this:

  1. Read many small files
  2. Rewrite them
  3. Commit everything at once at the end

The problem: if the job is very large, a failure means losing ALL progress. When Airflow retries the task, Iceberg re-scans the table, skips already-compacted data, and only processes the remaining small files. Partial progress makes this much more robust.

Running compaction too often can be a waste of compute and cause rewrite churn. Don't forget delete files either; see section 5 on full table maintenance.

3.5 AWS: Athena, EMR, Glue

On AWS you can run compaction through Amazon Athena, or using Spark on Amazon EMR / AWS Glue:

  • Athena uses the OPTIMIZE statement. One nice thing: Athena pricing is based on data scanned, so if there's nothing to compact, there's no cost. Run it per partition to avoid timeouts.
  • For large compaction jobs on EMR or Glue, use dynamic scaling so the cluster adjusts to the workload.
  • For Z-order sorting specifically, prefer EMR or Glue: Z-order is expensive and may need to spill data to disk. Athena is less suited for it.

Rule of thumb: compaction frequency and size

  • Daily → for batch pipelines
  • Every few hours → for streaming pipelines
  • Target file size: 256 MB to 1 GB

4. Compaction strategies

Not all compaction is the same. Iceberg offers three strategies, each with a different tradeoff between cost and query performance. Here's a comparison:

StrategySpeedBest forQuery gain
BinpackFastStreaming SLAs, frequent runsFewer file opens
SortMediumFilters on one columnData skipping via min/max stats
Z-orderSlowMulti-column filtersHighest for multi-dimensional queries

Binpack (default)

Takes many small files and packs them into larger files close to the target size. It doesn't rearrange the data at all, it just packs files together.

Use this when you need fast compaction with a short SLA. For example, running compaction every hour for streaming data to keep read performance acceptable. You can always run a more thorough sort-based compaction on a longer window (e.g. nightly) if needed.

Sort compaction

Compacts files AND sorts the data inside them based on one or more columns. The main benefit: query engines can use Iceberg's per-file min/max statistics to skip entire files when filtering, which can dramatically reduce the amount of data scanned.

To get the most out of sort compaction, you need to understand how your end users are querying the data: what columns do they filter on most? Sort on those.

Slower than binpack but pays off in query performance, especially on large tables with selective filters.

Z-order compaction

Z-order is for when you regularly filter on multiple columns at once. It reorganizes data using a Z-order space-filling curve to cluster rows that are close in multiple dimensions together on disk.

Z-order over age and height. The Z path visits one region at a time: young and tall, old and tall, young and short, old and short. Rows in the same region go into the same file. HEIGHT AGE Age 1-50 Height 5-10ft file 1 Age 51-100 Height 5-10ft file 2 Age 1-50 Height 1-5ft file 3 Age 51-100 Height 1-5ft file 4
The Z path visits one region at a time, so rows close in both age and height land in the same file.

In the diagram above: instead of storing rows randomly, the Z-path groups records from the same region (say ages 1 to 50 and heights 5 to 10 ft) into the same files. A query filtering on both age and height can then skip most of the files entirely.

CALL catalog.system.rewrite_data_files(
  table => 'people',
  strategy => 'sort',
  sort_order => 'zorder(age, height)'
);

If an Iceberg table has a sort order defined in its table properties, even binpack will respect it within each task, so the two aren't mutually exclusive.

Z-order is the most expensive strategy. On AWS, use EMR or Glue rather than Athena for this.

5. Full table maintenance

Compaction is important, but it's only one part of keeping an Iceberg table in good shape. There are three other operations that go hand in hand with it:

5.1 Compacting delete files

Every time you run a DELETE or UPDATE on an Iceberg table (without a full file rewrite), Iceberg writes a delete file instead of touching the datafiles. Over time these accumulate and slow down reads: every query has to check them.

The command to clean them up:

CALL system.rewrite_position_deletes(
  table => 'my_db.my_table'
);

Run this after or alongside regular compaction, ideally on the same schedule.

5.2 Expiring snapshots

Every operation on an Iceberg table creates a new snapshot. These are useful for time travel, but if you never clean them up they pile up forever, bloating your metadata and your storage.

CALL system.expire_snapshots(
  table => 'my_db.my_table',
  older_than => TIMESTAMP '2026-05-19 00:00:00'
);

A typical policy: keep the last 1 to 2 days of snapshots. For most use cases that's more than enough for time travel and rollback.

5.3 Removing orphan files

Sometimes files end up on storage without being referenced by any snapshot. This can happen due to failed writes, aborted jobs, or just bugs. These orphan files just waste space.

CALL system.remove_orphan_files(
  table => 'my_db.my_table'
);

Be careful with this one. Make sure the retention period is long enough to not accidentally delete files that belong to in-progress jobs.

5.4 When NOT to compact

It's easy to set up compaction and forget about it, but running it too aggressively can cause problems:

  • Running compaction on huge tables entirely every time is wasteful; prefer partition-based compaction where possible
  • Compacting too frequently causes rewrite churn: you end up rewriting files that just got compacted
  • Z-order on very large datasets can OOM your executors; use max-file-group-size-bytes to limit the blast radius
  • Always test your compaction schedule against your actual workload before locking it in

My Iceberg notes were giving small files problem energy, so I compacted them into an article. 🧊


Sources

Originally published on LinkedIn on May 26, 2026.

Back home