skip to content

After rewrite_data_files, an Iceberg table still has thousands of small files. Why?

level: seniorimportance: should knowfreq 56%

answer

  1. compaction has thresholds, not just a target
  2. groups are formed inside one boundary only
  3. five candidate files is the default bar
  4. a partition cannot be merged with its neighbour
  5. min-input-files, target size, where, rewrite-all

basics

~20 s

Usually the job skipped most groups: file groups are built per partition, and a partition with fewer than min-input-files candidates or with files already near the target size is left alone. A where filter, an unchanged small target size, or fresh ingest after the run explain the rest.

solid answer

~50 s

`rewrite_data_files` does not compact a table globally. It plans **file groups within each partition**, and a group is only rewritten if it clears the thresholds: `min-input-files` (default 5) small files in that partition, or files far enough from `target-file-size-bytes` (defaulting to the table's `write.target-file-size-bytes`, 512 MB). So a table partitioned by hour with three tiny files per hour is skipped everywhere, and no setting merges files across partition boundaries — that needs a partition-spec change, not compaction. Check the rest of the list too: a `where` argument limited the scope, `rewrite-all` was not set when you wanted unconditional rewriting, `partial-progress` committed only some groups before the job ended, or a streaming writer created new small files right after the run. Read the procedure's returned counts of rewritten files and groups rather than assuming it did nothing.

code

sql · 12 lines
sql
-- lower the bar and cap the target so small partitions become eligible
CALL prod.system.rewrite_data_files(
  table => 'db.events',
  strategy => 'binpack',
  where => 'event_date >= date ''2026-08-01''',
  options => map(
    'min-input-files', '2',
    'target-file-size-bytes', '268435456',
    'partial-progress.enabled', 'true',
    'max-concurrent-file-group-rewrites', '10'
  )
);

go deeper

for a junior

Recall that compaction combines many small data files into fewer larger ones, and that it is a maintenance job you run rather than something the table does by itself.

for a middle

Explain the planning rules: candidates grouped per partition, a minimum candidate count before a group is rewritten, and a target file size the packer aims for. Name the options that change each.

for a senior

Diagnose a no-op run end to end — partition granularity versus data volume, thresholds, where scope, partial progress, concurrent-commit conflicts — and know that space is only released once snapshots expire.

for a principal

Frame it as layout policy: partition granularity chosen for real volume per partition, a compaction budget proportional to the data rewritten, and whether ingest batching should change so compaction stops being a treadmill.

## What the procedure actually plans ```sql CALL prod.system.rewrite_data_files( table => 'db.events', strategy => 'binpack', options => map( 'min-input-files', '2', 'target-file-size-bytes', '268435456', 'partial-progress.enabled', 'true' ) ); ``` The action scans the table's current snapshot, groups candidate data files **per partition**, and bin-packs each group toward the target size. Two thresholds decide whether a group is rewritten at all: `min-input-files`, default 5, and the size window around `target-file-size-bytes`, which defaults to the table property `write.target-file-size-bytes` (512 MB) with `min-file-size-bytes` and `max-file-size-bytes` derived from it. Files already inside that window are not rewritten, because rewriting them would burn IO for nothing. Setting `rewrite-all` to true overrides the thresholds and rewrites every selected file — useful when you are changing sort order or materialising deletes rather than fixing size. ## The partition boundary is the usual culprit Compaction never merges data files from different partitions, because a data file belongs to exactly one partition tuple. A table partitioned at hourly granularity that receives a few megabytes an hour has, by construction, several small files per hour and nothing to combine them with. The procedure will report almost nothing rewritten, and it is behaving correctly. The real fix is upstream: coarsen the partition transform (hours to days), evolve the partition spec, or change the writer's `write.distribution-mode` so each partition receives fewer, larger files in the first place. Compaction cannot repair a partition layout that is too fine for the data volume. ## Other reasons the run looks like a no-op - **Scope was narrowed.** A `where` argument restricts rewriting to matching partitions. Passing yesterday's predicate leaves every older partition untouched, which is often intentional but easy to forget. - **The target size never changed.** If the table sets nothing, groups are packed toward 512 MB. Complaining about "small" 200 MB files while the target is 512 MB and `min-input-files` is 5 means the files were simply not eligible. - **Partial progress.** With `partial-progress.enabled` the job commits in batches capped by `partial-progress.max-commits` (default 10) instead of one big commit at the end. If the job hit the cap or was killed, only some groups landed. Without partial progress, a failure at the end discards everything — on a huge table that is the difference between slow progress and none. - **Concurrency limits and group size.** `max-concurrent-file-group-rewrites` (default 5) and `max-file-group-size-bytes` bound how much runs at once; a long job may still be working through the backlog. - **New files arrived.** A streaming or micro-batch writer keeps appending small files; compaction fixes the past, and its effect is invisible if ingest granularity is the underlying problem. - **Commit conflict.** The rewrite commits by replacing the files it read. A concurrent write that deleted or replaced those same files can fail validation, so the run does work but does not commit. ## Storage did not drop either Even a fully successful rewrite frees no space immediately. The pre-compaction files are still referenced by earlier snapshots, so the table holds both sets until `expire_snapshots` removes those snapshots. Maintenance order is therefore rewrite, then expire. ## Beyond size: sort and deletes `strategy => 'sort'` with a `sort_order` argument — including `'zorder(col_a, col_b)'` — rewrites while clustering rows, which improves data skipping rather than just file count. On a merge-on-read table, `rewrite_data_files` also applies the delete files that cover the files it rewrites, so deleted rows are materialised away; the `delete-file-threshold` option makes files carrying many deletes eligible even when their size is fine. Where only the delete files themselves are the problem, `rewrite_position_delete_files` compacts those without touching the data. ## How to check rather than guess The procedure returns the number of rewritten data files and file groups. Compare file counts and sizes before and after with `SELECT partition, count(*), avg(file_size_in_bytes) FROM db.events.files GROUP BY partition`, and look at `db.events.snapshots` to confirm a `replace` snapshot was actually committed.

  • When would you use strategy 'sort' instead of the default binpack?
    When the problem is scan volume rather than file count. Binpack only packs files toward the target size and preserves nothing about row order, while `sort` rewrites rows clustered by a sort order — including `zorder(col_a, col_b)` for multiple dimensions — so min/max stats in the manifests become selective and readers skip whole files. It costs a full shuffle, so it is a periodic job, not a per-batch one.
  • Why did the rewrite job fail to commit even though it did all the work?
    The rewrite commits by replacing exactly the data files it read. If a concurrent writer removed or replaced any of them first, the commit's validation fails and the work is discarded. Enabling `partial-progress.enabled` limits the blast radius by committing groups incrementally, and scheduling compaction away from heavy write windows avoids the conflict rather than retrying into it.

saying these in an interview costs you the question

  • Assumes compaction merges files across partition boundaries
  • Never checks min-input-files or the target file size
  • Thinks storage shrinks immediately after a rewrite
  • Blames the procedure without reading its returned counts
  • Believes a where filter compacts the whole table anyway

context