skip to content

Partitioning and Partition Evolution

You will learn how table formats physically lay out data so a query touches a fraction of it, and — the part that separates old Hive tables from modern ones — how to change that layout later without rewriting history. Interviewers love the 'your partition key was wrong, now what?' question because Hive-era tables had no good answer.

on this pageshow

explore

questions

6

What does Hive-style directory partitioning do to a table's file layout on object storage?

level: juniorimportance: must knowfreq 72%

answer

  1. the value lives in the path, not the file
  2. folders named column=value
  3. the filter picks folders before any read
  4. pruning removes I/O instead of filtering rows

basics

~10 s

Hive-style partitioning splits a table's files into directories named column=value, such as dt=2024-05-01. A query that filters on that column reads only the matching directories, so it touches a fraction of the table.

solid answer

~40 s

Instead of one flat directory of data files, the table's storage is organised into nested directories whose names encode a column and its value: `events/dt=2024-05-01/country=DE/part-0000.parquet`. That value is normally **not** stored inside the data files at all — it is recovered by parsing the path, which is why a partition column is sometimes called virtual. When a query filters on a partition column, the engine can decide *before reading anything* which directories are relevant and skip the rest; that is partition pruning, and it removes I/O rather than filtering rows after paying for it. The tradeoff is that pruning only works for the columns you partitioned by, the layout is fixed by whatever the writer chose, and each partition needs enough data to produce reasonably sized files.

code

text · 5 lines
text
s3://warehouse/events/
  dt=2024-05-01/country=DE/part-0000.parquet
  dt=2024-05-01/country=US/part-0000.parquet
  dt=2024-05-02/country=DE/part-0000.parquet
  dt=2024-05-02/country=US/part-0000.parquet

go deeper

for a junior

Be ready to describe the column=value directory layout and say that a filter on that column lets the engine skip whole directories without reading them.

for a middle

Explain that the partition value is usually absent from the data files and reconstructed from the path, and that pruning happens before I/O rather than as a row filter afterwards.

for a senior

An interviewer expects you to connect partition choice to file sizing and planning cost, and to name the classic operational trap of unregistered partitions being invisible to queries.

for a principal

Own the argument for why partition layout became table metadata rather than a directory convention: it is what makes layout a decision you can revise instead of a table you must rebuild.

## The problem partitioning solves A table on object storage is ultimately a set of immutable data files — Parquet, ORC or Avro — sitting under some prefix. Without any organising scheme, a query such as `SELECT count(*) FROM events WHERE dt = '2024-05-01'` has no way to know which files could possibly contain that day, so it must open all of them. On a table with years of history that is an enormous amount of I/O to answer a question about one day. Partitioning is the coarsest and oldest fix: physically group the files so that the value of a chosen column determines *where* a file lives. ## The Hive-style convention The convention that came out of Apache Hive, and which nearly every engine still understands, encodes the partition column and its value directly in the directory name: ``` s3://warehouse/events/ dt=2024-05-01/country=DE/part-0000.parquet dt=2024-05-01/country=US/part-0000.parquet dt=2024-05-02/country=DE/part-0000.parquet ``` Two properties matter: 1. **The partition column's value lives in the path, not in the file.** Writers usually strip it from the data itself — every row under `dt=2024-05-01` has the same value, so storing it a million times is waste. Readers reconstruct the column by parsing the directory name. 2. **The layout is a physical commitment.** The directory tree *is* the partitioning. Changing it means moving or rewriting files. ## How pruning actually happens Given `WHERE dt = '2024-05-01'`, the engine evaluates that predicate against partition values first, selects the surviving directories, and only then lists and reads files inside them. This is fundamentally cheaper than a row filter: a row filter still costs you the read. Two consequences follow directly: - A filter on a **non**-partition column prunes nothing at the partition level. It may still be helped by file-level statistics, but that is a different, finer mechanism. - Once a partition is selected, partitioning does nothing more. All of that partition's data is a candidate. ## Who knows which partitions exist In the classic Hive setup, a metastore holds the list of partitions. A directory that appears in storage but was never registered is invisible to queries until someone runs `ALTER TABLE ... ADD PARTITION` or a repair command. This is a frequent source of "the data is there but the query returns nothing". Modern table formats change this: the table's own metadata records, per data file, which partition it belongs to. Pruning becomes a filter over that metadata rather than a directory listing, no registration step is needed, and the on-disk path layout becomes a convention rather than the source of truth. ## Choosing partition columns Good partition columns are: **low cardinality** relative to the table, **present in most queries' filters**, and **evenly distributed**. Time (day, sometimes month or hour) is the archetype, because analytical queries almost always bound a time range and because old data ages out as whole partitions. The sizing rule of thumb is the one that catches people out: each partition should hold enough data to justify at least one full-sized data file. Partitioning by something with millions of distinct values produces millions of directories with a few kilobytes each, and the cost of enumerating and planning over them swamps any pruning benefit. ## What partitioning is not - It is **not an index**. Nothing is built or maintained; the layout *is* the mechanism. - It is **not sorting**. Rows within a partition are in whatever order the writer produced. - It does **not** speed up filters on other columns. - It does **not** reduce the data read once a partition is selected. ## The two failure modes to remember First, filtering on a column the partition value was *derived* from (filtering a raw timestamp when the table is partitioned by a date string) prunes nothing, because the engine does not know the two are related. Second, over-partitioning trades a scan problem for a planning-and-tiny-files problem. Both are the reason table formats moved partition tracking into metadata and made the partition layout something you can change later.

  • If the partition value is not stored in the data files, how does a SELECT still return that column?
    The reader reconstructs it. The engine parses the directory name into the declared partition column and its type, then materialises a constant for every row it reads from files under that path. That is also why a type mismatch between the declared column and the path text produces surprising comparison behaviour.
  • Why can a directory full of correct data return zero rows in a Hive-style table?
    Because the metastore, not storage, is the list of partitions. A writer that dropped files into a new path without registering the partition leaves it invisible until an ADD PARTITION or repair command runs. Formats that track files in table metadata do not have this failure mode, since a commit registers the files.
  • Does partitioning help a query that filters on a non-partition column at all?
    Not at the partition level — every partition remains a candidate. Such a query can still be narrowed by finer mechanisms, notably per-file minimum and maximum statistics that let the reader skip individual files, but that is independent of the directory layout and is much weaker than pruning whole partitions.

It is filing paperwork in one drawer per day rather than one giant pile: if you know the date you only open one drawer, but if you file one drawer per customer you end up with thousands of drawers holding a single sheet each.

saying these in an interview costs you the question

  • Says partitioning is an index the engine builds and maintains
  • Thinks the partition column is stored inside every data file
  • Claims more partitions always means faster queries
  • Expects a filter on any column to prune partitions
  • Believes partitioning also sorts rows within each partition

context

open as a page

Why does a table partitioned by dt scan every partition when a query filters only on event_ts?

level: middleimportance: must knowfreq 66%

basics

~20 s

Pruning only works on the partition column itself. A predicate on event_ts says nothing about dt, so the engine lists every partition and discards rows afterwards. Filter dt too, or use a format that derives dt from event_ts.

open as a page

How can a table format change a table's partitioning without rewriting existing data?

level: seniorimportance: must knowfreq 58%

basics

~20 s

The table records which partition layout each data file was written under. Changing the layout is a metadata commit that applies to new writes only; existing files keep their old partition values, so nothing is rewritten and no downtime is needed.

open as a page

Why partition a lakehouse table by a hash bucket of user_id rather than by user_id itself?

level: middleimportance: should knowfreq 45%

basics

~10 s

Hash bucketing maps a high-cardinality key onto a fixed number of partitions. Equality lookups on user_id still skip almost everything, while the table keeps a bounded number of partitions holding properly sized files.

open as a page

A lakehouse table partitioned by hour and customer_id has millions of partitions — what breaks?

level: seniorimportance: should knowfreq 50%

basics

~20 s

Cost moves from scanning to planning: the engine must enumerate and filter millions of partition entries before reading anything, files fall far below target size, and skew grows. Coarsen the time grain and hash the customer key into buckets.

open as a page

After changing a table's partition spec, when do you rewrite historical data under the new layout?

level: principalimportance: nice to knowfreq 30%

basics

~20 s

Rewrite history only when old data is queried often enough that the mixed layout measurably hurts, and the rewrite fits your compute budget. Otherwise let the new layout govern new writes and let retention retire the old data.

open as a page