Why do columnar engines target data files of hundreds of megabytes rather than a few megabytes?
answer
- fixed cost per file, variable cost per byte
- storage round trips have a latency floor
- encodings need volume to pay off
- too few files starves parallel workers
- statistics get coarser as files grow
basics
~20 sLarge files amortise per-file overhead — metadata, catalog entries, storage round trips, task scheduling — and give encodings enough rows to compress well. Files that are too large hurt too: less parallelism, coarser skipping, and every mutation rewrites more bytes.
solid answer
~50 sFile size is a two-sided optimisation. **Too small** and fixed per-file costs dominate: footer metadata may rival the data, the catalog grows huge, each read is a separate storage round trip whose latency dwarfs the transfer, and encodings have too few values to amortise a dictionary or a run. **Too large** and you lose granularity: fewer files means fewer independent units of parallel work, so a wide cluster idles; skipping becomes coarser because statistics summarise more rows per file; any mutation rewrites more bytes than it needed to; and writers need more memory to buffer a file before publishing it. The sweet spot most engines converge on is tens to a few hundred megabytes compressed — large enough that a sequential read amortises the open cost, small enough that a query still has hundreds or thousands of independent chunks to spread across workers. Sizing is enforced at write time by batching, and repaired afterwards by compaction.
code
text · 2 linespartition=2026-08-20 files= 12 avg_size= 210 MB plan_ms= 40 scan_ms= 2100
partition=2026-08-21 files=8420 avg_size= 0.3 MB plan_ms=3900 scan_ms= 5400go deeper
Know that analytical engines want a moderate number of reasonably large files, not thousands of tiny ones, and that batching your loads is how you get there.
Be able to argue both directions: what small files cost (metadata, round trips, weak compression) and what oversized files cost (lost parallelism, coarser skipping, expensive rewrites).
Diagnose from evidence — average file size per partition, file counts, planning time as a share of query time — and separate a writer-batching problem from a partitioning-granularity problem.
Treat file sizing as a platform standard rather than a per-table fix: define targets, decide who enforces them at write time versus via compaction, and account for the compute those merges consume.
## The two failure modes Data-file sizing in a columnar engine is not a single-direction rule of thumb. Both extremes have concrete, explainable costs, and a good interview answer names both. ### Too small Every file carries fixed overhead that is independent of how many rows it holds: - **Metadata.** The footer describes every column: encoding, offsets, sizes, per-column min/max statistics, sometimes null counts and dictionaries. In a file with a few hundred rows, this can exceed the data payload. - **Catalog pressure.** The table's file list is itself data the planner must read. Millions of entries make query planning slow before scanning even starts. - **I/O round trips.** On object storage a read has a latency floor per request. Fetching 2 MB costs nearly as much wall-clock as fetching 100 MB when latency dominates transfer; ten thousand small requests are hugely more expensive than a hundred large ones for the same bytes. Where the store charges per request, this also shows on the bill. - **Lost compression.** Dictionary, run-length and delta encodings amortise their own overhead across many values. Small blocks compress badly, so small files store the same logical data in more bytes. - **Scheduling overhead.** Engines assign work per file or per chunk; a tiny file still costs a task to schedule, a thread to run and a result to collect. ### Too large - **Parallelism collapses.** The unit of parallel scan work is bounded by file (or sub-file chunk) granularity. If a 1 TB table is eight files and the cluster has two hundred workers, most workers have nothing to do. Rule of thumb: you want comfortably more independent chunks than you have parallel workers, so stragglers can be balanced. - **Coarser skipping.** Per-file statistics summarise the whole file. A wider range in a file's min/max makes it less likely the file can be eliminated, so scans read more bytes than a smaller-file layout would. - **Expensive mutations.** Because replacement happens at file granularity, a delete or update of a handful of rows rewrites an entire file. Bigger files mean a bigger rewrite for the same logical change. - **Writer memory and failure cost.** A writer buffers and encodes a file before publishing it; larger targets need more memory, and a failed write throws away more work. ## Where the balance lands Engines and table formats converge on a target measured in tens to a few hundred megabytes **compressed**, with internal sub-divisions (row groups, blocks, granules) that give finer-grained parallelism and statistics inside a single file. That layering is the real answer to the tension: the file is sized for the storage layer and the metadata layer, while the sub-file chunk is sized for the scan and skipping layer. Compressed size matters more than row count for the storage side; row count matters more for the encoding and vectorized-read side. In practice you tune batches by both — flush at N megabytes or M rows, whichever trips first. ## Achieving the target Two mechanisms, and both are usually needed: 1. **At write time**, batch. Buffer rows until the batch would produce a file near the target, then commit. This is the cheap path — bytes are written once. 2. **After the fact**, compact. Background or scheduled merges read many small files and write fewer large ones. This is the expensive path — the same bytes are written at least twice — but it is the only remedy for a stream whose natural batch is small, and for tables whose partitioning produces uneven file sizes. A related trap: over-partitioning. Splitting a table into very fine partitions caps the amount of data any single file can hold. A table partitioned by hour and by tenant across thousands of tenants may physically be unable to produce large files, because each partition receives only a trickle. When a small-file problem resists compaction, the partitioning scheme is often the real cause. ## Signals in production Average bytes per file, file count per partition, and the ratio of planning time to execution time are the metrics that expose this. When planning time grows as a share of total query time and the average file is a few megabytes, you have a sizing problem — not a compute problem, and adding workers will not fix it.
- If a file is a few hundred megabytes, how does a query still get fine-grained parallelism?Files are internally divided into chunks — row groups, blocks or granules — each with its own statistics and offsets. Workers can be assigned individual chunks rather than whole files, so parallelism and skipping operate below file granularity while the file itself stays sized for the storage and metadata layers.
- A table has a persistent small-file problem that compaction keeps re-creating. What would you look at?The partitioning scheme first. If partitions are fine enough that each receives only a trickle of data, no batching strategy can produce large files inside them and compaction has nothing to merge across partition boundaries. Coarsening the partition granularity usually fixes it more effectively than tuning the writer.
- Why is compressed size the more useful target than row count?Storage round trips, catalog entries and scan throughput all scale with bytes, and bytes per row varies enormously between a narrow numeric table and a wide table of text. A row-count target that yields good files on one table produces tiny or enormous ones on another. Practical writers trip on either threshold, whichever comes first.
Shipping freight: one pallet per parcel wastes the truck, and one truck-sized crate cannot be split across two lorries. You want boxes big enough to be worth the paperwork and small enough that many hands can carry them at once.
saying these in an interview costs you the question
- Believes bigger files are always better
- Ignores that huge files starve parallel workers
- Measures targets in rows without regard to bytes
- Thinks compaction can fix an over-partitioned table
- Assumes more compute nodes fix a small-file problem