In a ClickHouse MergeTree table, what does the PARTITION BY clause actually do?
answer
- it groups parts, it does not sort rows
- merges never cross this boundary
- dropping a month should be instant
- think retention, not query speed
- finer is not better here
basics
~20 sPARTITION BY splits a ClickHouse MergeTree table into independent groups of parts, usually by month. It enables dropping or detaching data as a whole unit, per-partition TTL, and coarse pruning — it is a data-management tool, not a substitute for the sorting key.
solid answer
~40 s`PARTITION BY` assigns every row to a partition, typically `toYYYYMM(event_date)`. Parts belong to exactly one partition and are never merged across partitions, so each partition is an independent bucket of data. What it buys you is data lifecycle management: `ALTER TABLE ... DROP PARTITION` deletes a month instantly by unlinking files, `DETACH`/`ATTACH PARTITION` moves data between tables, TTL rules and tiered storage operate per partition, and mutations touch fewer parts. It also gives coarse pruning — a query filtering on the partition expression, or on the underlying column, skips whole partitions before any index analysis. What it does not do is order rows or act as an index: fine-grained skipping is the job of `ORDER BY` and the sparse primary index. Over-partitioning is the classic mistake.
code
sql · 11 linesCREATE TABLE events
(
event_date Date,
tenant_id UInt32,
user_id UInt64,
payload String
)
ENGINE = MergeTree
PARTITION BY toYYYYMM(event_date)
ORDER BY (tenant_id, event_date, user_id)
TTL event_date + INTERVAL 13 MONTH DELETE;go deeper
Know the typical form PARTITION BY toYYYYMM(date), that it groups data so a whole month can be dropped cheaply, and that it is not what makes filtered queries fast.
Explain that parts never merge across partitions, that partition pruning happens before index analysis, and why a high-cardinality partition key multiplies parts.
Choose the partition granularity from the retention and archival policy, verify part counts per partition in system.parts, and know when no partition key at all is correct.
Frame partitioning as the lifecycle contract for the table — retention, tiering, backfill and restore all key off it — and set an organisation-wide default rather than leaving each team to invent one.
## What a partition is Every row inserted into a `MergeTree` table is assigned a partition id by evaluating the `PARTITION BY` expression. An insert produces one new part per distinct partition value in the block. Parts belong to exactly one partition, and background merges only ever combine parts **within** the same partition — never across. So a partition behaves like an independent sub-table sharing the schema, the sorting key, and the indexes. Without a `PARTITION BY` clause, the table has a single partition. That is a perfectly valid and often correct choice. ## What it is for Four things, in rough order of importance: 1. **Whole-unit data lifecycle.** `ALTER TABLE events DROP PARTITION '202401'` removes a month of data by dropping its parts. It is close to instantaneous and costs no merge work, unlike `ALTER TABLE ... DELETE`, which is a mutation that rewrites every affected part. If your retention policy is "keep 13 months", partitioning by month is what makes it cheap. 2. **Moving data around.** `DETACH PARTITION` and `ATTACH PARTITION` (and `ATTACH PARTITION FROM` another table with the same schema) let you move or copy a slice of data at file level — the standard backfill and archival mechanism. 3. **TTL and tiered storage.** `TTL event_date + INTERVAL 6 MONTH TO VOLUME 'cold'` is evaluated per part, and partitioning aligned with the TTL expression means whole partitions age out together instead of leaving each part half-expired. 4. **Coarse pruning.** Before touching any index, ClickHouse discards partitions whose partition-key values cannot match the predicate. Each part additionally carries min/max values for the columns used in the partition expression, so a filter on `event_date` prunes even though the partition key is `toYYYYMM(event_date)`. ## What it is not It is not an index and it is not a sort order. Within a partition, everything about performance is decided by `ORDER BY` and the sparse primary index. A common junior mistake is to reach for a finer partition key to "make queries faster" — a partition key of `toYYYYMMDD(event_time)` plus a filter on a customer id prunes nothing on that customer id and creates thirty times as many parts as monthly partitioning would. It is also not required. Small or medium tables with no retention story are usually better off with no partition key at all, letting merges consolidate everything into a few large parts. ## Choosing the granularity The rule of thumb the ClickHouse community converged on is: partition by the unit you **delete or move**, and keep the number of partitions modest — typically monthly for time-series, occasionally weekly or daily for very high-volume tables with short retention. Two costs push back on fine partitions: - **Part count.** Every insert that spans N partitions writes N parts. Merges cannot consolidate across partitions, so the steady-state part count scales with the number of active partitions. High part counts slow query planning (index analysis runs per part), consume file handles and memory, and eventually trigger insert throttling and the `Too many parts` error. - **Metadata and startup.** Server start and many system operations enumerate parts; tables with tens of thousands of partitions are noticeably worse to operate. A high-cardinality expression — a user id, a hash, a raw timestamp — is the anti-pattern. ClickHouse guards the worst case with `max_partitions_per_insert_block`, which rejects an insert block touching too many partitions at once. ```sql CREATE TABLE events ( event_date Date, tenant_id UInt32, user_id UInt64, payload String ) ENGINE = MergeTree PARTITION BY toYYYYMM(event_date) ORDER BY (tenant_id, event_date, user_id) TTL event_date + INTERVAL 13 MONTH DELETE; ``` ## Inspecting partitions `system.parts` is the operational view: one row per part, with `partition`, `active`, `rows`, and `bytes_on_disk`. Grouping active parts by partition tells you both the data distribution and whether any partition has accumulated an unhealthy number of parts. The virtual columns `_partition_id` and `_part` let a query see which partition and part a row came from. ```sql SELECT partition, count() AS parts, sum(rows) AS rows FROM system.parts WHERE active AND table = 'events' GROUP BY partition ORDER BY partition; ``` ## The interview answer Say what it does (independent buckets of parts, no cross-partition merges), why you'd want it (drop/detach/TTL as whole units, plus coarse pruning), and what it is not (an index — that's the sorting key's job). Then name over-partitioning as the failure mode. That covers the whole surface.
- When is having no PARTITION BY clause the right choice?When the table has no retention or archival story that operates on whole slices, and queries already prune well through the sorting key. A single partition lets merges consolidate everything into a few large parts, which keeps part counts low and index analysis cheap. Partitioning that exists only 'because tables are partitioned' adds cost without buying anything.
- If the partition key is toYYYYMM(event_date), does a filter on event_date still prune partitions?Yes. Each part stores min/max values for the columns used in the partition expression, so ClickHouse can evaluate a predicate on the raw `event_date` column against those bounds and discard partitions before any index analysis. You do not have to write the filter in terms of the partition expression itself.
- How does DROP PARTITION differ from a DELETE on the same rows?`DROP PARTITION` unlinks the partition's parts — it is a metadata-and-files operation that completes almost immediately and reclaims all the space. `ALTER TABLE ... DELETE` is a mutation: it rewrites every part containing matching rows in the background, costing full I/O over that data and leaving the space reclaimed only once the rewrites finish.
Think of it as which filing cabinet a document goes in, while the sorting key decides its position within the drawer. You throw out a whole cabinet at year end; you find a document using the ordering inside it.
saying these in an interview costs you the question
- Treating the partition key as an index for query filters
- Choosing a daily or hourly partition key by default
- Believing merges can combine parts from different partitions
- Thinking every table needs a PARTITION BY clause
- Confusing partitions with shards across servers