skip to content

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

level: middleimportance: must knowfreq 80%

answer

  1. writes and merges race each other
  2. the engine slows you down before it stops you
  3. the limit is counted per partition, not per table
  4. average rows per part tells the whole story
  5. raising the ceiling only hides the decay

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.

solid answer

~40 s

Every insert block writes a new part, and background merges retire them at a finite rate. When the producer inserts small batches many times a second — or one insert sprays rows across many partitions — parts accumulate. MergeTree defends itself in two stages: past `parts_to_delay_insert` it artificially sleeps incoming inserts to apply backpressure, and past `parts_to_throw_insert` it fails them outright with `TOO_MANY_PARTS`. The fix is almost never to raise the thresholds. It is to insert **fewer, bigger** batches (tens of thousands of rows or more, roughly one insert per second per table), enable `async_insert` so the server batches for you, and coarsen an over-granular `PARTITION BY` so a single insert stops touching hundreds of partitions. Check `system.parts` for where parts pile up and `system.merges` for whether the merge pool is saturated.

code

sql · 12 lines
sql
-- which partitions hold the most parts, and how thin are they?
SELECT
    table,
    partition,
    count()                       AS parts,
    sum(rows)                     AS rows,
    round(sum(rows) / count())    AS rows_per_part
FROM system.parts
WHERE active
GROUP BY table, partition
ORDER BY parts DESC
LIMIT 10;

go deeper

for a junior

Recognise the error and know its one-line cause: too many small inserts creating parts faster than background merges combine them. The remedy to state is bigger, less frequent batches.

for a middle

Explain the two-stage defence — delay then reject — that the thresholds are counted per partition, and how block sizing and an over-granular partition key multiply parts per insert.

for a senior

Demonstrate the diagnosis: average rows per part from system.parts, merge activity from system.merges, and the judgment that raising the threshold trades a loud failure for permanent read degradation.

for a principal

Own the ingestion contract: where batching happens architecturally (producer, async_insert, or a Kafka buffer tier), how partition-key review is enforced at schema-design time, and how backfills are isolated from live streams on a shared merge pool.

## The failure in one sentence `Too many parts` means the write side is winning a race against the merge side: inserts are minting new data parts faster than the background merge pool can fold them into larger ones, and the engine has decided to push back rather than let the table degrade indefinitely. ## How the engine pushes back MergeTree watches the number of **active parts in a single partition** and applies two thresholds: - `parts_to_delay_insert` — above this, the server deliberately sleeps each incoming insert before acknowledging it. This is backpressure: the producer slows down, merges catch up. Symptomatically, insert latency climbs from milliseconds to seconds with no change in payload size. - `parts_to_throw_insert` — above this, the server stops accepting inserts for that table and raises the `TOO_MANY_PARTS` error with a message about merges processing significantly slower than inserts. Both are MergeTree-level settings and both are per partition, not per table — which is the detail candidates most often miss. A table with a thousand healthy partitions can still throw if one hot partition has accumulated too many parts. ## Why the parts pile up Four causes cover nearly every real incident: 1. **Small, frequent inserts.** The application inserts one row (or one HTTP request's worth) per statement, thousands of times a second. Each is a part. Merges cannot keep up at that rate no matter how much CPU you give them, because merge cost is superlinear in the number of tiny inputs relative to the useful data volume. 2. **An over-granular partition key.** `PARTITION BY toStartOfHour(ts)` or `PARTITION BY user_id` multiplies parts per insert: one insert covering many partition values writes one part per value. The guard `max_partitions_per_insert_block` exists precisely to catch this, and hitting it is a signal the partition key is wrong, not that the guard is wrong. 3. **Merges are throttled or blocked.** The background merge pool is sized finitely; if it is undersized, saturated by a huge merge, starved of disk I/O, or stopped (`SYSTEM STOP MERGES` left on after maintenance), parts accumulate even with sane inserts. 4. **Concurrent bulk backfill on top of live traffic.** A historical load writes into old partitions while the stream writes into today's, and the merge pool is shared. ## Diagnosing it ```sql -- where are the parts? SELECT table, partition, count() AS parts, sum(rows) AS rows FROM system.parts WHERE active GROUP BY table, partition ORDER BY parts DESC LIMIT 10; -- is the merge pool busy or idle? SELECT table, elapsed, progress, num_parts, total_size_bytes_compressed FROM system.merges; ``` If part counts are high and `system.merges` is idle, merges are blocked or throttled — investigate the pool and disk. If merges are constantly running and still losing, the insert pattern is the problem. `system.part_log` gives the historical rate of part creation, which is the number to compare against your batching assumptions. Average rows per part (`sum(rows) / count()` from `system.parts`) is the single most diagnostic figure: a few hundred rows per part is a smoking gun. ## Fixing it, in order of preference - **Batch at the producer.** Accumulate tens of thousands to hundreds of thousands of rows and insert once per second per table rather than continuously. This is the real fix and everything else is a workaround. - **Use `async_insert`.** Set `async_insert = 1` so the server buffers rows from many small concurrent inserts and flushes one part per buffer. Ideal when the producers are many independent processes that cannot coordinate a batch between them. - **Put a buffer in front.** Kafka plus the Kafka table engine, or an aggregation tier, converts a spray of small writes into steady large blocks. - **Coarsen `PARTITION BY`.** Move from hourly to daily, or daily to monthly. Partitioning exists for data lifecycle management (dropping old partitions cheaply), not for query pruning — the sorting key does the pruning. - **Give merges room.** Check background pool sizing and disk throughput; make sure merges were not left stopped. - **Raising the thresholds is the last resort.** It converts a loud, early failure into a silent slow decay: more parts means more per-part scan overhead, worse compression, and more memory per query. Raise them only as a temporary shield while you fix the write path. ## The interview framing This is the canonical ClickHouse production failure, and the answer interviewers are listening for has two halves: the mechanism (append-only parts versus finite background merge throughput, with two thresholds implementing delay-then-reject) and the judgment (fix the batching, not the threshold).

  • Insert latency has risen from 20 ms to 3 seconds but nothing is erroring yet. What is happening?
    You are in the delay band: the active part count in some partition has crossed `parts_to_delay_insert`, so the server is deliberately sleeping each insert to apply backpressure before it reaches the hard `parts_to_throw_insert` limit. Treat it as the early warning it is — check average rows per part and fix batching now, because the reject threshold is next.
  • Would raising parts_to_throw_insert be a reasonable fix?
    Only as a temporary shield. It removes the error without removing the cause, and a table carrying tens of thousands of tiny parts pays for it on every read: more file handles and index reads per scan, worse compression ratios, higher query memory. You have traded a loud failure for slow, permanent degradation.
  • Why can a single INSERT statement trigger this even when the batch is large?
    Because parts are per partition. An insert of a million rows spanning 400 hourly partitions writes at least 400 parts, not one. That is why `max_partitions_per_insert_block` guards the case, and why an over-granular `PARTITION BY` produces part explosion from batches that look perfectly healthy in row count.
  • How do you tell whether merges are the bottleneck rather than the inserts?
    Look at `system.merges`. If it is busy continuously and part counts still climb, merge throughput is genuinely saturated — check disk I/O and background pool sizing. If it is idle while parts accumulate, merges are blocked or stopped, so look for a leftover `SYSTEM STOP MERGES` or a resource constraint rather than blaming the write rate.

saying these in an interview costs you the question

  • Raises parts_to_throw_insert and calls the incident resolved
  • Thinks the part limit is per table rather than per partition
  • Runs OPTIMIZE TABLE FINAL on a schedule as the standard fix
  • Blames disk space rather than merge throughput versus insert rate
  • Adds finer partitioning to 'spread the load' and worsens it

context