Partitioning splits a large dataset into pieces (by date, region, customer hash, etc.) so jobs and queries can skip irrelevant data.
s3://lake/orders/ dt=2026-09-01/ part-000.parquet dt=2026-09-02/ part-000.parquet dt=2026-09-03/ part-000.parquet Query "WHERE dt = '2026-09-03'" reads ONE folder, not the whole lake.
Why partition
- Faster queries (partition pruning)
- Cheaper incremental loads (overwrite one day)
- Parallelism (workers process different partitions)
- Safer backfills (rebuild
dt=...only)
Good partition keys
Often time (dt, hour) plus maybe region. Cardinality should not explode into millions of tiny folders.
Bad partitioning
- Partition by high-cardinality id (
user_id) → millions of tiny files - Too coarse (one huge
year=) → little pruning benefit
Interview tip: "Partitioning = physical organization for prune + parallel process." Give a date-partition example and warn about small files.