skip to content

Data Ingestion

Inserts should be large and infrequent, or made async, because every insert creates a part and too many parts choke the merge pipeline. Interviewers ask about batching and the Kafka table engine since small-insert part explosion is the classic ClickHouse production failure.

part ofClickHouseoverview, primer and where to startread it →
on this pageshow

questions

6

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

open as a page

Why does a ClickHouse MergeTree table start rejecting inserts with "Too many parts"?

level: middleimportance: must knowfreq 80%

basics

~20 s

Because inserts are creating parts faster than background merges can combine them. ClickHouse first delays and then rejects inserts once a partition's active part count crosses the parts_to_delay_insert and parts_to_throw_insert thresholds, protecting query performance and the merge pool.

open as a page

What does ClickHouse's async_insert setting change about the INSERT write path?

level: middleimportance: should knowfreq 58%

basics

~20 s

With async_insert enabled, the server collects rows from many small concurrent inserts into an in-memory buffer per table and shape, then flushes one part when the buffer hits a size or time limit. It turns server-side batching on for clients that cannot batch themselves.

open as a page

How do you bulk-load Parquet files from S3 into ClickHouse with the s3 table function?

level: middleimportance: should knowfreq 45%

basics

~20 s

Run INSERT INTO target SELECT ... FROM s3(url, format), where the URL may contain glob patterns matching many files. ClickHouse reads and parses the files in parallel; s3Cluster spreads that work across every node of a cluster instead of one.

open as a page

How do you ingest a Kafka topic into ClickHouse using the Kafka table engine?

level: seniorimportance: should knowfreq 55%

basics

~20 s

Create three objects: a Kafka engine table that consumes the topic, a MergeTree table that stores the data, and a materialized view that reads from the first and inserts into the second. The Kafka table is a consumer, not storage — selecting from it consumes messages.

open as a page

When loading into ClickHouse, how do Native, Parquet and JSONEachRow formats differ in cost?

level: middleimportance: nice to knowfreq 32%

basics

~20 s

Native is ClickHouse's own columnar block format and needs almost no parsing, so it is the cheapest to ingest. Parquet is columnar and compact but must be decoded and converted. JSONEachRow is the most expensive: text parsing and type conversion per field.

open as a page