What does Hive-style directory partitioning do to a table's file layout on object storage?
answer
- the value lives in the path, not the file
- folders named column=value
- the filter picks folders before any read
- pruning removes I/O instead of filtering rows
basics
~10 sHive-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 sInstead 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 liness3://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.parquetgo deeper
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.
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.
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.
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