A write grouped into directories by a 200-value country column runs over 2,000 pieces. Why can it leave 400,000 files?
answer
- two numbers multiply, not one
- pieces times groups
- directories named for a column value
- cardinality of the grouping column
- each piece can touch every directory
basics
~20 sBecause the two numbers multiply. Each worker thread writes what it holds into the directory each record's value names, so a piece holding all 200 values leaves 200 files. Across 2,000 pieces that is 2,000 x 200 files, each tiny.
solid answer
~40 sA grouped write does not reorganise the data; it only decides which directory each record lands in. Every piece therefore emits one file into every directory its own records touch, so the file count is the sum over pieces of the distinct values present in that piece — bounded below by the number of values and above by `pieces x values`. With records for all 200 countries scattered through all 2,000 pieces you sit at the upper bound: 400,000 files, each roughly a two-hundredth of the size it would otherwise be. The fan-out is driven by the grouping column's cardinality, which is why grouping by a per-second timestamp or an identifier produces a directory per record and a file count nobody predicted.
go deeper
Recall that grouping a write by a column creates a directory per value, and that each piece writes its own file into each directory it has records for.
Do the arithmetic out loud: sum over pieces of the distinct values present, floored at the value count and capped at pieces times values, and name cardinality as the risk.
Show judgment about the grouping column itself — which queries the directories actually serve, what cardinality the layout can carry, and which downstream reader inherits the result.
The call is a layout standard on a shared table: which columns may be grouped on, what file-count ceiling a write is allowed to leave, and what it costs to move existing jobs onto that rule.
## Two numbers, multiplied **Column-value directories** are stored files grouped into directories named for one column's value, so that a later filter on that column can drop whole directories before anything is read. It is a different sense of the word "partition" from **a piece of the input**, which is one slice of the stored input that a single worker thread reads and processes from start to finish — and the two senses colliding is what makes this question hard to answer cleanly. The key fact is that the grouped write does not move anything. Each worker thread still writes only the records it is already holding; the grouping just decides which directory each of those records goes into. So a thread holding records for 200 different countries opens 200 outputs and leaves 200 files, all of them small. The file count is therefore: > the sum, over every piece, of the number of distinct values of the grouping column present in that piece with a floor of one file per value and a ceiling of `piece count x value count`. Values spread evenly across pieces put you at the ceiling: 2,000 x 200 = 400,000 files, each about a two-hundredth of the size it would otherwise have been. ## Cardinality is the whole risk - **a handful of values** (a status, a region, a year) is what the mechanism was designed for: a small, stable number of directories that filters can use; - **a few hundred values** is where the multiplication starts to bite, as in the 200-country case; - **a near-unique column** — a timestamp recorded to the second, a user identifier, an order number — gives roughly one directory per record, and the directory listing becomes larger than the data; - **two grouping columns** multiply their cardinalities before the piece count is applied at all: 200 countries x 24 hours is 4,800 directories before you have written a single file. The interview tell is a candidate who reasons about the directory count and stops. The directory count is only one of the two factors. ## Keeping the two senses apart | | a piece of the input | a column-value directory | |---|---|---| | what it is | one slice a worker thread processes end to end | one on-disk directory named for a column's value | | who creates it | the run's division of the input | the write's grouping instruction | | how many | the piece count, set upstream | the cardinality of the grouping column | | effect on the write | one output per piece | one output per piece **per directory it touches** | Saying both names out loud once — "one piece of the input per worker thread, which is not the same thing as one column-value directory on disk" — is usually enough to keep the rest of the answer straight. ## What moves you between the floor and the ceiling 1. If a value's records are scattered through every piece, every piece writes into that value's directory, and you are at the ceiling. 2. If a value's records already sit in one piece, that directory receives one file, and you are at the floor. 3. Getting from the first state to the second means redistributing records so that each value lands in one piece before the write — a redistribution owned by `Moving Data Between Workers`, and it buys the file count at the price of a step where records travel between workers. For this leaf the point is narrower and more durable: the file count of a grouped write is *pieces touching a group*, not *groups*. ## What varies between engines - some runtimes require the records to be ordered by the grouping column before the write, so only one output is open at a time per piece; others hold many outputs open at once, and some cap how many and reopen a directory later — which produces **more** files, not fewer; - some runtimes measure the previous step and merge pieces before writing; others write exactly what they hold; - a continuous job multiplies again by commit intervals: declared width x directories touched x intervals, which is how a grouped continuous write reaches millions of files without anyone choosing a large number. ## Where the aftermath belongs One country's directory being a thousand times the size of every other is uneven data, owned by `The One Heavy Key`. The room a worker needs while several outputs are open is owned by `Memory, Spill and Caching`. Repairing a shredded table afterwards, and the declared partition spec that table carries, belong to `lakehouse table compaction` and `Table Format Concepts`. What the next job inherits from the files is `Layout the Reader Inherits`. This leaf owns only the multiplication that created them.
- What makes a column a bad choice to group a write by?High cardinality. A timestamp to the second or an identifier gives roughly one directory per record, so the listing outgrows the data and later filters gain nothing. Two grouping columns are worse: their cardinalities multiply before the piece count is applied. A grouping column wants few, stable values that queries actually filter on.
- Does the fan-out still happen if each group's records already sit in one piece?No — then each group is written by one worker thread and leaves one file, which is the floor. Getting records into that arrangement means a redistribution before the write, owned by `Moving Data Between Workers`. The point here is only that the file count counts pieces touching a group, not groups.
- Is the fan-out the same in a continuous job?It is worse, because a third factor appears. There is no final piece count; each instance of the declared operator width writes into every directory its records touch, once per commit interval. The file count therefore grows with elapsed time as well as with width and cardinality.
saying these in an interview costs you the question
- Expects one file per group, because the directories are named for the column.
- Groups a write by a near-unique column such as a per-second timestamp.
- Thinks the number of directories alone sets the file count.
- Treats a column-value directory and a piece of the input as the same thing.
- Assumes the grouped write reorganises records before writing them.