A ClickHouse table partitioned by toYYYYMMDD(event_time) now rejects inserts with "Too many parts" — why?
answer
- one part per partition touched, per insert
- merges never cross the boundary
- count the buckets, not just the rows
- partition by what you delete
- raising the limit hides the cause
basics
~20 sEach insert creates at least one part per partition it touches, and merges never combine parts across partitions. A daily partition key multiplies the active partitions, so parts accumulate faster than background merges can consolidate them and ClickHouse throttles then rejects inserts.
solid answer
~40 sEvery `INSERT` block writes one new part per distinct partition value it contains, and background merges only combine parts **within** a partition. With `toYYYYMMDD(event_time)`, a batch spanning a week creates seven parts instead of one, and each day's partition consolidates independently and slowly because it holds a fraction of the data. Part counts climb, and once a partition passes the `parts_to_delay_insert` threshold ClickHouse deliberately slows inserts, then throws `TOO_MANY_PARTS` past `parts_to_throw_insert`. The fix is usually the partition key, not the thresholds: move to `toYYYYMM(event_time)` unless daily `DROP PARTITION` is a genuine retention requirement, and check whether backfills are spraying old timestamps across many partitions at once. Raising the thresholds buys time and makes queries worse, since index analysis runs per part.
code
sql · 9 linesSELECT partition,
count() AS parts,
sum(rows) AS rows,
formatReadableSize(sum(bytes_on_disk)) AS size
FROM system.parts
WHERE active AND database = currentDatabase() AND table = 'events'
GROUP BY partition
ORDER BY parts DESC
LIMIT 20;go deeper
Know that each insert creates parts, that merges tidy them up in the background, and that too many small parts is a real failure mode with a clear error message.
Explain that parts never merge across partitions, so active-partition count multiplies fragmentation, and name the delay-then-throw thresholds that produce the error.
Diagnose from system.parts whether the cause is key cardinality or write batching, and fix the partition key or the loader rather than raising the threshold.
Set the rule that partition granularity follows the retention and archival unit, and own the migration cost of changing a partition key on a table already at scale.
## The mechanics behind the error A `MergeTree` insert is atomic per block and writes **one immutable part per distinct partition value in that block**. Background merges then consolidate small parts into larger ones — but only inside a single partition, never across. The number of active partitions therefore sets a floor on how fragmented the table can get, and it multiplies the effect of every other fragmentation source. ClickHouse protects itself with two per-partition thresholds. Once a partition's active part count crosses `parts_to_delay_insert`, inserts are artificially slowed to give merges time to catch up. Past `parts_to_throw_insert`, inserts fail outright with the `TOO_MANY_PARTS` error. A separate guard, `max_partitions_per_insert_block`, rejects a single insert block that touches too many partitions at once — the classic symptom of a partition key with far too high cardinality. ## Why a daily key makes it likely Compare monthly and daily partitioning of the same stream: - **Monthly.** Ongoing inserts land in one partition. Merges see a steady stream of parts in that one bucket and consolidate them into progressively larger parts. Steady-state active parts stay small. - **Daily.** Ongoing inserts land in today's partition — fine so far. But any batch whose rows span midnight, and every backfill or late-arriving-event load, writes into several partitions at once. Each of those partitions merges independently and slowly, because merge scheduling favours combining similarly-sized parts and each day holds only a thirtieth of the volume. The table's total part count is roughly the per-partition count times the number of active partitions. The pathological version is a partition key with unbounded cardinality — `PARTITION BY event_time` on a raw `DateTime`, `PARTITION BY user_id`, `PARTITION BY cityHash64(...)`. There, every insert block touches hundreds or thousands of partitions and the table falls over almost immediately. ## Diagnosis Start with `system.parts`, grouping active parts by partition. The shape of the answer tells you which story you are in: a single huge partition with many parts means an ingest-rate problem; many partitions each with a handful of parts means the partition key is too fine. ```sql SELECT partition, count() AS parts, sum(rows) AS rows, formatReadableSize(sum(bytes_on_disk)) AS size FROM system.parts WHERE active AND database = currentDatabase() AND table = 'events' GROUP BY partition ORDER BY parts DESC LIMIT 20; ``` If many partitions each hold few rows, the partition key's cardinality is the cause. If one recent partition holds hundreds of tiny parts, look at how the writer batches — that is an ingest-side problem, and the answer lies in fewer, larger inserts. ## Fixes, in order of preference 1. **Coarsen the partition key.** Move to `toYYYYMM(event_time)`, or to weekly if retention genuinely operates at that granularity. The test is simple: partition by the unit you *drop or move*, nothing finer. If nobody ever runs `DROP PARTITION` on a single day, daily partitions are buying nothing and costing a lot. Changing the partition key means creating a new table and copying the data — plan it as a migration. 2. **Stop backfills from spraying.** A load that spans many partitions at once is worth restructuring so each batch targets one partition, which both reduces part creation and keeps individual inserts atomic per partition. 3. **Batch harder on the write path.** Fewer, larger inserts create fewer parts regardless of the partition key. This is the ingest side of the same coin and is often the co-cause. 4. **Raise the thresholds — last, and knowingly.** Increasing `parts_to_delay_insert` and `parts_to_throw_insert` removes the error without removing the fragmentation. Since primary-index analysis runs per part, a table with thousands of small parts reads far more granules than the same data in a few large ones, so you trade a loud failure for slow queries. Do it to survive an incident, not as the fix. `OPTIMIZE TABLE ... PARTITION ...` can force consolidation of one partition and is useful as a one-off after a bad backfill. It is not a schedule-it-nightly answer: it does full-rewrite I/O over the partition and competes with the merges you actually need. ## The judgment the interviewer is after The candidate who says "raise parts_to_throw_insert" has treated the alarm. The candidate who says "the partition key should match the retention unit, and daily partitions on a table nobody drops daily are pure cost" has found the cause. Mentioning that the total part count drives query planning as well as insert health shows the connection between the ingest symptom and the query-side consequence — which is the reason this question keeps getting asked.
- How would you tell an over-fine partition key apart from an ingest-batching problem?Group active parts by partition in `system.parts`. Many partitions each holding a few small parts points at partition-key cardinality; one recent partition holding hundreds of tiny parts points at a writer inserting too frequently in small batches. The two often coexist, but the fixes differ — coarsen the key versus batch the writes.
- Why is raising parts_to_throw_insert a poor permanent fix?The threshold is a safety valve, not a cause. Primary-index analysis runs per part, so a table carrying thousands of small parts selects far more granules and reads far more data than the same rows in a handful of large parts. Raising the limit converts a clear insert failure into diffuse query slowness that is much harder to attribute.
- Is OPTIMIZE TABLE PARTITION a reasonable scheduled job to keep part counts down?No, as routine policy. It forces a full rewrite of the partition, consuming I/O and competing with the background merges that would have done the work anyway. As a one-off after a bad backfill left a partition fragmented, it is exactly the right tool. If you need it regularly, the partition key or the write batching is wrong.
saying these in an interview costs you the question
- Raising parts_to_throw_insert and calling the issue resolved
- Believing background merges will consolidate across partitions
- Scheduling OPTIMIZE TABLE FINAL nightly to control part counts
- Choosing daily partitions with no daily retention operation
- Assuming more partitions make queries faster