ClickHouse — A Columnar Analytics Database from the Inside
Parts and merges — INSERTs create parts, merges reduce them
In one line
In MergeTree, one INSERT creates at least one new part, and background merges combine small parts into larger ones. If parts pile up faster than they merge, the server rejects INSERTs (TOO_MANY_PARTS). So the heart of loading data into ClickHouse is not "how many rows you insert" but "how many parts you create".
Why this was needed
The previous module said that a part, once written, does not change. Thanks to that, a write is just creating one new directory without locks, and a read only has to scan files that are already sorted and compressed. But this design comes with a bill. If you insert rows one line at a time, one single-row part is created each time. A query has to look at each part's index separately and open its files separately, so the more parts there are, the slower it gets, and merging has to read those small parts again and write them out again.
People used to row-oriented databases easily write code that sends one INSERT line per event. In ClickHouse that code becomes a production outage. This module follows the process by which parts are created and merged through the events in system.part_log, and sees in numbers what gets blocked when you cross the threshold and what asynchronous INSERT, where the server gathers data for you, changes.
How it works
The part name is the history. all_3_7_2 means partition all, covering block numbers 3 through 7, at merge level 2. A part created by an INSERT receives one block number and becomes all_N_N_0, and as the official documentation (Part merges) says, the level rises by one with each merge. In the lab Pod, if you send 10 INSERT statements to a table with merges stopped, all_1_1_0 … all_10_10_0 remain as they are.
One INSERT ≠ one part. The server gathers the incoming rows into blocks and writes a part per block. INSERT ... SELECT appends blocks until min_insert_block_size_rows (default 1,048,449) have gathered. When I inserted 1 million rows with min_insert_block_size_rows = 250000, 4 parts were created (261,636 × 3 + 215,092). If the table has partitions, blocks are divided again per partition — the topic of the next module.
Events are left in part_log. This Pod has system.part_log turned on. When a part is created, one line of NewPart (together with which INSERT's query_id made it) is added, and when parts are merged, one line of MergeParts (what was merged is in merged_from). system.parts shows only "now", but part_log shows "how we got here". If you drop a table and recreate it with the same name, it mixes with the old records, so filter by table_uuid.
Merges can be stopped, but OPTIMIZE stops with them. SYSTEM STOP MERGES 표 (the Korean word stands for the table name) stops merges of that table (they are released when the server restarts). Measured, running OPTIMIZE ... FINAL on a stopped table was rejected with Cancelled merging parts (ABORTED). Right after releasing it with SYSTEM START MERGES, a background merge first combined the 10 parts into all_1_10_1, and FINAL rewrote that single part into all_1_10_2. As the documentation says, FINAL merges even when there is already one part. Which path it took may differ each time, so confirm with part_log.
Thresholds. If a partition's active parts exceed parts_to_delay_insert, INSERTs are deliberately delayed, and when they reach parts_to_throw_insert, they are rejected with Too many parts ... Merges are processing significantly slower than inserts (error 252). This value is a table setting. When I lowered it to 5 and inserted one row at a time with merges stopped, the first five went in and the sixth was rejected. The way to release it is to reduce parts — turn merges on, merge, and insert again and it goes in.
Asynchronous INSERT is gathered by the server. The 26.8 defaults are async_insert = 1 and wait_for_async_insert = 1. When you send INSERT ... VALUES in the lab Pod, two queries, Insert and AsyncInsertFlush, are left in query_log. But in waiting mode, when one client sends them one at a time, the buffer is emptied each time, so parts are created one at a time (5 sends → 5 parts). For several INSERTs to become one part, they must be in the buffer at the same time. When I sent 5 in non-waiting mode (wait_for_async_insert = 0) and emptied it with SYSTEM FLUSH ASYNC INSERT QUEUE, the 5 became one part. The documentation does not recommend non-waiting mode — because errors are not returned to the client. In the lab it is used to see gathering that does not depend on time. For reference, asynchronous INSERT does not apply to INSERT ... SELECT.
What it looks like in the field
The most common outage is "Too many parts". The cause is almost always one of two — a collector that inserts one line at a time, or partitions cut too finely. The documentation also says raising the threshold only delays the symptom. If you first count NewPart per query_id in part_log, you immediately see which INSERT is pouring out parts. The fix is either to gather on the client and send (the recommendation in the documentation (Selecting an insert strategy) is at least 1,000 rows at a time, preferably 10,000–100,000), or, if that is not possible, to leave it to asynchronous INSERT.
The second is an operation that runs OPTIMIZE ... FINAL from cron. The documentation (Avoid OPTIMIZE FINAL) says to treat it as an administrative task, not a routine one — it rewrites even when there is already one part, so it is expensive on a big table. The reason this lab uses FINAL is to fix the number of parts so you can see it.
What you will do in the next lab
You send 10 INSERT statements to parts.events, with merges stopped, to create 10 parts, and confirm with part_log that one 1-million-row INSERT becomes several parts depending on block size. You read the process of merging with OPTIMIZE after turning merges on, trigger TOO_MANY_PARTS on a table with the threshold lowered to 5, and then release it. Finally you gather 5 asynchronous INSERTs into one part and sum up the number of new parts per table in one table.