skip to content

How do you choose the partitionBy columns for a Spark-written table many teams will query?

level: principalimportance: nice to knowfreq 32%

answer

  1. what do the queries actually filter on?
  2. every distinct value becomes a directory
  3. divide total bytes by directory count
  4. cardinality multiplies, it does not add
  5. finer pruning without more directories

basics

~20 s

Partition by the low-cardinality columns that most queries filter on — usually a date — sized so each directory holds hundreds of megabytes. High-cardinality columns explode into tiny files and slow directory listing; use clustering inside a partition for finer pruning.

solid answer

~50 s

`partitionBy` creates a directory per distinct value combination, and directories are the only thing a plain file-based reader can prune without opening data. So the choice is driven by the **predicates your consumers actually write**, weighted by cardinality. Pick columns that appear in most `WHERE` clauses, are low-cardinality, and leave each directory holding a few hundred MB to a few GB. A date is almost always the right first choice; a second column is justified only when it is genuinely selective and does not multiply the directory count past what listing can bear. High-cardinality columns like `user_id` are a trap: thousands of directories, tiny files, slow planning, and no real gain. For finer pruning, cluster *within* a partition — `repartitionByRange` on a secondary column so row-group min/max statistics prune it — rather than adding directories.

code

python · 8 lines
python
spark.conf.set("spark.sql.sources.partitionOverwriteMode", "dynamic")

(df.repartitionByRange(64, "user_id")      # cluster inside the day
   .sortWithinPartitions("user_id")
   .write
   .partitionBy("dt")                      # prune across days
   .mode("overwrite")
   .parquet("s3://lake/events/"))

go deeper

for a junior

Know that partitionBy creates one directory per value and that filtering on that column lets Spark skip whole directories. Date is the usual choice.

for a middle

Explain how directory count multiplies across partition columns, why tiny directories make planning and compression worse, and how to estimate average bytes per directory before choosing.

for a senior

Show the operational judgment: overwrite semantics, compaction, and using clustering with row-group statistics for columns too granular to partition by.

for a principal

Own layout as a platform standard — target file sizes, directory-count budgets, enforcement through shared write tooling, and the migration cost of changing a scheme every consumer already depends on.

## What partitionBy actually does `df.write.partitionBy("dt", "country")` writes Hive-style directories: `.../dt=2026-08-20/country=PT/part-....parquet`. The values are encoded in the path and removed from the file contents. When a reader issues `WHERE dt = '2026-08-20'`, Spark prunes at the directory level and never lists, opens or reads anything under other dates. That is the whole benefit, and it is a large one — partition pruning is the cheapest possible filter. It is unrelated to in-memory partitions and task parallelism, despite sharing the word. ## The four forces in tension **Predicate match.** A partition column earns its place only if consumers filter on it. Partitioning by a column nobody filters on gives you all the costs and none of the pruning. This is an empirical question — read the query logs, don't guess. **Cardinality and directory count.** Every distinct value combination is a directory. Two years of daily data is ~730 directories: fine. Add `country` at 200 values and you are at 146,000 directories, most of them nearly empty. Add `user_id` and you have millions. Directory listing on object storage is a per-request cost paid on every query's planning phase, and metastore partition metadata grows with it. Multiplicative growth is the failure mode. **File size.** Directory count and file size are the same problem viewed from either end. Total bytes divided by directories gives average directory size; if that is below roughly 100 MB you have over-partitioned. Target directories in the hundreds-of-MB to low-GB range, and inside them files of similar magnitude — small Parquet files lose dictionary and run-length compression efficiency and give predicate pushdown weak statistics. **Write pattern.** Partitioning aligned to how data arrives makes incremental writes clean: a daily batch touches one `dt` directory, so it can be replaced atomically. Partition on something orthogonal to arrival and every batch touches every directory, which is both slow and dangerous to overwrite. ## Overwrite semantics — the operational trap With `spark.sql.sources.partitionOverwriteMode` at its default `STATIC`, `mode("overwrite")` on a partitioned write **deletes the entire table path** before writing, not just the partitions in the incoming DataFrame. A daily job that means to replace one day and instead removes the history is a well-known incident. Setting the mode to `DYNAMIC` restricts the overwrite to the partitions present in the data being written, which is what a daily reload almost always wants. Decide this deliberately and encode it in your job template rather than leaving it to whoever writes the next pipeline. ## Pruning below the directory level When consumers need to filter on a column too high-cardinality to partition by, the answer is **clustering, not more directories**. Sort or range-partition the data by that column so each Parquet row group covers a narrow value range; the row-group min/max statistics then let the reader skip most of the file. `df.repartitionByRange(n, "user_id")` before the write, or a `sortWithinPartitions`, gets you this on plain Parquet. Table formats add their own mechanisms — Iceberg's hidden partitioning with transforms such as bucketing and truncation, Delta's clustering — which decouple the physical layout from the query predicate and let it evolve without rewriting every consumer's SQL. Spark's own `bucketBy` serves a narrower purpose: it pre-hashes a table by a join key so repeat joins skip the shuffle, at the price of a metastore-managed table and a fixed bucket count. ## Deciding, and revisiting A workable procedure: sample real consumer queries and rank the columns by how often they appear as an equality or range predicate; compute, for each candidate scheme, the resulting directory count and the average bytes per directory; discard any scheme whose average directory falls below your file-size floor; among what survives, prefer the one that matches the write cadence. Usually this lands on date alone, sometimes date plus one coarse dimension. Then accept that it will change. Layout is a physical choice that outlives the reasoning behind it, and the cost of a bad one is paid by every reader on every query while the cost of fixing it is one rewrite job. Budget for periodic compaction and for occasional re-layout, keep the partition scheme out of consumer code where the storage format lets you (hidden partitioning, views), and measure the directory count as a first-class metric of a dataset's health. ## The multi-tenant angle On a shared platform, one team's over-partitioned table degrades everyone: metastore pressure, listing storms against shared object-store buckets, and planning-time driver load. Treating target file size and maximum directory count as platform standards — enforced by a shared write helper rather than by review comments — is usually more effective than educating each pipeline author individually.

  • A team wants to partition a fact table by user_id because their dashboards filter on it. What do you tell them?
    That the directory count would equal the user count — millions of near-empty directories, listing-dominated planning, and unreadably small files. Partition by date, and get user-level pruning from layout instead: range-partition or sort by `user_id` within each day so Parquet row-group statistics skip most of the file. If the access pattern is genuinely point lookups by user, a file-based lake is the wrong store.
  • Why is mode("overwrite") on a partitioned Spark write dangerous by default?
    With `spark.sql.sources.partitionOverwriteMode` at its default `STATIC`, overwrite clears the whole table path before writing, so a job intending to replace one day's partition wipes the entire history. Setting it to `DYNAMIC` limits the overwrite to the partitions present in the DataFrame, which is what an incremental reload means. Set it explicitly in the job.
  • When is bucketBy the right tool instead of partitionBy?
    When the goal is skipping a shuffle rather than pruning a scan. `bucketBy` pre-hashes the table on a join key into a fixed number of buckets, so joins and aggregations on that key can proceed without an exchange. It requires a metastore-managed table, both sides bucketed identically, and it fixes the bucket count at write time — much less flexible than partitioning, and worth it only for a repeated, expensive join.
  • How do you know an existing table is over-partitioned?
    Divide total dataset bytes by directory count and by file count. Average files well under about 100 MB, or directories holding a few thousand rows, mean the scheme is too fine. Corroborate with query planning time: if listing and metadata dominate the runtime of a selective query, the layout is the bottleneck, not the scan.

saying these in an interview costs you the question

  • Partitions by a high-cardinality column such as user id or timestamp
  • Adds partition columns without checking what consumers filter on
  • Ignores that directory count multiplies across partition columns
  • Assumes overwrite mode only replaces the partitions being written
  • Treats partitionBy and in-memory partition count as the same thing

context