ClickHouse is an open-source column-oriented database built for analytical queries over very large tables, such as percentiles and group-bys over billions of rows. Teams usually hit the same surprises in the first months: a query that should be instant reads the whole table, an insert pattern that worked in testing fails with a too-many-parts error, a deduplicating engine returns duplicates, or an UPDATE takes an hour.

Each surprise follows from one design, the MergeTree storage engine, and this page explains it from the inside. The general case for columnar storage, encodings and vectorised execution is made in columnar database architecture; this page assumes it and concentrates on what is specific to ClickHouse: how parts and the sparse index work, how to choose the sort key, what the engine variants do at merge time, how materialized views really behave, how to ingest without breaking the merge process, and how replication and sharding fit together.

Advertisement

Parts, granules and the sparse primary index

A MergeTree table is a set of immutable parts. Every INSERT sorts its block by the table's ORDER BY key and writes it as a new part on disk: one directory holding, for each column, a compressed data file and a marks file, plus the primary index and some metadata. Background merges continually pick several parts from the same partition and combine them into a larger sorted part, then drop the originals. It resembles the LSM tree described in LSM trees, but with no write-ahead log or memtable: each insert goes straight to a part.

Within a part, rows are grouped into granules, by default 8,192 rows (index_granularity), with adaptive sizing that also caps granules at about 10 MiB of data (index_granularity_bytes). The primary index stores the sort key of the first row of each granule, not of every row, so it is small enough to keep in memory even for tables with trillions of rows. That is why it is called sparse. The marks file maps each granule to an offset in each column's compressed file. To answer a query, ClickHouse prunes whole parts by partition and min-max metadata, binary-searches the primary index to find the range of granules whose keys can match, and reads only those granules, and only for the columns the query uses.

The unit of reading is therefore the granule. A point lookup by key still reads at least 8,192 rows of each needed column, which is fast for analytics and wasteful for key-value access; ClickHouse is not built for the latter.

A MergeTree table: inserts create parts, merges combine themINSERT blocksorted in memorynew partpart 1part 2part 3mergemerged part1 + 2 + 3Inside one partprimary.idx: one key per granulecolumn.bin: compressed blockscolumn.mrk: granule to offsetgranule = 8,192 rows by defaultchecksums, count, minmax of partitionQuery: WHERE tenant_id = 42 AND ts in one day1. Partition pruningminmax per part2. Primary indexbinary search marks3. Read granulesonly needed columnsSkip indexes and projections can prune further at step 2
Inserts create sorted parts that merges combine. A query prunes parts by partition, then uses the sparse primary index to choose granules, then reads only the needed columns.

Designing the sort key and partitions

The ORDER BY clause is the most important decision in a ClickHouse schema, because it determines both what the index can prune and how well columns compress. Put first the columns that most queries filter on with equality, ordered from lower to higher cardinality, and put the time column after them. For a multi-tenant events table queried by tenant and time range, (tenant_id, event_type, ts) lets the index jump to one tenant's rows and then one type's rows within it. Leading with the timestamp would make every per-tenant query read a slice across all tenants.

PRIMARY KEY defaults to the ORDER BY expression and may be a prefix of it. Making it shorter keeps the in-memory index smaller while the data stays sorted by the full key, which helps when a long sort key is wanted for compression or for the merge semantics described below. PARTITION BY is a data-management tool, not an indexing tool: partitions let you drop, detach or move old data cheaply, and TTL rules operate on them efficiently. Partition by month or day of the time column, and keep the number of partitions small; merges never cross partitions, so a high-cardinality partition key, such as a customer id, multiplies the parts on disk and leads straight to the failure described in the ingestion section.

CREATE TABLE events
(
    tenant_id   UInt32,
    ts          DateTime64(3),
    event_type  LowCardinality(String),
    user_id     UInt64,
    duration_ms UInt32 CODEC(T64, ZSTD),
    props       String CODEC(ZSTD(3)),
    INDEX idx_user user_id TYPE bloom_filter GRANULARITY 4
)
ENGINE = MergeTree
PARTITION BY toYYYYMM(ts)
ORDER BY (tenant_id, event_type, ts)
TTL toDateTime(ts) + INTERVAL 13 MONTH DELETE;

-- How much did the index prune?
EXPLAIN indexes = 1
SELECT event_type, count(), quantile(0.99)(duration_ms)
FROM events
WHERE tenant_id = 42 AND ts >= '2026-09-01' AND ts < '2026-09-02'
GROUP BY event_type;

The schema uses LowCardinality for the event type, a dictionary encoding that the documentation recommends for columns with fewer than about 10,000 distinct values, explicit codecs for a numeric and a text column, a bloom-filter skip index for occasional lookups by user, and a TTL that deletes data after 13 months. EXPLAIN indexes = 1 reports, for each index, how many parts and granules survived, which is the quickest way to check that a query uses the key you designed. If a query that should be selective reads most of the granules, the filter does not match a prefix of the sort key.

Advertisement

The engine family: semantics applied at merge time

Variants of MergeTree apply extra logic when parts merge. ReplacingMergeTree keeps only the last row (or the row with the highest version column) among rows with the same sort key. SummingMergeTree sums numeric columns across rows with the same key. AggregatingMergeTree combines partial aggregate states, such as those from uniqState or quantileState. CollapsingMergeTree and VersionedCollapsingMergeTree cancel pairs of rows marked with a sign column, which lets you model updates as a cancel row plus a new row.

The crucial property is that these semantics apply only when parts merge, and merges happen in the background at an unspecified time, never across partitions. Until then, duplicates or unsummed rows are visible to queries. Correct queries therefore aggregate again at read time (GROUP BY, argMax by version, -Merge combinators) or use FINAL, which applies the merge logic during the query at extra cost. Treat these engines as storage optimisations that make read-time aggregation cheaper, not as constraints the database enforces.

Materialized views are insert triggers

A ClickHouse materialized view is not a cached query result that refreshes. The classic kind is a trigger: when a block is inserted into the source table, the view's SELECT runs over that block alone and the result is inserted into a target table. Combined with AggregatingMergeTree this gives incremental rollups that cost almost nothing at query time.

CREATE TABLE events_daily
(
    tenant_id  UInt32,
    day        Date,
    event_type LowCardinality(String),
    events     SimpleAggregateFunction(sum, UInt64),
    users      AggregateFunction(uniq, UInt64),
    p99_ms     AggregateFunction(quantile(0.99), UInt32)
)
ENGINE = AggregatingMergeTree
ORDER BY (tenant_id, day, event_type);

CREATE MATERIALIZED VIEW events_daily_mv TO events_daily AS
SELECT tenant_id, toDate(ts) AS day, event_type,
       count() AS events,
       uniqState(user_id) AS users,
       quantileState(0.99)(duration_ms) AS p99_ms
FROM events
GROUP BY tenant_id, day, event_type;

-- Read it: rows for one key may still be unmerged, so always aggregate again.
SELECT day, sum(events), uniqMerge(users), quantileMerge(0.99)(p99_ms)
FROM events_daily
WHERE tenant_id = 42
GROUP BY day ORDER BY day;

Three consequences follow from the trigger model. The view sees only the inserted block, so a GROUP BY in the view aggregates within each block, and the target table holds many partial rows per key until merges combine them; hence the final query merges states again. A view created after data exists does not backfill; you must insert historical data into the target yourself, carefully, to avoid double counting rows that arrive during the backfill. And a failure in the view's insert can fail the original insert, so an expensive or fragile view slows down or breaks ingestion for the source table. Refreshable materialized views, which rerun a full query on a schedule, suit joins that cannot be computed block by block.

Ingestion and the too-many-parts failure

Because every insert creates a part, the number of inserts per second, not the number of rows, is what stresses ClickHouse. A thousand single-row inserts per second create a thousand parts per second, and merges cannot keep up. ClickHouse protects itself with limits on active parts. When a single partition exceeds parts_to_delay_insert, 1,000 by default, inserts are artificially slowed; above parts_to_throw_insert, 3,000 by default since version 23.6 (it was 300 before), inserts into that partition fail with a TOO_MANY_PARTS error explaining that merges are slower than inserts. A separate table-wide limit, max_parts_in_total, defaults to 100,000 across all partitions.

The fixes are on the client side. Batch inserts into large blocks, commonly tens of thousands to a few hundred thousand rows, and insert at most about once per second per table from each writer. When many small producers cannot batch, enable asynchronous inserts with the async_insert setting: the server buffers small inserts and flushes them as one part when a size, a time or a query-count threshold is reached. Leave wait_for_async_insert at its default of 1 so the client is acknowledged only after the flush is written; setting it to 0 acknowledges before data is durable. The queries below are the first thing to run when inserts slow down.

-- Parts per partition: the number that leads to "too many parts"
SELECT table, partition, count() AS parts, sum(rows) AS rows
FROM system.parts
WHERE active AND database = currentDatabase()
GROUP BY table, partition
ORDER BY parts DESC
LIMIT 10;

-- Merges in flight and unfinished mutations
SELECT table, elapsed, progress, num_parts FROM system.merges;
SELECT table, mutation_id, command, parts_to_do, latest_fail_reason
FROM system.mutations WHERE NOT is_done;

Updates and deletes

Parts are immutable, so changing rows means rewriting them. ALTER TABLE ... UPDATE and ALTER TABLE ... DELETE are mutations: they are asynchronous by default, and they rewrite every part that contains affected rows, in the background, which for a broad condition means rewriting most of the table. They are meant for occasional corrections such as a privacy deletion, not for application traffic. A failing mutation can stay stuck retrying; system.mutations shows why, and KILL MUTATION cancels it.

Lightweight DELETE FROM ... WHERE is cheaper: it marks rows as deleted through a hidden mask column, queries skip them immediately, and the space is reclaimed by later merges. It still runs as a mutation underneath and its cost grows with the number of parts touched. For data that changes routinely, model change as new rows instead, with ReplacingMergeTree and a version column, and resolve the latest version at read time.

Replication and sharding

Replication is per table. ReplicatedMergeTree, or the replicated variant of any engine in the family, keeps a log of parts in ClickHouse Keeper, a coordination service using the Raft consensus protocol and compatible with ZooKeeper's protocol. When a replica writes a part, it records the part in the shared log; the other replicas fetch the part itself directly from a replica that has it. Merges are also coordinated through the log so replicas end up with identical parts. Replication is asynchronous and multi-master: any replica accepts inserts. Replicated tables also deduplicate retried inserts by hashing each inserted block, so a client that retries an insert with identical data after a timeout does not create duplicates, provided the retry sends exactly the same block.

Sharding is a separate layer. A Distributed table holds no data; it routes inserts to shards by a sharding key and fans queries out to one replica of each shard, then merges the partial results. Distributed inserts are queued locally and forwarded in the background by default, so an acknowledged insert is not yet on its shard; inserting into the shard tables directly avoids that. Keeper is a dependency to run with care: if it loses quorum, replicated tables go read-only for inserts while queries continue.

Failure modes

FailureCauseFix
TOO_MANY_PARTSSmall, frequent inserts or a high-cardinality partition keyBatch, use async inserts, coarser partitions
Query reads whole tableFilter is not a prefix of ORDER BYCheck EXPLAIN indexes = 1; redesign the key or add a projection
Duplicates from ReplacingMergeTreeRows not merged yetAggregate at read time or use FINAL
Rollup double counts after backfillView trigger and manual backfill both insertedBackfill a bounded time range before the view's start
Mutation never finishesError in the command or a huge rewriteInspect system.mutations; KILL MUTATION; avoid broad updates
Inserts fail, queries workKeeper lost quorumMonitor Keeper; run three or five nodes on separate hosts

Trade-offs

ClickHouse buys scan speed and compression with a write model that wants large batches and a sort key fixed at design time; changing ORDER BY means creating a new table and copying data. Its merge-time engines make deduplication and rollups cheap but eventually consistent. Point lookups, frequent updates and multi-row transactions are poor fits, which is why it usually sits beside an OLTP database rather than replacing one, a split explained in OLTP vs OLAP architectures. For logs specifically, its SQL and compression trade against the full-text search of other stores, compared in Loki vs Elastic vs ClickHouse for logs.

What to do next

  1. Write down your top five queries and choose ORDER BY so their equality filters form a prefix, lowest cardinality first and time last.
  2. Partition by month or day only; never by a high-cardinality id.
  3. Run EXPLAIN indexes = 1 on each top query and confirm most granules are skipped.
  4. Batch inserts into large blocks or turn on async_insert with wait_for_async_insert left at 1.
  5. Alert on active parts per partition from system.parts, well before the delay threshold.
  6. Build rollups with materialized views into AggregatingMergeTree, and always aggregate again when reading them.
  7. Model changing data as versioned rows; reserve mutations for rare corrections and watch system.mutations.
  8. If you replicate, run Keeper on three or five dedicated nodes and alert on its quorum.
Key takeaway: ClickHouse stores each insert as an immutable sorted part and merges parts in the background. A sparse primary index over 8,192-row granules lets queries skip most data, but only for filters that match a prefix of ORDER BY, so the sort key is the central design decision. Engine variants and materialized views apply their logic at merge or insert time, so queries must aggregate again. Ingest in large batches or async inserts to stay under the parts limits, treat mutations as rare, and run Keeper carefully if you replicate.