ClickHouse
A columnar OLAP database for fast scans and aggregates over append-heavy event data.
Also worth naming: ClickHouse Cloud · Altinity.Cloud · Managed ClickHouse
ClickHouse is the concrete answer when a dashboard, metric, or investigation query needs to scan billions of events quickly. The useful habit is to think from sorted parts, sparse indexes, batches, and pre-aggregated views, not from row-by-row updates.
What it is
ClickHouse is a column-oriented OLAP database. It is built for facts you append and analyse: clicks, logs, metrics, traces, payments, auctions, feature events, and security events. A query usually reads a few columns across many rows, filters by time or tenant, and aggregates by dimensions. That matches the analytics and OLAP shape: keep the product database safe, move events into an analytical path, and let the analytical store do broad scans.
It is not a drop-in replacement for Postgres. Postgres is the system of record for transactions, constraints, and row-level correctness. ClickHouse is the serving and investigation store for analytical facts. It can ingest from Kafka, object storage, files, CDC, or application batches, then answer SQL over columnar data.
The core storage engine family is MergeTree. The official MergeTree docs describe tables as sorted data parts, with inserts creating parts and background merges compacting them. The same docs say the primary key is a sparse index over granules, not a uniqueness constraint, and that the ORDER BY key usually defines the primary key. That one design explains both the speed and the sharp edges: choose the sort key for your filters, batch inserts so you do not create too many parts, and avoid treating it like an OLTP table.
Reach for ClickHouse when the interview needs real-time analytics with SQL and predictable low-latency aggregates. Reach for a managed warehouse when the priority is broad BI governance and you can tolerate slower freshness. Reach for Elasticsearch when text relevance, log search, or inverted indexes are the product. Reach for Postgres when the data is still the transactional source of truth.
When to reach for it
Reach for this when…
- Append-heavy events need fast SQL scans, GROUP BY queries, and dashboard aggregates.
- Fresh analytical tables or materialized rollups must update seconds after events arrive.
- You can model data as denormalised facts sorted by the filters users actually apply.
- You need a self-hosted or managed OLAP engine beside Kafka, Flink, S3, or a warehouse.
Not really this pattern when…
- You need row-level transactions, foreign keys, and frequent point updates: use Postgres.
- The product is full-text relevance, typo tolerance, or faceted search: use Elasticsearch or OpenSearch.
- The company wants a broad governed BI warehouse with many teams and batch freshness: a warehouse may be simpler.
- The workload is a small app with a few reports: Postgres plus read replicas may be enough.
How it works
1. Columnar storage plus vectorised execution
ClickHouse stores values by column and executes on chunks of columns. The official architecture overview calls this vectorized query execution: functions process arrays rather than individual rows, which lowers per-value overhead and can use SIMD. The performance overview adds that query operators pass batches of rows, parallelise across CPU cores, and read only the columns needed by the query. That is why a revenue dashboard over event_time, country, and amount does not pay to read a JSON payload column it does not use.
2. MergeTree writes parts, then merges them later
An insert into a MergeTree-family table writes a new data part. Within each part, rows are sorted by the ORDER BY key. Later, background merges combine small parts into larger parts. The MergeTree docs state that parts belonging to different partitions are not merged, and that the merge mechanism does not guarantee all rows with the same primary key are in one part. So query correctness must not depend on a background merge having happened unless you use the right engine and query logic.
Each insert creates a sorted part. Background merges combine small parts into larger parts, and merge-time engines can deduplicate, sum, aggregate, or apply TTL work while the user query path keeps reading stable snapshots.
3. The primary key is sparse, and the ORDER BY key is the real design
In ClickHouse, the primary key does not make rows unique. It stores marks for granules. The MergeTree docs use the default index_granularity=8192, and explain that sparse indexes stay small enough to keep in memory. A filter on the leftmost sort-key columns can skip ranges before reading column files. A filter on a column not aligned with the sort key often becomes a broad scan.
The ORDER BY key sorts each part. ClickHouse stores primary-key marks for granules, not for every row, so the index stays small and can skip ranges before reading column files.
Partitioning is a coarse data-management tool, not the main query accelerator. The same docs say you often do not need a partition key, and if you do, it is usually no more granular than month. Over-partitioning by tenant or user creates many small parts and harms merges. Put high-cardinality filter columns early in ORDER BY instead.
4. Materialized views pre-aggregate as data arrives
A ClickHouse materialized view is insert-triggered. The official incremental materialized view guide says it runs a query over inserted blocks and writes the result into a target table. With SummingMergeTree, rows with the same sorting key are summed during merges; with AggregatingMergeTree, aggregate states are merged. The guide also warns that the target table ORDER BY should line up with the view GROUP BY when those engines do merge-time aggregation.
This is how you make hot dashboards cheap. Store raw facts in a normal MergeTree table, then write minute, tenant, or product rollups into a smaller target. You still keep raw facts for backfill and new questions.
5. Replication and distribution are separate
ReplicatedMergeTree copies parts between replicas inside a shard. The replication docs state that replication is table-level, asynchronous, multi-master, and uses ClickHouse Keeper for replica metadata. By default, an insert waits for confirmation from one replica; use insert_quorum when you need acknowledgement from multiple replicas.
Distributed tables are the query and insert routing layer across shards. The Distributed engine docs say these tables store no data themselves; reads are automatically parallelised over remote servers, and remote indexes are used when present.
ReplicatedMergeTree uses ClickHouse Keeper for replica metadata. Distributed tables route reads and inserts across shards, while each shard owns its own replicated local tables.
Performance envelope
ClickHouse performance envelope: the rows below are sourced anchors, not a promise for every cluster. Hardware, schema, ORDER BY key, compression, batching, parts, joins, and concurrency decide the result.
| Dimension | Sourced anchor | Design meaning |
|---|---|---|
| Column scan | [Docs example](https://clickhouse.com/docs/materialized-view/incremental-materialized-view): 238.98M rows scanned in 0.133s for a simple aggregate | Shows why broad analytical scans can be fine when the schema and query match the engine. |
| Pre-aggregation | [Same docs example](https://clickhouse.com/docs/materialized-view/incremental-materialized-view): a SummingMergeTree target answered in 0.004s over 8.97K rows | Hot dashboard metrics move work from query time to insert and merge time. |
| Sparse index granule | [MergeTree docs](https://clickhouse.com/docs/engines/table-engines/mergetree-family/mergetree): default index granularity is 8192 rows | Primary-key filtering skips granules, but it is not a per-row B-tree lookup. |
| Synchronous inserts | [Insert strategy docs](https://clickhouse.com/docs/best-practices/selecting-an-insert-strategy): use at least 1,000 rows, ideally 10,000-100,000 rows per batch, and around one insert query per second | Many tiny inserts create many parts and overload background merges. |
| Async inserts | [Async insert docs](https://clickhouse.com/docs/optimize/asynchronous-inserts): flush at 100 MiB, 200 ms or 1000 ms on Cloud, or 450 queries by default | Server-side batching helps agents and apps that cannot batch client-side. |
| JOINs | [JOIN best practices](https://clickhouse.com/docs/best-practices/minimize-optimize-joins): keep JOINs to a minimum and avoid more than 3-4 joins per query for high-performance workloads | Denormalise facts or use dictionaries for hot many-to-one lookups. |
Capabilities in interviews
Event analytics and dashboards
Store raw facts and scan the columns a dashboard needs, with the product database off the path.
This is the default shape:
product DB / services → Kafka or CDC → ClickHouse MergeTree → dashboard SQLUse a table sorted by the filters users apply most: tenant, event type, time, or product id. Keep the product database as truth for orders and users; ClickHouse stores the analytical copy. This agrees with the OLAP split in analytics and OLAP: OLTP protects invariants, OLAP scans facts.
Choose this variant when
- Dashboards over clicks, orders, logs, traces or payments
- Ad hoc investigation over many rows
- Append-heavy facts with mostly analytical reads
Materialized pre-aggregation
Use materialized views into SummingMergeTree or AggregatingMergeTree for hot known metrics.
The official materialized view guide frames views as insert triggers into a target table. The target is commonly:
ENGINE = SummingMergeTree
ORDER BY (bucket_minute, tenant_id)for simple counters and sums, or:
ENGINE = AggregatingMergeTree
ORDER BY (bucket_minute, tenant_id)for aggregate states such as uniqState, avgState, or quantiles. Query with sum(...) or uniqMerge(...) so you remain correct while background merges catch up.
Choose this variant when
- The same metric is read constantly
- You know the grain and dimensions in advance
- Seconds-old dashboards matter more than arbitrary flexibility
Kafka and streaming ingestion
Feed ClickHouse from Kafka, ClickPipes, connectors, or the Kafka table engine.
Kafka stays the retained event log, as described on the Kafka page. ClickHouse is the analytical sink. The official Kafka table engine docs say the engine subscribes to topics and processes streams as they become available, and recommend ClickPipes on ClickHouse Cloud.
A common self-managed pattern is a Kafka engine table plus a materialized view that inserts parsed rows into a MergeTree table. This keeps raw event retention and replay in Kafka, while ClickHouse stores query-ready columnar facts.
Choose this variant when
- Kafka is already the durable ingestion backbone
- You need fresh analytics from a stream
- A replayable log remains the repair path
Corrections, upserts and deletes
Prefer append-based correction patterns; use mutations and lightweight deletes with care.
ClickHouse supports updates and deletes, but the storage model is append and merge first. The update overview recommends specialised engines such as ReplacingMergeTree for large volumes of row changes, and declarative updates for smaller, less frequent changes.
ALTER TABLE ... UPDATE and ALTER TABLE ... DELETE are mutations. The ALTER UPDATE docs explicitly call them heavy operations not designed for frequent use, and the ALTER DELETE docs warn that affected parts are rewritten. Lightweight deletes mark rows with a hidden mask and remove data physically during later cleanup, according to the lightweight delete docs.
Choose this variant when
- Late corrections are occasional and batched
- Latest-state tables can tolerate eventual deduplication
- Deletes are infrequent or align with partitions
Replicated and sharded OLAP
Replicate for availability inside a shard; use Distributed tables to fan queries across shards.
Use ReplicatedMergeTree when a shard needs multiple replicas and automatic part copying. The replication docs state that replication is independent of sharding and each shard has its own replication.
Use the Distributed engine when a query should fan out to several shards. The Distributed engine docs say a Distributed table stores no data of its own and parallelises reads on remote servers. The shard key then becomes a real design choice: co-locate filters and avoid scattering every query when you can.
Choose this variant when
- One node no longer holds the hot working set
- Availability needs replicas per shard
- Queries can aggregate partial results from several shards
Operating knobs
ORDER BY key
This is the first design decision. Put the columns used by common filters, especially tenant, product, event type, and time, early enough to prune granules. The key also controls compression and merge-time engines such as ReplacingMergeTree, SummingMergeTree and AggregatingMergeTree. Do not copy an OLTP primary key unless it matches the analytical access pattern.
Partitioning
Partition for lifecycle and coarse pruning, often by month or day at very high volume. The official docs warn that partitioning is not the same as ORDER BY and should not be too granular. If you partition by tenant or user, background merges cannot combine across those partitions and part counts can explode.
Insert strategy
Synchronous inserts need client-side batching: at least 1,000 rows and ideally 10,000-100,000 rows per batch. If edge agents or app workers cannot batch, enable async inserts with wait_for_async_insert=1 so the server batches and only acknowledges after flush.
Pre-aggregation grain
A materialized view is worth it when you know the hot question and the grain: for example, revenue per tenant per minute. Store raw facts as well, because a new dimension or a bug in the view needs a backfill. Align the view GROUP BY with the target ORDER BY when using SummingMergeTree or AggregatingMergeTree.
Mutability strategy
For corrections, decide whether you append a new version, run a lightweight update, run a mutation, or drop a partition. ReplacingMergeTree gives eventual deduplication by sorting key; use FINAL or equivalent query logic when exact latest state is required before merges finish. Frequent OLTP-style row changes are a warning sign.
JOIN and denormalisation boundary
ClickHouse can join, but hot low-latency analytics usually denormalises dimensions into facts. If a small dimension changes slowly, dictionaries can act as fast key-value lookups. If the query needs many relational joins with full fidelity, that part may belong in Postgres or the warehouse model, not the live ClickHouse serving path.
Versus the alternatives
ClickHouse versus the alternatives you are likely to name.
| Dimension | ClickHouse | Warehouse / lakehouse | Elasticsearch / OpenSearch | Postgres |
|---|---|---|---|---|
| Best job | Fresh OLAP over append-heavy facts | Governed BI, batch modelling, many analysts | Text relevance, faceted search, log search | Transactions, constraints, joins, source of truth |
| Query shape | Scans and GROUP BY over selected columns | SQL analytics across curated models | Inverted-index lookup, filters and facets | Point reads, transactions, relational joins |
| Freshness | Seconds to minutes when ingest is designed well | Often minutes to hours, depending on pipeline | Near-real-time search refresh | Immediate within a transaction |
| Writes | Batched inserts and append-first corrections | Batch or micro-batch loads | Async indexing from a source | Row inserts, updates and deletes |
| JOINs | Supported, but denormalise hot paths | Usually strong for modelled analytics | Not relational joins as the main job | A core strength |
| Use when | Dashboards and investigations need fast event aggregates | Organisation wants managed warehouse semantics | Users search text or logs by relevance | Correctness and relationships matter first |
Failure modes & gotchas
Every synchronous insert creates a part. The official insert guide recommends at least 1,000 rows per batch, ideally 10,000-100,000, and around one insert query per second for synchronous inserts. One-row inserts make merges fall behind and can trigger "too many parts" errors. Batch client-side or use async inserts.
If users filter by tenant and time but the table is ordered only by event id, the sparse primary index cannot prune much data. Pick the sort key from actual WHERE clauses and groupings, not from the OLTP key.
Partitioning by user or tenant creates many small partitions, and parts from different partitions are not merged. Use partitioning for lifecycle boundaries such as month or day, and use ORDER BY for query pruning.
Heavy mutations rewrite affected parts. Lightweight updates and deletes fit smaller changes but add SELECT overhead; neither makes ClickHouse a shopping-cart row store. If every request updates a row, keep that state in Postgres or another OLTP store.
ReplacingMergeTree removes duplicates only during background merges; duplicates can remain visible. Use FINAL or query-time aggregation when exact latest state matters, and budget for its cost.
ClickHouse supports several join algorithms, but the official JOIN guide still recommends denormalisation for latency-sensitive analytics and avoiding more than 3-4 joins per query. Flatten hot dimensions, use dictionaries for small many-to-one lookups, or move relational exploration elsewhere.
ReplicatedMergeTree is asynchronous and multi-master; recent inserts can lag on other replicas unless you use quorum settings. Distributed tables scatter across shards, so a poor shard key can make every query touch every shard. State consistency and fan-out costs explicitly.
In production
Cloudflare
HTTP analytics at 6M requests per second
Cloudflare’s 2018 engineering post says its HTTP analytics pipeline had grown to an average of 6M HTTP requests per second, with peaks up to 8M requests per second. The old pipeline used 106 Kafka brokers with 3x replication and 106 partitions, then aggregated into Postgres and Citus; Cloudflare wrote that Flink could not keep up with ingestion on all 6M HTTP requests per second.
The replacement kept Kafka consumers but moved aggregation into a 36-node ClickHouse cluster with 3x replication and materialized views. At the time, Cloudflare said the system served 7M+ customer domains, 2.5B monthly unique visitors, and 1.5T monthly page views. The new Zone Analytics API moved from struggling above 15 queries per second to about 40 queries per second, with load tests reaching about 150 queries per second on that setup.
The design lesson is not just “ClickHouse is fast.” Cloudflare got there by matching schema and materialized views to the API shape, removing old single points of failure, and keeping the analytical path separate from the product systems that generated the events.
Good vs bad answer
Numbers in this section are illustrative.
Interviewer probe
“You need analytics for an ad-click platform: high event volume, live dashboards, fraud slices, and ad hoc SQL. Where does ClickHouse fit?”
Weak answer
"I would put all clicks in Postgres, add indexes, and maybe cache dashboard queries. If it gets too slow, I will move old rows to a warehouse."
Strong answer
"Postgres stays the source of truth for accounts, campaigns and billing. Click events go through Kafka and land in ClickHouse as a MergeTree fact table. I batch inserts or use async inserts, because small inserts create too many parts. The table is ordered by (advertiser_id, event_date, campaign_id) or the filter pattern we measure, so the sparse primary index can skip ranges.
For hot dashboard metrics, I add materialized views into SummingMergeTree or AggregatingMergeTree, for example clicks and spend per advertiser per minute. Raw facts stay available for new dimensions and backfills. I denormalise stable campaign attributes into the fact or use dictionaries for small lookups, because a low-latency ClickHouse design should not join a snowflake schema on every request.
Updates are append-first: late corrections arrive as new versions or as rare mutations, not per-click OLTP updates. Replication is ReplicatedMergeTree per shard with ClickHouse Keeper; Distributed tables fan reads across shards. I would choose a warehouse instead for broad company BI with slower freshness, Elasticsearch for text/log search, and Postgres for the transactional campaign state."
Why it wins: It separates OLTP from OLAP, names the ClickHouse table shape, explains batching and ORDER BY, uses materialized views for hot aggregates, avoids join-heavy modelling, and scopes updates and replication honestly.
Interview playbook
When it comes up
- The prompt asks for dashboards, metrics, log analytics, clickstream analysis or observability slices.
- A database table is being asked to scan months of events for every user-facing report.
- You need a concrete OLAP store beside Kafka, Flink, S3 or Postgres.
- The interviewer asks how to make analytical reads fast without corrupting the product database.
Order of reveal
- 11. Place it off the OLTP path. Postgres owns transactions; ClickHouse owns analytical copies of events and facts.
- 22. Name MergeTree and the ORDER BY key. I store facts in a MergeTree table ordered by the filters the dashboard uses, because the sparse index prunes sorted granules.
- 33. Batch ingestion. I batch inserts, or enable async inserts with wait_for_async_insert=1 when agents cannot batch, so I do not create too many parts.
- 44. Pre-aggregate hot metrics. For known dashboards I use materialized views into SummingMergeTree or AggregatingMergeTree, while keeping raw facts for repair and ad hoc analysis.
- 55. Be honest about joins and updates. I denormalise hot dimensions, keep mutations rare, and use ReplacingMergeTree only with eventual deduplication in mind.
Signature phrases
- “The ORDER BY key is the ClickHouse schema decision.” — It ties storage layout to query speed.
- “Batch inserts, or the parts will beat you.” — Shows you know the MergeTree write path.
- “Materialized views move a known aggregate from query time to insert and merge time.” — Explains pre-aggregation without magic.
- “Replicas are for availability; Distributed is for shards.” — Separates two commonly confused cluster layers.
Likely follow-ups
?“Why not use the warehouse instead?”Reveal
I would use the warehouse for broad governed BI, finance models, and slower batch freshness. I would use ClickHouse when the product needs fresh analytical serving, such as a customer dashboard or incident slice that must scan recent events quickly. Many systems use both: ClickHouse for the live serving surface, warehouse or lakehouse for long-range modelling and organisation-wide reporting.
?“How do you handle late corrections?”Reveal
Prefer append-based correction patterns. For latest-state rows, insert a new version into ReplacingMergeTree and use FINAL or a query pattern that deduplicates when exactness matters before merges finish. For rare bulk corrections, run mutations in partition-bounded batches. For large expiry, drop partitions. I do not put high-frequency OLTP updates in ClickHouse.
?“Your ClickHouse dashboard slows down. What do you inspect first?”Reveal
I check whether the query uses the ORDER BY prefix, how many parts and partitions it reads, whether inserts are too small, whether a materialized view should serve the metric, and whether joins or high-cardinality aggregations are exploding memory. If it is a cluster query, I check shard scatter and replica lag. The fix is usually schema and ingestion shape before buying more nodes.
Worked example
Numbers in this section are illustrative.
Numbers in this example are illustrative.
Setup. A ticketing platform needs live analytics during a stadium on-sale: orders started, payment failures, queue depth, and suspicious checkout patterns by event, seller, country and minute. The OLTP system still owns seat holds, payments and bookings. ClickHouse serves the analytics surface, not the seat-claim transaction.
Ingestion and table shape. Product services write events to Kafka. A connector or Kafka engine table feeds a raw ticket_events MergeTree table. Inserts are batched; if many edge workers send tiny payloads, use async_insert=1, wait_for_async_insert=1. The hot filters are event_id, minute and event type, so the fact table is ordered by (event_id, bucket_minute, event_type, seller_id) and partitioned by month or day.
Hot rollups. A materialized view writes orders_by_event_minute into SummingMergeTree, keyed by (event_id, bucket_minute, event_type). Another AggregatingMergeTree table stores unique buyer states with uniqState(user_id) and latency quantile states. Dashboards read these small tables first, while raw facts remain for new questions and backfills.
Corrections and cluster. Payment processors can send late status changes. For current-state tables, insert a higher-version row into ReplacingMergeTree and query with FINAL only where exact state is required. Each shard uses ReplicatedMergeTree with Keeper metadata, and a Distributed table fans queries across shards. Keep most event-scoped dashboards local when the shard key allows it.
Result. Seat correctness remains in Postgres and the booking service. ClickHouse answers live analytical questions from batched, sorted facts and materialized rollups. Kafka remains the replay path, so a bad rollup can be rebuilt without inventing data.
Cheat sheet
- •ClickHouse = columnar OLAP for append-heavy facts, not the OLTP source of truth.
- •MergeTree inserts create sorted parts; background merges compact and transform them.
- •ORDER BY is the schema decision. The primary key is a sparse index, not uniqueness.
- •Partition for lifecycle, usually coarse time; avoid high-cardinality partitions.
- •Batch inserts: docs say [at least 1K rows, ideally 10K-100K](https://clickhouse.com/docs/best-practices/selecting-an-insert-strategy), or use async inserts.
- •Materialized views write into target tables on insert; Summing/AggregatingMergeTree serve rollups.
- •ReplacingMergeTree dedupe is eventual. Use FINAL or query-time logic when exactness matters.
- •Denormalise hot facts. Keep JOINs minimal; dictionaries help many-to-one lookups.
- •ReplicatedMergeTree copies parts inside a shard; Distributed tables fan across shards.
Drills
Numbers in this section are illustrative.
Why is the ORDER BY key more important than the primary key name in ClickHouse?Reveal
The ORDER BY key sorts rows inside each part, and the primary key is usually that same expression stored as sparse marks. Queries filter fastest when they can use the left side of that order to skip granules. It is not a uniqueness constraint like an OLTP primary key.
A service sends one-row inserts to ClickHouse thousands of times per second. What breaks?Reveal
Each synchronous insert creates a part, so ClickHouse must track and merge many tiny parts. Background merges can fall behind, query performance drops, and the table can hit too-many-parts errors. Batch client-side, use a queue, or enable async inserts with wait_for_async_insert=1 so the server batches before writing parts.
When do you choose a materialized view into SummingMergeTree?Reveal
When the hot metric is known and additive, such as count or sum per event per minute. The view computes partial rows as data arrives, and SummingMergeTree merges rows with the same sorting key. You still query with aggregation or FINAL when needed because background merging is asynchronous.
What it is