Parts & Background Merges
How writes become storage
Every INSERT into a MergeTree table creates a new part — never
modifies an old one. Understanding parts and how they merge explains
almost every ClickHouse operational quirk: why small, frequent inserts
are bad, why SELECT COUNT(*) is instant, and why OPTIMIZE TABLE
exists.
What a part actually is
A part is a directory on disk. Inside it: one compressed file per column, a sparse primary index, checksums, and metadata like row count and column min/max values. A part is immutable — once written, it is only ever read or deleted, never edited in place.
This is where compressed blocks and
granules meet. Each column's .bin file is
just its compressed blocks written back to back. The paired .mrk2
("marks") file is the missing link between them: one mark per granule,
pointing at the exact byte offset in .bin where that granule's block
starts. primary.idx is the sparse index itself — one entry per
granule, holding the primary key of its first row. A lookup
binary-searches primary.idx for the right granule, reads the
matching mark, and seeks straight to that byte offset — no scanning,
no decompressing anything it doesn't need.
The 2 in .mrk2 is a format version — ClickHouse has used .mrk,
.mrk2, and .mrk3 as the on-disk marks format evolved (e.g. to
support adaptive granule sizes). The role never changes: granule →
byte offset.
Why merges are necessary
If every insert became a permanent, separate part, a table that receives thousands of small inserts a day would end up scanning thousands of tiny files per query — index overhead per part, open file handles per part, no sorting across inserts. Background merges continuously combine smaller sorted parts into fewer, larger sorted parts, restoring the property that the whole table behaves like one big sorted structure.
each INSERT writes its own immutable part — nothing sorted across them yet
This is the same idea as compaction in an LSM-tree (RocksDB, Cassandra, LevelDB): accept writes fast by appending, then pay the reorganization cost later, in the background, off the write path.
Merges are a suggestion, not a promise
ClickHouse decides when and which parts to merge based on part size and count — you don't control the schedule. You can force it:
OPTIMIZE TABLE events FINAL;OPTIMIZE ... FINAL forces a full merge of every part into one,
rewriting the entire table's data on disk. It is expensive and
I/O-heavy — reasonable for a one-off cleanup on a small/medium table,
dangerous to run routinely on a large one.
The practical consequence: batch your inserts
Because every INSERT statement creates at least one new part
regardless of size, sending one row per INSERT is one of the most
common ways to misuse ClickHouse — it creates a huge number of tiny
parts faster than the background merge process can keep up, and the
server starts rejecting inserts with Too many parts.
- Batch inserts client-side: thousands of rows per
INSERT, not one. - Or insert through a
Buffertable / async insert queue that batches for you. - Fewer, larger inserts → fewer parts → less merge pressure → faster queries.
SELECT count() FROM table without a WHERE is instant because
each part already stores its row count in metadata — ClickHouse just
sums those numbers without reading a single row of actual data.