skip to content

Compaction and the Small-Files Problem

You will learn why streaming and micro-batch writes shred a table into thousands of tiny files, what that does to planning time, object-store request cost and scan throughput, and how compaction, clustering and file-size targets fix it. This is the most common real operational question about lakehouse tables.

on this pageshow

explore

questions

5

Why do thousands of tiny data files slow down queries on a lakehouse table?

level: juniorimportance: must knowfreq 78%

answer

  1. overhead is charged per file, not per byte
  2. planning must look at every file entry
  3. object stores bill and delay every request
  4. small files mean small row groups and useless statistics
  5. a thousand tiny tasks instead of a few large ones

basics

~20 s

Each file costs fixed overhead: a metadata entry to plan, an object-store request to open, and a footer to parse. With thousands of tiny files that per-file overhead dominates, so the engine spends its time on bookkeeping rather than reading rows.

solid answer

~40 s

A lakehouse table is a set of immutable data files plus a metadata layer that lists them. Cost is roughly *per file*, not per byte. Query planning must read metadata for every file to decide which to skip, so planning time grows linearly with file count. Then each surviving file needs at least one object-store GET (often two — footer, then column data), and every request carries latency measured in tens of milliseconds. Columnar compression and encoding also work better with more rows per file: a 4 MB file has tiny row groups, weak dictionary reuse and per-file statistics too coarse to prune. Ten thousand 1 MB files and a hundred 100 MB files hold the same data, but the first is dominated by overhead and the second by useful I/O.

go deeper

for a junior

Be ready to say that cost is charged per file — planning entry, storage request, footer parse — so many tiny files means mostly overhead. Name streaming or frequent small writes as the usual source.

for a middle

Explain the four distinct costs: planning that scales with file count, per-request object-store latency and billing, degraded columnar encoding and useless min/max statistics, and per-file task scheduling. Connect each to a concrete symptom.

for a senior

An interviewer expects you to diagnose from evidence — long planning times before any data is read, request counts far exceeding bytes read, poor pruning — and to tie the file count back to a specific writer's commit interval or parallelism.

for a principal

Own the policy question: what commit cadence, writer parallelism and partition grain a platform allows, and what compaction service runs against every table by default so no team has to remember it.

## What the small-files problem is An open lakehouse table is two layers: a **data-file layer** (immutable columnar files, typically Parquet or ORC, sitting in object storage) and a **table layer** (metadata that says which files currently constitute the table, plus per-file statistics). The small-files problem is what happens when the same volume of data is spread over far more files than it needs to be — thousands or millions of files of a few kilobytes to a few megabytes each, instead of hundreds of files of a few hundred megabytes. Nothing is *incorrect* about a table in this state. Every row is present and readable. The damage is entirely in cost and latency, and it shows up in four separate places. ## Cost 1 — query planning grows with file count Before reading a single row, the engine must decide which files it needs. To do that it reads the table's metadata: the list of current files, their partition values, and their column-level min/max statistics. That list has one entry per file. Planning a query over ten thousand files means evaluating ten thousand entries; over a million files it means evaluating a million. Planning that should take milliseconds starts taking tens of seconds, and it takes that long *even for a query that ends up reading almost nothing*, because pruning still has to inspect every candidate. ## Cost 2 — object storage charges and delays per request Object stores are not filesystems. Every file open is an HTTP request with latency typically in the tens of milliseconds, and cloud providers bill per thousand GET requests. A columnar file usually needs at least two round trips: one to fetch the footer holding the schema and row-group statistics, and one or more to fetch the actual column chunks. So a scan over 50,000 small files issues on the order of 100,000 requests before it has read anything useful. The wall-clock cost is real even with high parallelism, and the request bill can quietly exceed the storage bill. ## Cost 3 — columnar encoding degrades Columnar formats earn their speed from having many rows per column chunk: dictionary encoding reuses values, run-length encoding compresses repeated runs, and per-chunk statistics are meaningful. A 2 MB file might hold a single small row group covering a few tens of thousands of rows. Dictionaries barely amortize, compression ratios fall, and — most damaging — the per-file min/max statistics span nearly the whole value range because a small random slice of the data tends to contain both very low and very high values. The result is that *pruning stops working*: the engine cannot skip files, so it reads all of them. ## Cost 4 — task scheduling overhead Distributed engines usually assign work per file or per split. Thousands of tiny files become thousands of tiny tasks, each with scheduler dispatch cost, JVM or process warm-up, and a result to collect. The job spends more time launching tasks than executing them, and the cluster looks busy while doing very little. ## Where the tiny files come from They are almost never authored deliberately. The usual sources are: - **Streaming or micro-batch ingestion.** A job committing every 30 seconds writes at least one file per partition per commit. A table with 20 partitions written every 30 seconds produces roughly 57,600 files a day. - **High writer parallelism.** If 200 tasks each write their own output and the batch is small, each file is small — file count is bounded below by the number of writers. - **Over-partitioning.** Partitioning by a high-cardinality column shreds each batch across many directories, so every partition gets a sliver. - **Row-level updates and deletes.** Each modification rewrites or adds files; frequent small changes multiply the file count. ## The fix, in one sentence Because data files are immutable, you do not append into an existing file — you run a **compaction** job that reads many small files, writes a smaller number of appropriately sized ones, and atomically swaps the table's file list to point at the new files. The old files stay on disk until a separate retention-aware cleanup removes them. ```sql -- A frequently-committed streaming sink and a compaction job that -- runs against the same table on a schedule are the normal pairing. SELECT count(*) FROM events; -- unchanged before and after compaction ``` The key mental model: **the number of files is an operational property of a table, independent of its contents, and it must be managed.** A table is not "done" when the data is correct; it is done when the data is correct *and* laid out in files large enough that overhead is negligible.

  • Why can't the engine just concatenate small files at read time instead of compacting?
    Reading is where the cost already lands — the engine would still issue every request and parse every footer before it could concatenate anything. Merging at read time repeats the whole cost on every query, whereas compaction pays it once and every subsequent query benefits. Compaction is caching the merge result durably.
  • Does the small-files problem affect writes as well as reads?
    Yes, indirectly. Every commit must produce a new metadata state that describes the current file set, so as file count grows the metadata itself gets larger and more expensive to write and to read on the next commit. Very large file lists also make concurrent-commit conflict checking slower.
  • Is one giant file per table the logical conclusion?
    No. Files that are too large hurt too: they limit read parallelism (fewer splits than available workers), make any row-level change rewrite far more data than necessary, and lengthen the tail of a failed task's retry. The target is a band — commonly a few hundred megabytes — not a maximum.

Shipping a warehouse of goods in ten thousand envelopes instead of a hundred pallets: the contents are identical, but you pay postage and handling ten thousand times.

saying these in an interview costs you the question

  • Claiming small files only waste storage space, not query time
  • Saying compression makes file count irrelevant
  • Believing the engine transparently merges small files at read time
  • Assuming more files always means more parallelism and therefore faster
  • Blaming the query engine rather than the table's physical layout

context

open as a page

What does bin-packing compaction do to a table's data files, and how is the target size chosen?

level: middleimportance: must knowfreq 72%

basics

~20 s

Bin-packing compaction groups many small files into bins that add up to a target size, rewrites each bin as one new file, and atomically swaps the table's file list. It preserves row content and order-agnostic semantics; the target is a band, commonly a few hundred megabytes.

open as a page

What is write amplification in a lakehouse table, and what raises it?

level: middleimportance: should knowfreq 50%

basics

~20 s

Write amplification is the ratio of bytes physically written to bytes logically changed. Immutable files force whole-file rewrites, so changing one row in a 512 MB file writes 512 MB. Larger files, frequent small updates, and repeated compaction of the same data all raise it.

open as a page

After compaction commits, why do the old data files still occupy storage, and how are they removed?

level: seniorimportance: should knowfreq 48%

basics

~20 s

Compaction only removes files from the table's current file list; the objects stay because older table states still reference them and in-flight readers are still reading them. A separate retention-aware cleanup deletes files that no retained state references, after a safety window.

open as a page

When is sorting or clustering during a rewrite worth its cost compared with plain bin-packing?

level: seniorimportance: should knowfreq 56%

basics

~20 s

Bin-packing fixes file count; sorting fixes pruning. Pay for a sort-based rewrite only when queries filter on a column whose values are scattered across every file, so per-file statistics overlap and nothing can be skipped. It costs a full shuffle.

open as a page