Merge — and the 'too many parts' crash

A background worker continuously fuses small adjacent parts into bigger sorted parts; if the write rate exceeds the merge rate, the part count crosses parts_to_throw_insert and the engine rejects new INSERTs.

Previously

Many small parts pile up fast. If the count grows unbounded, every query has to interrogate every part — there must be a background process that fuses them, or the engine collapses.

Scene 07

Merge — and the 'too many parts' crash

  1. Watch
  2. Try it
  3. Predict
  4. Capture
ACTIVE PARTS40 / 3,000delay_insert (1000)throw_insert (3000)WRITE FIREHOSE5 batches/secINSERT … VALUES (10k rows)PART STACKnewest on top+20 more parts belowMERGE WORKER3 in → 1 outmerge rate: 5 parts/secMERGED PARTlevel L+1rows: sum of inputsold parts dim and deleteBalanced: merge rate ≥ write rate. The per-partition part count hovers.
What to watch for

Watch the steady-state baseline. Batched INSERTs land as new parts on the stack; the background merge worker picks 3-5 adjacent parts and fuses them into one bigger part at the next level. The active-parts counter hovers around 40 — well below the two thresholds.

Continue unlocks when the animation finishes.
Implementation

Highlighted lines are the ones running in the diagram right now.

MergeWorker.run
background loop: pick adjacent same-level parts, sort-merge, swap
def merge_worker():
while True:
for partition in active_partitions():
# candidates = adjacent parts at the SAME level
picks = pick_adjacent_same_level(partition)
if len(picks) < 2:
continue
merged = sort_merge(picks) # in sorted-key order
emit_larger_part(merged) # level += 1
delete(picks) # originals dim and go
MergeTree.insert_path
admission control: throttle at delay, reject at throw
def insert_path(block, partition):
active_parts = count_active_parts(partition)
if active_parts > parts_to_throw_insert: # default 3000
raise "Too many parts (" + str(active_parts) + ")"
if active_parts > parts_to_delay_insert: # default 1000
sleep(backoff(active_parts)) # artificial throttle
new_part = write_part(block, partition) # 1 INSERT = 1 part
register(new_part)
MergeTree.pick_adjacent_same_level
partition-key effect: merge is bounded by partition walls
def pick_adjacent_same_level(partition):
parts = list_parts(partition) # ONE partition only
# merges NEVER cross partition boundaries —
# a high-cardinality partition key (e.g. hourly)
# leaves each partition with too few adjacent
# same-level parts to fuse, so merge starves.
return longest_adjacent_run_at_same_level(parts)

Where this sits in Build a columnar OLAP store (ClickHouse / Druid style)

Scene 07 of 13, in the Write side act — Bulk inserts → immutable parts → background merge → too-many-parts cliff.. Background worker fuses small parts into bigger ones; when write rate exceeds merge rate, parts_to_throw_insert fires and inserts get rejected.

Up next. Merge keeps the part count bounded, but a query still has to descend into surviving parts. What it does inside each part — scan everything, or jump to a range — depends on how the rows were sorted when the part was written.

All 13 scenes in Build a columnar OLAP store (ClickHouse / Druid style) · Every curriculum

Built with Arqly
Every scene in Build a columnar OLAP store (ClickHouse / Druid style) builds on the one before it.All 13 Build a columnar OLAP store (ClickHouse / Druid style) scenes