skip to content

Immutable Files, Compaction & Small Files

Analytical storage is written once as immutable files or segments and reshaped only by background merges, never updated in place. Interviewers probe it because the small-file problem and compaction lag are the top operational pains of every warehouse and lake.

on this pageshow

questions

6

Why is inserting rows one at a time into a columnar analytical table so much worse than batching?

level: juniorimportance: must knowfreq 70%

answer

  1. writes produce whole files, never rows
  2. fixed cost per file, not per row
  3. encodings need many rows to pay off
  4. merges must rewrite what you just wrote
  5. batch by size or time before committing

basics

~20 s

Columnar engines write immutable files, and the smallest unit a write can produce is a whole file. One row per insert means millions of tiny files, terrible compression, metadata bigger than data, and a background merge queue that can never catch up.

solid answer

~50 s

Analytical storage is append-only: a write buffers rows, transposes them into per-column blocks, encodes and compresses them, and publishes one new **immutable file**. There is no free space in an existing file to append into. So a single-row insert pays the full per-file cost — footer and statistics metadata, a catalog entry, one object-store write — for one row, and the encodings that make columnar storage cheap (dictionary, run-length, delta) have nothing to amortise over. On the read side the scan becomes metadata-bound: thousands of tiny reads instead of a few large sequential ones. The engine then tries to fix it with background merges, which rewrite the same data repeatedly; if inserts arrive faster than merges complete, file count grows without bound and the engine eventually throttles or rejects writes. The fix is to buffer client-side and commit batches sized in tens or hundreds of megabytes, or seconds-to-minutes of data.

code

sql · 7 lines
sql
-- anti-pattern: one immutable file per statement
INSERT INTO events VALUES ('2026-08-21 10:00:00', 42, 'click');
INSERT INTO events VALUES ('2026-08-21 10:00:01', 43, 'view');

-- one file, millions of rows, encodings amortised
INSERT INTO events
SELECT ts, user_id, action FROM staging_events;

go deeper

for a junior

Remember the one-line rule: analytical tables want few large writes, not many small ones, because each write creates a whole new file. Be ready to say what you would do instead — buffer and commit in batches.

for a middle

Explain the mechanics: rows are buffered, transposed into column blocks, encoded, and published as one immutable file, so per-file overhead and lost compression are paid per insert. Mention that background merges then have to rewrite everything.

for a senior

Show you have operated this: describe the file-count backlog, the throttling or write rejection that follows, and how you sized batches against a freshness requirement in a real pipeline.

for a principal

Own the trade: the batch interval is a contract between data freshness and storage/compute cost. Be ready to argue where in the pipeline buffering belongs and what latency the business is actually paying for.

## What an insert physically does An OLTP row store has mutable pages with free space, so an insert usually means "find a page with room, write the row into it". Columnar analytical storage does not work that way. Data is written once into a **self-contained immutable file** — Snowflake calls these micro-partitions, ClickHouse calls them parts, lake tables call them data files — and that file is never modified afterwards. Producing one involves buffering rows in memory, transposing them from rows into per-column blocks, encoding and compressing each block, writing a footer holding per-column offsets and statistics, and publishing the file atomically so readers either see all of it or none of it. The consequence that matters for this question: **the smallest thing a write can produce is a file**, not a row. A one-row insert produces a one-row file. ## Fixed costs paid per file, not per row Every file carries overhead that does not shrink with row count: - footer/metadata bytes describing every column, its encoding, its offsets and its min/max statistics; - an entry in the table's metadata or catalog that the planner must read; - one write operation against the storage layer (an object-store PUT is a network round trip with its own latency floor); - at read time, one open/GET plus a metadata parse plus a scheduled task. With a hundred rows in a file the metadata can be larger than the data itself. At a modest hundred single-row inserts per second you create over eight million files a day, and the table's metadata alone becomes a scalability problem before the data does. ## Compression and encodings need volume Columnar compression is cheap because a column block holds many similar values: a dictionary is worth building when thousands of rows share a few hundred distinct strings, run-length encoding is worth it when values repeat, delta encoding is worth it over a long increasing sequence. A block containing one value compresses to roughly "the value plus a header". Batches also matter for execution: vectorized readers process values in batches of thousands, so one-row column chunks defeat the whole read path. ## The background-merge tax Because engines expect this failure mode, they run background merges that combine small files into larger ones. That helps readers but is not free: merging rewrites data that was already written, which is **write amplification** — the same row is read and written again on every merge level it passes through. Merge throughput is finite. If the insert rate produces files faster than merges can consolidate them, the backlog grows, query planning degrades, and the engine's defence is to slow down or reject writes rather than melt (ClickHouse's "too many parts" error is the well-known example of exactly this guardrail). ## What to do instead - **Batch on the client.** Accumulate rows and commit on a size or time trigger — for example flush every 10–100 MB or every few seconds, whichever comes first. This is a deliberate trade of freshness for file size. - **Use the engine's buffering ingest path** if it has one: managed streaming ingest, an asynchronous insert buffer, or a write-optimised staging area that the engine consolidates for you. - **Put a queue in front.** A stream processor or message broker consumer that micro-batches turns an unbounded stream of tiny writes into a bounded stream of reasonable files. - **Stage then merge.** Land raw arrivals in a small landing table and periodically `INSERT INTO target SELECT * FROM landing` in large chunks, then clear the landing table. ## How to talk about it in an interview A strong answer names immutability as the root cause rather than saying "inserts are just slow", then gives both sides: the write side (per-file fixed cost, no compression amortisation, merge backlog) and the read side (metadata-bound planning, many tiny I/Os, poor scan throughput). It also quantifies the fix — "batch to roughly file-size targets, seconds of latency, not per-row" — instead of vaguely recommending "bulk loading".

  • How would you choose the flush interval for a client that buffers rows before committing?
    Work backwards from two constraints: the freshness the consumers actually need, and the file size the engine wants. Flush on whichever trigger fires first — a size threshold near the target file size, or a time threshold set by the freshness requirement. If the incoming rate is low enough that the size trigger never fires, accept that the time trigger produces small files and rely on compaction to clean them up.
  • If batching is impossible because the source is a per-event webhook, what are your options?
    Put a buffer between the source and the table: a message queue or stream consumed by a micro-batching writer, or the engine's own asynchronous insert buffer if it offers one. Alternatively write events to a small write-optimised staging table and merge it into the main table on a schedule. The goal is the same — decouple event arrival rate from file creation rate.
  • Does a bigger batch always mean better performance?
    No. Past the engine's target file size the gains flatten, while the writer needs more memory to buffer, a failed commit wastes more work, and the data becomes staler. Very large files also give coarser granularity for parallel scans. Batch up to the target file size, then commit.

It is like mailing one sheet of paper per envelope: the envelope, the address and the postage cost the same whether it holds one page or five hundred, and the recipient has to open every one.

saying these in an interview costs you the question

  • Thinks the engine appends rows into an existing file
  • Blames insert latency only, missing the file-count explosion
  • Believes compression works equally well on tiny blocks
  • Assumes background merges make any insert pattern fine
  • Suggests adding an index instead of batching writes

context

open as a page

In a columnar table whose data files are immutable, how does an UPDATE or DELETE actually work?

level: middleimportance: must knowfreq 75%

basics

~20 s

The engine never edits a file. It finds the files containing affected rows, writes new files holding the surviving and modified rows, and atomically swaps the file set in the table metadata. Old files stay until retention expires, then are dropped.

open as a page

Why do columnar engines target data files of hundreds of megabytes rather than a few megabytes?

level: middleimportance: should knowfreq 55%

basics

~20 s

Large 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.

open as a page

A columnar table on object storage got tiny appends every minute for a year — why did queries slow down?

level: seniorimportance: should knowfreq 65%

basics

~20 s

The table accumulated hundreds of thousands of tiny files. Planning must enumerate them all, each scan task pays a storage round trip for a few kilobytes of data, and compression is poor — so the query becomes metadata and latency bound rather than throughput bound.

open as a page

How would you set ingest batching and compaction policy for a continuously loaded, dashboard-queried columnar table?

level: principalimportance: should knowfreq 40%

basics

~20 s

Start from the freshness the dashboards actually need, size ingest batches to that interval, and let compaction absorb whatever fragmentation remains. Then bound write amplification by choosing how many merge passes data goes through and how large the final files get.

open as a page

In an immutable-file columnar table, what is a delete vector and when does it beat rewriting the file?

level: seniorimportance: nice to knowfreq 35%

basics

~20 s

A delete vector is a small side file marking which row positions in a data file are logically deleted. Readers filter them out on the fly. It turns an expensive file rewrite into a tiny write, at the cost of extra work on every read until compaction.

open as a page