HeadlinesBriefing favicon HeadlinesBriefing.com

Compacting 1,000 Iceberg Files: Query Performance Impact

Towards Data Science •
×

When you deal with large analytical datasets stored in table formats such as Apache Iceberg, there is a known issue called the small files problem. This happens when data is written in lots of tiny files rather than a smaller number of reasonably sized ones. This increases metadata overhead and can slow down query planning and execution. To combat this, a process called compaction is used which combines all the smaller files into a small number of larger files.

I’m focussing on Apache Iceberg as it’s rapidly growing into one of the leading open table formats for large-scale analytical data, bringing features such as schema evolution, time travel, partition evolution and reliable transactions. Importantly, although Iceberg supports compaction it doesn’t do it for us automatically. We as the system admins still have to execute that procedure or configure another system to trigger it.

We’ll create an Iceberg table containing 50 million rows spread across 1,000 tiny files. Those files will remain untouched until we issue the compaction command. We’ll measure three SQL workloads before and after the rewrite. Everything runs locally. You won’t need a cloud account, Docker, a Hadoop cluster, or a paid service.

Suppose a streaming job writes a small batch every minute. Each batch may produce one or more new files. After a month, a modest quantity of data can be scattered across tens of thousands of objects. Iceberg’s write.target-file-size-bytes property is a “best endeavours” operation, not an absolute promise. The Iceberg documentation explicitly notes that Spark can’t write a file larger than the task producing it.