Partitioning done right
When to partition a table, how to choose the column, and how over-partitioning backfires.
On this page
Show code in
Every code block on the page follows this.
You will learn
- What table partitioning does physically
- How to choose a partition column, and how much data each partition should hold
- How over-partitioning backfires
- Alternatives: clustering, hidden partitioning, liquid clustering
Read first
- Parquet and columnar storage · 4 min read
- The small file problem · 6 min read
Comfortable with these? Read on.
What partitioning does
events/ ├── event_date=2025-03-01/ │ ├── part-00000.parquet │ └── part-00001.parquet ├── event_date=2025-03-02/ │ └── ... └── event_date=2025-03-03/
A filter on the partition column, WHERE event_date = '2025-03-02', reads only that folder: partition pruning, decided before any file is opened. The partition column's values live in folder names (or metadata), not in the data files.
df.write.format("delta").partitionBy("event_date").saveAsTable("events")
CREATE TABLE events (event_id BIGINT, user_id BIGINT, event_date DATE, payload STRING) USING delta PARTITIONED BY (event_date);
Choosing a partition column
- Used in filters by most queries: otherwise pruning rarely helps and every query pays for more files.
- Low cardinality: hundreds to a few thousand values over the table's life, not millions.
- Enough data per partition: a common guideline is at least about 1 GB per partition. Delta's guidance suggests not partitioning tables under about 1 TB at all and relying on clustering instead.
- Stable: changing the partition scheme later means rewriting the table (except with Iceberg partition evolution).
| Column | Good partition? | Why |
|---|---|---|
| event_date on a 20 TB table | Yes | Used in nearly every query, ~55 GB per day |
| event_date on a 50 GB table | Probably not | About 140 MB per day: many small partitions |
| country (60 values) | Maybe, as a second level only if queries filter on it and partitions stay large | Multiplies the partition count |
| user_id | No | Millions of partitions with tiny files |
| event_timestamp | No | Practically unique values |
How over-partitioning backfires
Plan · the idea
event_date and country "to make country filters fast". 3 years × 60 countries = 65,700 partitions.Reality · the files
Effect · every query
Better · the fix
Alternatives and complements
Clustering
Hidden partitioning (Iceberg)
Liquid clustering (Delta)
Common mistakes
Partitioning small tables
Partitioning by high-cardinality columns
Filtering on a function of the partition column
Key takeaways
- Partitioning writes folders per value and lets filters skip whole folders.
- Choose low-cardinality columns used in most filters, with about 1 GB or more per partition.
- Over-partitioning creates huge numbers of small files.
- Use clustering for high-cardinality filters; partition only by a date-like column.
Check yourself
3 questions1. Which is the worst partition column?
Show the answer
user_id. Millions of distinct values create millions of tiny partitions.
2. A 40 GB table is partitioned by day over 3 years. Roughly how large is each partition?
Show the answer
About 36 MB. 40 GB / ~1,100 days ≈ 36 MB: too small to be worth partitioning daily.
3. How do you make user_id filters fast without partitioning by it?
Show the answer
Cluster files by user_id (Z-order or liquid clustering). Clustering narrows per-file ranges so statistics skip files.
Practice it
Interview problems that use this: write the PySpark, run it, and get graded on hidden tests.
Go deeper
Hidden partitioning and partition evolutionLakehouse
Liquid clusteringSpark internals
The small file problem
Primary sources: Delta: best practices (choose the right partition column) · Iceberg: partitioning