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. This reduces metadata overhead and gives the query engine fewer files to open, scan and manage, which can improve performance significantly.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 to data stored in object storage or distributed file systems. Importantly, although Iceberg supports compaction it doesn’t do it for us automatically. Iceberg supplies the rewrite_data_files procedure and the metadata needed to select and process files for compaction but we as the system admins still have to execute that procedure or configure another system to trigger it. But, really, does compaction make that much of a difference? That’s the question this article will try to answer. 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. That will tell us if compaction is worth it.Everything runs locally. You won’t need a cloud account, Docker, a Hadoop cluster, or a paid service. The downside of this setup is that we can’t reproduce truly massive datasets that are used in real-world systems but hopefully our results will give us some useful insights.Table of contentsWhat is Apache Iceberg?Setting up a dev environment and installing the required softwareWriting the Python codeConfiguring a local Iceberg catalogueCreating the Spark sessionGenerating deterministic test dataCreating the tableAsking Iceberg about its filesEstablishing our SQL benchmarksPrettify the outputMain driver codeRunning the demoSummaryWhat is Apache Iceberg?Apache Iceberg is an open table format for accessing huge analytic datasets, bringing database-like features such as schema evolution, partition evolution, time travel, and reliable transactions to data stored in files such as Parquet. Parquet is a file format. It determines how rows and columns are encoded inside an individual file. Iceberg operates one level above that.An Iceberg table normally contains Parquet, Avro or ORC data files, plus metadata describing which files currently belong to the table. Its snapshots provide a history of table changes, while manifest files help query engines locate relevant data files without recursively listing every directory.This extra metadata lets engines treat a collection of files more like a database table. Iceberg supports atomic changes, schema evolution, partition evolution and time-travel queries without converting the underlying data into a proprietary storage format.It can’t, however, prevent every poor write pattern.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. A query engine must plan work for those files, open them, read their metadata and close them again.The same thing can happen with batch processing. Spark writes output from its tasks independently, and a file can’t span an Iceberg partition boundary. Iceberg’s write.target-file-size-bytes property is therefore 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. A 512 MB target is irrelevant if a task only has enough data to create a 70 KB Parquet file.Compaction fixes this by reading small files and rewriting their rows into fewer, larger files. This generally means less work is rquired to read and process those files. For this demo, we’re going to test Iceberg’s default bin-pack strategy, which changes the packaging without deliberately sorting the rows.The Apache Iceberg documentation classifies data-file compaction as optional maintenance. It tells us to inspect the files metadata table and run rewriteDataFiles when appropriate. The format provides the operation, but deciding when to run it remains part of operating the table.Setting up a dev environment and installing the required softwareBefore writing any code, let’s set up a development environment to keep the project isolated. I use the uv tool for this, but use whichever method you know best. This is what we’ll need for our experiment. If you already have some or all of these, leave well alone and install just the ones you need.Python 3.10–3.13Java 17 or Java 21PySpark 4.0.3Apache Iceberg 1.11.0Iceberg has no separate Python installation step here. When the Spark session is created, Spark resolves the iceberg-spark-runtime-4.0_2.13:1.11.0 dependency from Maven Central and caches the JAR locally. You’ll see that happen when we run the script later.Writing the Python codeCreate a file named iceberg_compaction_demo.py. The code blocks below form one script and should be added in the order shown. Configuring a local Iceberg catalogueThe 64 MiB target is large enough to make compaction meaningful on a laptop without turning our local test into an all-day job. Iceberg’s normal data-file target is 512 MB.Creating the Spark sessionNext, create the Spark session (you can see the reference to Iceberg 1.11.0 that I talked about earlier):local[*] tells Spark to use the available logical processors. Our catalogue is named local, uses Iceberg’s HadoopCatalog, and stores everything beneath the iceberg_lab_warehouse directory.Despite its name, this catalogue doesn’t require Hadoop to be running. It uses Hadoop’s filesystem interface to manage an ordinary local directory. A Hadoop catalogue on a local filesystem isn’t safe for concurrent writers, but that limitation is OK for this single-process experiment.Adaptive Query Execution is disabled because Spark might otherwise combine our deliberately small tasks. That would be sensible behaviour in a real workload but would spoil the demonstration.On its first run, Spark downloads the approximately 46 MB Iceberg runtime from Maven Central. It caches the JAR, so later runs don’t normally download it again.Generating deterministic test dataAdd this function to your code file:This data is deliberately non-random, so each run creates the same rows and query result. The final repartition forces the DataFrame through the requested number of Spark tasks. Because each task writes independently, asking for 1,000 tasks gives us 1,000 small data files. Setting 64 shuffle partitions keeps the aggregation benchmarks from creating 1,000 result-side tasks.Creating the tableFor repeatability, the code deletes the iceberg_lab_warehouse folder at the beginning of every run. Don’t point the WAREHOUSE environment variable at a directory containing anything you need to keep!The distribution mode is disabled so Iceberg doesn’t reorganise our carefully fragmented input before writing it.Asking Iceberg about its filesCounting files in the directory is unreliable because Iceberg retains old files for historical snapshots. Instead, query the table’s files metadata table, which describes files belonging to the current snapshot:As expected, my run produced 1,000 data files containing 50 million records.Establishing our SQL benchmarksOne query would tell us very little, so the test suite uses three workloads:A filtered aggregation for a range of customersA full-table aggregation grouped by event dateA narrow lookup covering 10,000 consecutive IDsAdd these queries and benchmark function:The unreported first execution of each query lets the JVM initialise and loads Iceberg’s metadata. Five measured executions are more informative than selecting whichever single run supports the argument.Prettify the outputMain driver codeThe result check in this section is important. Compaction must change the table’s physical files without changing a row returned by any query.Running the demoOn my reasonably well-specc’ed desktop, all three median times improved between 31% and 63% after compaction. Not too shabby!Your numbers will differ. We’re testing local file layout, CPU scheduling and filesystem caching as well as Iceberg. The useful comparison is before versus after on the same machine.This isn’t evidence that compaction makes every query 63% faster. The compressed table is still only 380 MiB, and Spark’s fixed job overhead accounts for some of its runtime. The experiment establishes the mechanism: the query engine has fewer files to plan and open. The benefit in a real table depends on its storage, file count, filters, partitions, engine and workload.Compaction also performed a complete read and rewrite of the table. That work isn’t free. A table queried once may never recover the cost of compacting it.SummaryCompaction is an important part of operating Iceberg tables that accumulate large numbers of small files. When small-file fragmentation becomes significant, compaction can substantially improve query performance - but whether and how often it should run depends on the workload.Bear in mind that as soon as you start to write more data and files, the benefits of compaction will be lost over time. This is why compaction requires a policy.You’ll probably need to spend some time and effort ensuring you’re compacting at the right frequency for your specific workloads.A sensible policy depends entirely on your business needs. You could schedule it every night during a quiet period or after a certain number of files are created. Iceberg’s rewrite_data_files procedure accepts a where parameter, so maintenance can target recent or particularly fragmented partitions instead of rewriting an entire table.You also need to consider file size. Larger files reduce file-opening and metadata overhead, but they also reduce read parallelism and make each rewrite more substantial.One final consideration that often causes confusion: those “old” files you just compacted are still there. That’s because Iceberg needs them for doing time-travel queries. Compaction changes the current snapshot; it isn’t snapshot expiration or orphan-file removal. Those are separate maintenance operations with separate retention decisions.The small-file problem therefore has no one-off fix. Iceberg provides the compaction operations, but it can’t know how often a table is queried, how much maintenance capacity is available or how long historical snapshots must survive.Compaction turns many small files into fewer large ones. Data engineering begins with deciding when that rewrite is worth doing.
I Compacted 1,000 Apache Iceberg Files Into 6. Here’s What Happened to Query Performance.
Full Article
Original Source
Read the full article at Towardsdatascience →KhanList aggregates and links to publicly available news content. We do not host full articles from third-party sources. Always verify important information with original sources.