Whenever you cope with massive analytical datasets saved in desk codecs reminiscent of Apache Iceberg, there’s a recognized challenge referred to as the small information downside. This occurs when knowledge is written in a lot of tiny information relatively than a smaller variety of moderately sized ones. This will increase metadata overhead and might decelerate question planning and execution.
To fight this, a course of referred to as compaction is used which mixes all of the smaller information right into a small variety of bigger information. This reduces metadata overhead and offers the question engine fewer information to open, scan and handle, which may enhance efficiency considerably.
I’m focussing on Apache Iceberg because it’s quickly rising into one of many main open desk codecs for large-scale analytical knowledge, bringing options reminiscent of schema evolution, time journey, partition evolution and dependable transactions to knowledge saved in object storage or distributed file techniques.
Importantly, though Iceberg helps compaction it doesn’t do it for us robotically. Iceberg provides the rewrite_data_files process and the metadata wanted to pick and course of information for compaction however we because the system admins nonetheless need to execute that process or configure one other system to set off it.
However, actually, does compaction make that a lot of a distinction? That’s the query this text will attempt to reply. We’ll create an Iceberg desk containing 50 million rows unfold throughout 1,000 tiny information. These information will stay untouched till we challenge the compaction command. We’ll measure three SQL workloads earlier than and after the rewrite. That may inform us if compaction is price it.
All the pieces runs regionally. You gained’t want a cloud account, Docker, a Hadoop cluster, or a paid service. The draw back of this setup is that we will’t reproduce really huge datasets which might be utilized in real-world techniques however hopefully our outcomes will give us some helpful insights.
What’s Apache Iceberg?
Apache Iceberg is an open desk format for accessing large analytic datasets, bringing database-like options reminiscent of schema evolution, partition evolution, time journey, and dependable transactions to knowledge saved in information reminiscent of Parquet. Parquet is a file format. It determines how rows and columns are encoded inside a person file. Iceberg operates one stage above that.
An Iceberg desk usually incorporates Parquet, Avro or ORC knowledge information, plus metadata describing which information at the moment belong to the desk. Its snapshots present a historical past of desk adjustments, whereas manifest information assist question engines find related knowledge information with out recursively itemizing each listing.
This additional metadata lets engines deal with a set of information extra like a database desk. Iceberg helps atomic adjustments, schema evolution, partition evolution and time-travel queries with out changing the underlying knowledge right into a proprietary storage format.
It could’t, nevertheless, forestall each poor write sample.
Suppose a streaming job writes a small batch each minute. Every batch could produce a number of new information. After a month, a modest amount of information will be scattered throughout tens of 1000’s of objects. A question engine should plan work for these information, open them, learn their metadata and shut them once more.
The identical factor can occur with batch processing. Spark writes output from its duties independently, and a file can’t span an Iceberg partition boundary. Iceberg’s write.target-file-size-bytes property is subsequently a “greatest endeavours” operation, not an absolute promise. The Iceberg documentation explicitly notes that Spark can’t write a file bigger than the duty producing it. A 512 MB goal is irrelevant if a activity solely has sufficient knowledge to create a 70 KB Parquet file.
Compaction fixes this by studying small information and rewriting their rows into fewer, bigger information. This typically means much less work is rquired to learn and course of these information. For this demo, we’re going to check Iceberg’s default bin-pack technique, which adjustments the packaging with out intentionally sorting the rows.
The Apache Iceberg documentation classifies data-file compaction as non-compulsory upkeep. It tells us to examine the information metadata desk and run rewriteDataFiles when acceptable. The format offers the operation, however deciding when to run it stays a part of working the desk.
Organising a dev setting and putting in the required software program
Earlier than writing any code, let’s arrange a improvement setting to maintain the challenge remoted. I take advantage of the uv device for this, however use whichever methodology greatest.
REM Create a brand new challenge folder and change to itc:> mkdir C:iceberg-compactionc:> cd /d C:iceberg-compactionREM Replace uv deviceC:iceberg-compaction> uv self replacedata: Checking for updates…success: You are already on model v0.12.5 of uv (the newest model).
That is what we’ll want for our experiment. If you have already got some or all of those, go away effectively alone and set up simply those you want.
Python 3.10–3.13
Java 17 or Java 21
PySpark 4.0.3
Apache Iceberg 1.11.0
REM Set up JAVAc:iceberg-compaction> winget set up EclipseAdoptium.Temurin.21.JDKDiscovered Eclipse Temurin JDK with Hotspot 21 [EclipseAdoptium.Temurin.21.JDK] Model 21.0.12.101This utility is licensed to you by its proprietor.Microsoft isn't liable for, nor does it grant any licenses to, third-party packages.Downloading https://github.com/adoptium/temurin21-binaries/releases/obtain/jdk-21.0.12.1+1/OpenJDK21U-jdk_x64_windows_hotspot_21.0.12.1_1.msi ██████████████████████████████ 171 MB / 171 MBEfficiently verified installer hashBeginning bundle set up...Efficiently put inC:iceberg-compaction> java -versionopenjdk model "21.0.12.1" 2026-08-18 LTSOpenJDK Runtime Atmosphere Temurin-21.0.12.1+1 (construct 21.0.12.1+1-LTS)OpenJDK 64-Bit Server VM Temurin-21.0.12.1+1 (construct 21.0.12.1+1-LTS, blended mode, sharing)REM Create and change to a brand new setting with Python 3.13C:iceberg-compaction> uv initInitialized challenge `iceberg-compaction`C:iceberg-compaction> uv venv --python 3.13C:iceberg-compaction> .venvScriptsactivate(iceberg-compaction) C:iceberg-compaction> python --versionPython 3.13.1REM Set up Spark(iceberg-compaction) C:iceberg-compaction> uv add pyspark==4.0.3Resolved 3 packages in 42.96s Constructed pyspark==4.0.3 Constructed iceberg-compaction @ file:///C:/Customers/thoma/initiatives/iceberg-compaction Ready 3 packages in 33.86s░░░░░░░░░░░░░░░░░░░░ [0/3] Putting in wheels... warning: Did not hardlink information; falling again to full copy. This will result in degraded efficiency. If the cache and goal directories are on completely different filesystems, hardlinking will not be supported. If that is intentional, set `export UV_LINK_MODE=copy` or use `--link-mode=copy` to suppress this warning.Put in 3 packages in 2.28s + iceberg-compaction==0.1.0 (from file:///C:/Customers/thoma/initiatives/iceberg-compaction) + py4j==0.10.9.9 + pyspark==4.0.3
Iceberg has no separate Python set up step right 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 regionally. You’ll see that occur once we run the script later.
Writing the Python code
Create a file named iceberg_compaction_demo.py. The code blocks under kind one script and must be added within the order proven.
The 64 MiB goal is massive sufficient to make compaction significant on a laptop computer with out turning our native take a look at into an all-day job. Iceberg’s regular data-file goal is 512 MB.
Creating the Spark session
Subsequent, create the Spark session (you may see the reference to Iceberg 1.11.0 that I talked about earlier):
native[*] tells Spark to make use of the accessible logical processors. Our catalogue is known as native, makes use of Iceberg’s HadoopCatalog, and shops every little thing beneath the iceberg_lab_warehouse listing.
Regardless of its identify, this catalogue doesn’t require Hadoop to be working. It makes use of Hadoop’s filesystem interface to handle an peculiar native listing. A Hadoop catalogue on an area filesystem isn’t protected for concurrent writers, however that limitation is OK for this single-process experiment.
Adaptive Question Execution is disabled as a result of Spark would possibly in any other case mix our intentionally small duties. That will be smart behaviour in an actual workload however would spoil the demonstration.
On its first run, Spark downloads the roughly 46 MB Iceberg runtime from Maven Central. It caches the JAR, so later runs don’t usually obtain it once more.
This knowledge is intentionally non-random, so every run creates the identical rows and question consequence. The ultimate repartition forces the DataFrame by way of the requested variety of Spark duties. As a result of every activity writes independently, asking for 1,000 duties offers us 1,000 small knowledge information. Setting 64 shuffle partitions retains the aggregation benchmarks from creating 1,000 result-side duties.
Creating the desk
if WAREHOUSE.exists(): shutil.rmtree(WAREHOUSE)spark = build_spark()spark.sparkContext.setLogLevel("WARN")spark.sql("CREATE NAMESPACE IF NOT EXISTS native.lab")spark.sql( f""" CREATE TABLE {TABLE} ( id BIGINT, customer_id INT, event_type STRING, event_date DATE, quantity DECIMAL(10, 2) ) USING iceberg TBLPROPERTIES ( 'write.distribution-mode' = 'none', 'write.target-file-size-bytes' = '{TARGET_FILE_SIZE}' ) """)make_events(spark, 0, ROWS, INITIAL_FILES).writeTo(TABLE).append()
For repeatability, the code deletes the iceberg_lab_warehouse folder initially of each run. Don’t level the WAREHOUSE setting variable at a listing containing something you must hold!
The distribution mode is disabled so Iceberg doesn’t reorganise our fastidiously fragmented enter earlier than writing it.
Asking Iceberg about its information
Counting information within the listing is unreliable as a result of Iceberg retains outdated information for historic snapshots. As a substitute, question the desk’s information metadata desk, which describes information belonging to the present snapshot:
def file_statistics(spark: SparkSession, label: str) -> dict: print(f"n{label}") consequence = spark.sql( f""" SELECT COUNT(*) AS data_files, SUM(record_count) AS data, ROUND(SUM(file_size_in_bytes) / 1048576.0, 2) AS total_mib, ROUND(AVG(file_size_in_bytes) / 1024.0, 2) AS average_kib, ROUND(MIN(file_size_in_bytes) / 1024.0, 2) AS smallest_kib, ROUND(MAX(file_size_in_bytes) / 1024.0, 2) AS largest_kib FROM {TABLE}.information WHERE content material = 0 """ ) row = consequence.first() print( f"{int(row['data_files']):,} energetic information, " f"{int(row['records']):,} data, " f"{float(row['total_mib']):,.2f} MiB" ) return row.asDict()
As anticipated, my run produced 1,000 knowledge information containing 50 million data.
Establishing our SQL benchmarks
One question would inform us little or no, so the take a look at suite makes use of three workloads:
A filtered aggregation for a variety of shoppers
A full-table aggregation grouped by occasion date
A slim lookup protecting 10,000 consecutive IDs
Add these queries and benchmark perform:
QUERIES = { "Filtered buyer aggregation": f""" SELECT event_type, COUNT(*) AS occasions, ROUND(SUM(CAST(quantity AS DOUBLE)), 2) AS total_amount FROM {TABLE} WHERE customer_id BETWEEN 1000 AND 1999 GROUP BY event_type ORDER BY event_type """, "Full-table every day aggregation": f""" SELECT event_date, COUNT(*) AS occasions, ROUND(AVG(CAST(quantity AS DOUBLE)), 2) AS average_amount FROM {TABLE} GROUP BY event_date ORDER BY event_date """, "Slender ID-range lookup": f""" SELECT COUNT(*) AS occasions, ROUND(SUM(CAST(quantity AS DOUBLE)), 2) AS total_amount FROM {TABLE} WHERE id BETWEEN 500000 AND 509999 """,}def benchmark(spark: SparkSession, label: str, repetitions: int = 5): print(f"n{label}") measurements = {} for identify, question in QUERIES.gadgets(): # One unreported run warms the JVM and reads the desk metadata. anticipated = spark.sql(question).acquire() timings = [] for _ in vary(repetitions): spark.catalog.clearCache() began = time.perf_counter() precise = spark.sql(question).acquire() timings.append(time.perf_counter() - began) if precise != anticipated: elevate RuntimeError(f"{identify} returned inconsistent outcomes") median = statistics.median(timings) measurements[name] = {"rows": anticipated, "median": median} print(f"n{identify} ({len(anticipated)} consequence rows)") print("Occasions (seconds):", ", ".be part of(f"{worth:.3f}" for worth in timings)) print(f"Median: {median:.3f} seconds") return measurements
The unreported first execution of every question lets the JVM initialise and hundreds Iceberg’s metadata. 5 measured executions are extra informative than choosing whichever single run helps the argument.
Prettify the output
def format_average_file_size(value_kib) -> str: value_kib = float(value_kib) if value_kib >= 1024: return f"{value_kib / 1024:,.2f} MiB" return f"{value_kib:,.2f} KiB"def print_comparison(before_files, after_files, before_queries, after_queries): rows = [ ( "Active data files", f"{int(before_files['data_files']):,}", f"{int(after_files['data_files']):,}", ), ( "Information", f"{int(before_files['records']):,}", f"{int(after_files['records']):,}", ), ( "Complete energetic knowledge dimension", f"{float(before_files['total_mib']):,.2f} MiB", f"{float(after_files['total_mib']):,.2f} MiB", ), ( "Common file dimension", format_average_file_size(before_files["average_kib"]), format_average_file_size(after_files["average_kib"]), ), ] for identify in QUERIES: rows.append( ( identify, f"{before_queries[name]['median']:.3f} s", f"{after_queries[name]['median']:.3f} s", ) ) headers = ("Measurement", "Earlier than", "After") widths = [ max(len(headers[index]), *(len(row[index]) for row in rows)) for index in vary(3) ] border = "+" + "+".be part of("-" * (width + 2) for width in widths) + "+" def print_row(row): cells = [row[index].ljust(widths[index]) for index in vary(3)] print("| " + " | ".be part of(cells) + " |") print("nBefore/after comparability") print(border) print_row(headers) print(border) for row in rows: print_row(row) print(border)
Principal driver code
def major() -> None: print(f"PySpark goal model: {SPARK_VERSION}") print(f"Iceberg model: {ICEBERG_VERSION}") print(f"Warehouse: {WAREHOUSE}") print("The warehouse listing is deleted and recreated on each run.") if WAREHOUSE.exists(): shutil.rmtree(WAREHOUSE) spark = build_spark() spark.sparkContext.setLogLevel("WARN") attempt: print(f"Working Spark {spark.model}") if spark.model != SPARK_VERSION: print( f"WARNING: this experiment was written for Spark {SPARK_VERSION}, " f"however {spark.model} is working." ) spark.sql("CREATE NAMESPACE IF NOT EXISTS native.lab") spark.sql(f"DROP TABLE IF EXISTS {TABLE}") spark.sql( f""" CREATE TABLE {TABLE} ( id BIGINT, customer_id INT, event_type STRING, event_date DATE, quantity DECIMAL(10, 2) ) USING iceberg TBLPROPERTIES ( 'write.distribution-mode' = 'none', 'write.target-file-size-bytes' = '{TARGET_FILE_SIZE}' ) """ ) print(f"nWriting {ROWS:,} rows by way of {INITIAL_FILES} Spark duties...") make_events(spark, 0, ROWS, INITIAL_FILES).writeTo(TABLE).append() before_files = file_statistics(spark, "Earlier than compaction") earlier than = benchmark(spark, "Earlier than compaction") print("nCompacting the desk with Iceberg's bin-pack technique...") compaction = spark.sql( f""" CALL native.system.rewrite_data_files( desk => 'native.lab.occasions', technique => 'binpack', choices => map( 'target-file-size-bytes', '{TARGET_FILE_SIZE}', 'min-input-files', '2' ) ) """ ) compaction.present(truncate=False) after_files = file_statistics(spark, "After compaction") after = benchmark(spark, "After compaction") for identify in QUERIES: if earlier than[name]["rows"] != after[name]["rows"]: elevate RuntimeError(f"{identify} modified after compaction") print_comparison(before_files, after_files, earlier than, after) print( "nFinished. The present desk is undamaged in iceberg_lab_warehouse. " "Outdated information are additionally retained as a result of Iceberg snapshots nonetheless refer " "to them." ) lastly: spark.cease()if __name__ == "__main__": major()
The consequence examine on this part is necessary.
......for identify in QUERIES: if earlier than[name]["rows"] != after[name]["rows"]: elevate RuntimeError(f"{identify} modified after compaction")......
Compaction should change the desk’s bodily information with out altering a row returned by any question.
Working the demo
(iceberg-compaction) C:iceberg-compaction> python iceberg_compaction_demo.pyPySpark goal model: 4.0.3Iceberg model: 1.11.0Warehouse: D:iceberg-compactioniceberg_lab_warehouseThe warehouse listing is deleted and recreated on each run.WARNING: Utilizing incubator modules: jdk.incubator.vector......Writing 50,000,000 rows by way of 1000 Spark duties...Earlier than compaction1,000 energetic information, 50,000,000 data, 380.29 MiB......After compaction6 energetic information, 50,000,000 data, 357.54 MiBAfter compactionFiltered buyer aggregation (4 consequence rows)Occasions (seconds): 0.435, 0.448, 0.446, 0.574, 0.431Median: 0.446 secondsFull-table every day aggregation (31 consequence rows)Occasions (seconds): 0.582, 0.573, 0.644, 0.567, 0.563Median: 0.573 secondsSlender ID-range lookup (1 consequence rows)Occasions (seconds): 0.964, 0.402, 0.319, 0.300, 0.387Median: 0.387 secondsEarlier than/after comparability+-------------------------------+------------+------------+| Measurement | Earlier than | After |+-------------------------------+------------+------------+| Energetic knowledge information | 1,000 | 6 || Information | 50,000,000 | 50,000,000 || Complete energetic knowledge dimension | 380.29 MiB | 357.54 MiB || Common file dimension | 389.42 KiB | 59.59 MiB || Filtered buyer aggregation | 1.013 s | 0.446 s || Full-table every day aggregation | 1.580 s | 0.573 s || Slender ID-range lookup | 0.564 s | 0.387 s |+-------------------------------+------------+------------+
On my moderately well-specc’ed desktop, all three median instances improved between 31% and 63% after compaction. Not too shabby!
Your numbers will differ. We’re testing native file format, CPU scheduling and filesystem caching in addition to Iceberg. The helpful comparability is earlier than versus after on the identical machine.
This isn’t proof that compaction makes each question 63% sooner. The compressed desk continues to be solely 380 MiB, and Spark’s mounted job overhead accounts for a few of its runtime. The experiment establishes the mechanism: the question engine has fewer information to plan and open. The profit in an actual desk depends upon its storage, file depend, filters, partitions, engine and workload.
Compaction additionally carried out a whole learn and rewrite of the desk. That work isn’t free. A desk queried as soon as could by no means get well the price of compacting it.
Abstract
Compaction is a vital a part of working Iceberg tables that accumulate massive numbers of small information. When small-file fragmentation turns into important, compaction can considerably enhance question efficiency – however whether or not and the way usually it ought to run depends upon the workload.
Keep in mind that as quickly as you begin to write extra knowledge and information, the advantages of compaction will likely be misplaced over time. This is the reason compaction requires a coverage.
You’ll in all probability have to spend some effort and time guaranteeing you’re compacting on the proper frequency to your particular workloads.
A smart coverage relies upon completely on what you are promoting wants. You may schedule it each evening throughout a quiet interval or after a sure variety of information are created. Iceberg’s rewrite_data_files process accepts a the place parameter, so upkeep can goal current or significantly fragmented partitions as a substitute of rewriting a complete desk.
You additionally want to think about file dimension. Bigger information cut back file-opening and metadata overhead, however additionally they cut back learn parallelism and make every rewrite extra substantial.
One remaining consideration that usually causes confusion: these “outdated” information you simply compacted are nonetheless there. That’s as a result of Iceberg wants them for doing time-travel queries. Compaction adjustments the present snapshot; it isn’t snapshot expiration or orphan-file removing. These are separate upkeep operations with separate retention choices.
The small-file downside subsequently has no one-off repair. Iceberg offers the compaction operations, however it could possibly’t understand how usually a desk is queried, how a lot upkeep capability is out there or how lengthy historic snapshots should survive.
Compaction turns many small information into fewer massive ones. Knowledge engineering begins with deciding when that rewrite is price doing.