Skip to content
Good engineers know 4 min read · All formats 2 practice problems ↓

Partitioning done right

When to partition a table, how to choose the column, and how over-partitioning backfires.

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

Comfortable with these? Read on.

TL;DR Partitioning stores a table in folders by the values of one or more columns, so queries filtering on them skip whole folders. It works well for low-cardinality columns used in almost every query, like a date, and with each partition holding at least around 1 GB. Partitioning by a high-cardinality column creates millions of tiny files.

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.

PySparkSpark SQL
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).
ColumnGood partition?Why
event_date on a 20 TB tableYesUsed in nearly every query, ~55 GB per day
event_date on a 50 GB tableProbably notAbout 140 MB per day: many small partitions
country (60 values)Maybe, as a second level only if queries filter on it and partitions stay largeMultiplies the partition count
user_idNoMillions of partitions with tiny files
event_timestampNoPractically unique values

How over-partitioning backfires

Plan · the idea

A 200 GB table is partitioned by event_date and country "to make country filters fast". 3 years × 60 countries = 65,700 partitions.

Reality · the files

Average partition: 3 MB. Each daily load writes up to one file per task per partition, so after a year the table has hundreds of thousands of small files.

Effect · every query

Queries that do not filter on country now open 60 times more files than needed. Planning takes minutes, and the metastore or transaction log is bloated.

Better · the fix

Partition by month (or not at all), and cluster by country and date with Z-order or liquid clustering. File-level statistics then provide most of the skipping, with large healthy files.

Alternatives and complements

Clustering

Z-order or sorting within files: skipping via min/max statistics without extra folders. Works on high-cardinality columns.

Hidden partitioning (Iceberg)

Partition by a transform like day(ts) or bucket(16, id); queries filter on the raw column and still prune. The layout can evolve later.

Liquid clustering (Delta)

Replaces partitioning and Z-order with incremental clustering on chosen keys that can change over time.
Why this matters: partitioning is coarse and permanent; clustering is fine-grained and adjustable. Partition only by the one or two columns that define how data arrives and is deleted (typically a date), and cluster for everything else.

Common mistakes

Partitioning small tables

Creates small files with no benefit.

Partitioning by high-cardinality columns

Millions of folders and tiny files.

Filtering on a function of the partition column

WHERE year(event_date) = 2025 may not prune; filter on the raw column range, or use generated columns or hidden partitioning.

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 questions

1. 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.

Solve: How Much Does Each Query Read? →

Go deeper

Primary sources: Delta: best practices (choose the right partition column) · Iceberg: partitioning