skip to content

In ClickHouse, what happens on disk each time you INSERT into a MergeTree table?

level: juniorimportance: must knowfreq 78%

answer

  1. writes are append-only, never in place
  2. one directory per inserted block
  3. sorted, compressed, one file per column
  4. something in the background combines them
  5. parts are the unit queries must open

basics

~20 s

Each inserted block becomes a brand-new immutable data part: a directory of sorted, compressed per-column files plus index files. Existing files are never modified; a background merge process later combines parts into fewer, larger ones.

solid answer

~40 s

A MergeTree INSERT is append-only. The server buffers the incoming rows into blocks, sorts each block by the table's `ORDER BY` key, compresses each column into its own file, writes a complete new **part** directory, and then atomically publishes it into the set of active parts. Nothing already on disk is rewritten, so inserts never block concurrent readers. One statement can produce more than one part — you get roughly one part per block per partition touched, so an insert spanning many partition values creates many parts at once. Background merge threads then repeatedly merge small parts into bigger sorted parts. Because every query has to open and scan all active parts, the health of a ClickHouse table depends on inserts arriving in large batches rather than row-by-row.

code

sql · 15 lines
sql
CREATE TABLE events
(
    event_time DateTime,
    user_id    UInt64,
    action     LowCardinality(String)
)
ENGINE = MergeTree
PARTITION BY toYYYYMM(event_time)
ORDER BY (user_id, event_time);

INSERT INTO events VALUES (now(), 1, 'click');

SELECT partition, name, rows, level, bytes_on_disk
FROM system.parts
WHERE table = 'events' AND active;

go deeper

for a junior

Be able to say that ClickHouse writes a new immutable, sorted, compressed directory per inserted block and merges them later. Knowing that inserts should be batched is the expected takeaway.

for a middle

Explain the mechanics: block sizing, one part per block per partition, the sparse index inside a part, and how background merges bound the part count over time.

for a senior

Show you can reason about the read side too — per-part scan overhead, compression loss from tiny parts — and that you size batches and partition keys to control parts per second in production.

for a principal

Own the write-path contract across teams: who batches, what the per-table insert rate budget is, and how partition-key choices in one team's schema turn into merge-pool pressure for everyone on the cluster.

## The unit of storage: a data part In ClickHouse's MergeTree family, table data does not live in one big file. It lives in **parts**. A part is a directory on disk containing the rows of one write, already sorted by the table's `ORDER BY` (sorting) key, with each column stored in its own compressed file, alongside a sparse primary-index file and mark files that map index granules to offsets inside the column files. Every part belongs to exactly one partition — the value of the `PARTITION BY` expression, if the table declares one. Part directory names encode that partition plus a block-number range and a merge level, which is why you see names like `all_1_1_0` or `202408_15_20_1`. Parts are **immutable**. Once written, a part's files are never edited in place. Everything else in the engine follows from that choice. ## What an INSERT actually does When an INSERT arrives, the server: 1. accumulates the incoming rows into one or more in-memory blocks; 2. sorts each block by the table's sorting key; 3. applies the per-column compression codecs; 4. writes a complete new part directory and flushes it to disk; 5. atomically adds that part to the in-memory list of active parts, at which point it becomes visible to queries. No existing file is touched, so readers see either the whole part or none of it, and a crashed or failed insert leaves at most an unreferenced temporary directory that is cleaned up later. The count of parts produced is **not** one per statement. It is roughly one per block per partition. Two things drive it: block sizing settings such as `max_insert_block_size` and `min_insert_block_size_rows` / `min_insert_block_size_bytes`, and how many distinct partition values the inserted rows span. A single INSERT of a million rows spread over 300 daily partitions writes hundreds of parts, not one — this is why ClickHouse guards inserts with `max_partitions_per_insert_block` and why a badly chosen partition key turns ordinary loads into part storms. ```sql INSERT INTO events VALUES (...); SELECT partition, name, rows, level FROM system.parts WHERE table = 'events' AND active; ``` ## Background merges Because each write leaves a new part, a background pool continuously merges parts within the same partition into larger sorted parts, re-running the merge sort over their sorting keys and re-compressing the result. The old parts stay on disk, inactive, until they are dropped. Merging is what keeps the part count bounded, keeps compression ratios good (long runs of similar values compress far better inside one big part than split across many small ones), and, for the specialised engines, is where row collapsing or aggregation actually happens. Merging is background work with finite throughput. Inserts are foreground work with none. If inserts create parts faster than merges retire them, the part count grows without bound — the classic ClickHouse production failure. ## Why the part count matters for reads A query must consider **every active part** of every partition it did not prune away. For each part the engine opens files, reads the sparse primary index, decides which granules to read, and spins up read tasks. That per-part overhead is small but not free, and it is paid regardless of how few rows the part holds. A table with 40 well-merged parts and a table with 40,000 tiny parts can hold identical data and differ by orders of magnitude in scan time and memory. ## Practical consequences - **Batch your writes.** Aim for tens of thousands to hundreds of thousands of rows per INSERT rather than one row per statement, and keep the per-table insert rate low (on the order of one insert per second, not thousands). If the producer cannot batch, let the server do it with `async_insert`. - **Do not partition finely.** Partitioning by hour or by a high-cardinality id multiplies parts per insert. Monthly or daily partitioning is the usual choice; the sorting key, not the partition key, is what makes queries fast. - **Expect eventual, not immediate, tidiness.** Right after a big load the table will hold many parts; they shrink in number over the following minutes as merges run. `OPTIMIZE TABLE ... FINAL` forces the issue but rewrites data and is expensive, so it is a maintenance tool, not part of the load path. - **Watch it directly.** `system.parts` shows active parts and their sizes; `system.merges` shows merges in flight; `system.part_log` records each part's creation and merge history. The mental model to carry into an interview: an INSERT in ClickHouse is a file write, not a row update, and everything about ingestion tuning is about controlling how many files you create per second.

  • Does one INSERT statement always create exactly one part?
    No. It creates roughly one part per block per partition touched. A large insert can be split into several blocks by `max_insert_block_size`, and rows spanning many `PARTITION BY` values produce a separate part per partition. Inserting a day's data across 300 daily partitions writes hundreds of parts from a single statement, which is why `max_partitions_per_insert_block` exists as a guard.
  • If parts are immutable, how does ClickHouse handle a DELETE or UPDATE?
    Through mutations, which rewrite whole parts in the background rather than editing rows, or through lightweight deletes that mark rows with a hidden mask and let a later merge physically drop them. Either way the work is proportional to the parts touched, not the rows changed, so frequent small updates are a poor fit for this engine.
  • Where do you look to see how many parts a table currently has?
    `system.parts` filtered on the table with `active = 1` gives one row per live part, with row counts and sizes; aggregating by partition shows where parts are piling up. `system.merges` shows merges currently running, and `system.part_log` gives the historical trail of part creation, merging and removal.

saying these in an interview costs you the question

  • Thinks an INSERT appends rows into existing column files
  • Believes one INSERT statement always produces exactly one part
  • Assumes parts are merged synchronously before the insert returns
  • Says row-by-row inserts are fine because ClickHouse is fast
  • Confuses partitions with parts and uses the terms interchangeably

context