concept

Small File Problem

The severe query degradation caused by a table stored as very many tiny files, where per-file overhead dominates actual data reading.

performancestoragecompaction

The classic way to create it is streaming ingestion writing a file per micro-batch: a job committing every minute produces 1,440 files a day, 525,000 a year, per partition. The table is a normal size; the file count is catastrophic.

The cost is per-file rather than per-byte, which is why it surprises people. Each file requires listing, opening, reading a footer and planning, and on object storage each of those is a network round trip with latency measured in tens of milliseconds. A query touching 50,000 small files spends almost all its time on overhead. Metadata operations degrade first and worst, and eventually the planning phase alone times out.

The fix is compaction: periodically rewriting many small files into files in the region of 128 MB to 1 GB, which is where columnar formats and object storage both perform well. Open table formats support this as a maintenance operation, including compacting while queries continue to run.

The design-time avoidance is better than the cure: buffer streaming writes to a sensible file size before committing, accept the resulting latency, and choose partitioning that does not shatter each batch across hundreds of directories.