Twelve million Parquet files at 3 MB each is 36 TB of data. It is also a table that will make your query planner sweat before it reads a single row. The failure mode is not a crash. It is a slow, expensive drift: planning times climb, manifest lists grow, and every compaction job you run seems to cost more than the last one.
This article prices rewrite_data_files for that shape of table and gives you a diagnostic you can run on your own catalog before you schedule anything.
What the procedure actually does
Iceberg’s rewrite_data_files is a Spark action that reads groups of small data files and writes them back as larger files. The Iceberg maintenance documentation describes it as combining small files into larger files to reduce metadata overhead and runtime file open cost. It is not a metadata-only operation. It reads Parquet bytes and writes new Parquet bytes.
The procedure exposes options that control how much it reads and writes in one run. The documented options include target-file-size-bytes, min-input-files, max-concurrent-file-group-rewrites, partial-progress, and delete-file-threshold. The Iceberg docs show target-file-size-bytes being set explicitly in the Spark action example, and the Snowflake Iceberg documentation points at delete-file-threshold and rewrite-all as the options that matter when position deletes accumulate.
Two of those options decide your bill.
min-input-files sets the smallest group of files the procedure will consider worth rewriting. If you leave it low, the procedure will happily rewrite pairs of 3 MB files into 6 MB files. That is a bad trade. You paid a full read and write to move from 12 million files to 6 million files, and you still have a small-file problem.
target-file-size-bytes sets the size of the output. If you set it to 512 MB, each rewrite group produces roughly one 512 MB file per partition. That is the number that determines how many files you have after the job, and therefore how many times you need to run the job again.
partial-progress controls whether the procedure commits work as it goes or holds everything until the end. Partial progress reduces the blast radius of a failed run, but it also produces more snapshots and more commits. Each commit is a metadata write. On a table with millions of files, metadata writes are not free.
The arithmetic of 12 million files
Assume 36 TB of data in 12 million files, average 3 MB. Assume you set target-file-size-bytes to 512 MB and min-input-files high enough that the procedure only touches groups that will produce at least one full target file.
At 512 MB per output file, 36 TB becomes roughly 72,000 files. That is a 166x reduction in file count. It is also a full read of 36 TB and a full write of 36 TB.
On object storage, that is 72 TB of I/O. At S3 standard pricing, data transfer out is not the issue for same-region reads, but request costs are. A 3 MB file read is one GET. Twelve million GETs is twelve million requests. S3 charges per 1,000 requests, so the request cost alone is in the low hundreds of dollars for the reads, plus the PUTs for the new files. The dominant cost is not the request line item. It is the compute hours spent reading and writing 72 TB.
If your Spark cluster sustains 2 GB/s of aggregate read throughput, 36 TB takes about 5 hours. If it sustains 500 MB/s, it takes about 20 hours. The write side is usually slower because Parquet encoding and compression are CPU-bound. A realistic single full rewrite of this table is a multi-hour job on a cluster large enough to matter.
That is the number to put in front of whoever approves the cluster budget. It is not a nightly job. It is a quarterly job at best, and only if the write pattern that created the small files has been fixed.
Why the small files exist
Before you schedule a rewrite, find out what is producing 3 MB files. The Iceberg maintenance docs are explicit that streaming queries may produce small data files that should be compacted. If your ingestion is a Kafka Connect sink or a Flink job committing every minute, you will get small files by design. Rewriting them without changing the commit interval is a treadmill.
The diagnostic is the files metadata table. The Iceberg docs call it useful for inspecting data file sizes and determining when to compact partitions. In Athena, the equivalent is SELECT * FROM "dbname"."tablename$files". In Trino, the Iceberg connector exposes metadata tables you can query directly.
Run this against your own table and look at the distribution, not the average:
SELECT
partition,
count(*) AS file_count,
sum(file_size_in_bytes) / 1024 / 1024 AS total_mb,
avg(file_size_in_bytes) / 1024 / 1024 AS avg_mb,
min(file_size_in_bytes) / 1024 / 1024 AS min_mb,
max(file_size_in_bytes) / 1024 / 1024 AS max_mb
FROM my_table.files
GROUP BY partition
ORDER BY file_count DESC
LIMIT 20;
If the top partitions have hundreds of thousands of files and the bottom partitions have a handful, you do not need a full-table rewrite. You need a partition-filtered rewrite. The Iceberg Spark action supports a filter expression, and the docs show it being used with an equality predicate on a date column. That is the difference between a 20-hour job and a 20-minute job.
Manifests are the other half of the cost
File count does not only affect data reads. The Iceberg spec describes the metadata tree as an index over the table’s data. Manifests track data files, and manifest lists track manifests. The spec’s stated goal is O(1) remote calls to plan a scan, not O(n) where n grows with the number of files.
That goal holds only if manifests are reasonably sized. The Iceberg maintenance docs note that more data files leads to more metadata stored in manifest files, and that small data files cause an unnecessary amount of metadata and less efficient queries from file open costs. They also document rewriteManifests as a separate operation for regrouping data files into manifests when the write pattern does not align with the query pattern.
This matters for your rewrite budget because rewrite_data_files does not fix manifest layout. It produces new data files, and those new files get recorded in new manifests. If your existing manifests are already fragmented, you may need rewriteManifests after the data rewrite. That is a second job with its own runtime.
The docs show rewriteManifests being called with a rewriteIf predicate on manifest file length, for example rewriting manifests smaller than 10 MB. That is a cheap operation compared to rewriting data, but it is not zero. It reads and writes manifest files, and on a table with millions of files, there are a lot of manifests.
Delete files change the calculus
If your table uses merge-on-read deletes, rewrite_data_files has a second job: applying delete files to data files. The Iceberg spec defines two delete types. Position deletes mark a row deleted by data file path and row position. Equality deletes mark a row deleted by column values.
The Snowflake Iceberg documentation is direct about the operational consequence: excessive position deletes, especially dangling position deletes, might prevent table creation and refresh operations, and the recommended fix is table maintenance using rewrite_data_files with delete-file-threshold or rewrite-all.
That means your rewrite frequency is not just a function of file size. It is a function of delete volume. A table with heavy update traffic accumulates position deletes faster than it accumulates small files. If you set delete-file-threshold too high, readers pay the merge cost on every query. If you set it too low, you rewrite data files that have only a few deletes, which is wasted I/O.
The diagnostic is the delete file count per data file. Query the metadata tables for delete file counts and compare them to data file counts. If the ratio is climbing, your rewrite cadence needs to be tied to delete accumulation, not to a calendar.
What it costs in hours and dollars
Here is a priced example. These are illustrative numbers based on the arithmetic above, not benchmarks from a specific cluster. Substitute your own throughput and instance pricing.
A full rewrite of 36 TB at 2 GB/s aggregate read and 1 GB/s aggregate write takes roughly 5 hours of read time and 10 hours of write time, overlapped to about 10 hours wall clock. A 100-node cluster of 16 vCPU instances at $0.50 per instance-hour costs $500 for the run. Add S3 request costs in the low hundreds. Call it $700 to $900 per full rewrite.
If you run that monthly, it is $8,400 to $10,800 per year. If you run it weekly, it is $36,000 to $47,000 per year. The difference between monthly and weekly is not a best-practice question. It is a budget question, and the answer depends on how fast your write pattern recreates the problem.
The cheaper path is to fix the write side. If you can raise the commit interval on the streaming job that produces 3 MB files, you reduce the rate at which small files appear. That is a config change, not a cluster run. The Iceberg docs describe the problem as streaming queries producing small data files. The fix is usually in the writer, not in the maintenance job.
How often to run it
There is no universal cadence. There is a decision procedure.
First, measure the file count and the average file size per partition. If the average is below your target and the count is growing, you have a compaction need.
Second, measure the query planning time. If planning is a small fraction of total query time, compaction is not urgent. If planning dominates, it is.
Third, measure the delete file ratio. If deletes are accumulating faster than files, your cadence is set by deletes.
Fourth, price the rewrite. If the rewrite costs more than the query slowdown it prevents, do not run it. That is the tradeoff. A table with 12 million files that is queried twice a day may not justify a $900 monthly rewrite. A table with 12 million files that is queried continuously by a dashboard may justify it weekly.
The Iceberg docs recommend regularly expiring snapshots and cleaning metadata files, and they note that tables with frequent commits may need to regularly clean metadata files. That is a separate maintenance tax from data compaction. Budget for both.
What to run before you schedule anything
Run the file distribution query above. Run a query against the manifests metadata table to see manifest count and sizes. Run a query against the snapshots table to see how many snapshots you are retaining. Check write.metadata.previous-versions-max and write.metadata.delete-after-commit.enabled in your table properties. The Iceberg docs give the defaults as 100 and false, and they show that with deletion disabled, orphaned metadata files accumulate and can only be cleaned with orphan file deletion.
Then run a partition-filtered rewrite on your worst partition and time it. Multiply by the number of partitions you actually need to fix. That is your real cost. It is almost always less than a full-table rewrite, and it is almost always more than the estimate you gave your manager.
FAQ
Does rewrite_data_files rewrite the whole table by default?
No. It operates on file groups selected by the procedure’s options. You can filter by partition. The Iceberg Spark action supports a filter expression, and the docs show it being used with an equality predicate.
What is a safe target-file-size-bytes?
The Iceberg docs show 500 MB in their example. Trino’s Iceberg connector defaults iceberg.target-max-file-size to 1 GB. The right value depends on your query engine’s split size and your object store’s read characteristics. Larger files reduce metadata overhead but increase the cost of rewriting any single file.
Can I run it while writes are happening?
Iceberg uses optimistic concurrency. The spec says a change that rewrites files can be applied to a new table snapshot if all of the rewritten files are still in the table. If a concurrent write removes or replaces a file in your rewrite group, the commit may fail and retry. On a busy table, expect retries.
Do I need rewriteManifests too?
Only if your manifest layout is misaligned with your query pattern. The Iceberg docs describe rewriteManifests as regrouping data files into manifests when the write pattern does not align with the query pattern. It is a separate operation from data compaction.
How do I know if delete files are the problem?
Query the metadata tables for delete file counts. The Snowflake Iceberg docs warn that excessive position deletes can prevent table creation and refresh, and recommend rewrite_data_files with delete-file-threshold or rewrite-all. If your delete file count is climbing faster than your data file count, compaction cadence should follow deletes.
Sources
Apache Iceberg maintenance documentation: https://iceberg.apache.org/docs/latest/maintenance/
Apache Iceberg Spark procedures documentation: https://iceberg.apache.org/docs/latest/spark-procedures/
Apache Iceberg table specification: https://iceberg.apache.org/spec/
Amazon Athena Iceberg table documentation: https://docs.aws.amazon.com/athena/latest/ug/querying-iceberg-table-data.html
Trino Iceberg connector documentation: https://trino.io/docs/current/connector/iceberg.html
Snowflake Iceberg tables documentation: https://docs.snowflake.com/en/user-guide/tables-iceberg





