skip to content

When a grouped step must have its group's values at once, is every column of those rows held or only the ones read?

level: middleimportance: should knowfreq 38%

answer

  1. positions recorded, not rows copied
  2. what forces the materialisation
  3. one column's values or the whole row
  4. rows in the group times bytes materialised

basics

~20 s

It depends on what the surface hands the step. Some pass only the values of the column being reduced; some materialise a table of the group's rows across every column, including ones the step never reads — and then the group's width is part of the bill.

solid answer

~50 s

Forming the groups is bookkeeping: recording which row positions carry which key value. That much copies nothing. A copy happens only when something demands the group as a materialised object, and what that object contains differs by surface. A reduction the tool implements over one column generally needs only that column's values for the group. A **hand-written per-group body** — a function you supply that the library cannot look inside — is often handed a table of that group's rows, every column of it, whether or not the body touches them. So the room for the biggest group is its row count multiplied by the bytes of what was materialised, and that second factor is the one you control: narrowing the columns before the split, or reducing one column at a time, changes the bill without changing the answer.

go deeper

for a junior

Know that grouping rows by a key does not by itself duplicate them; something later in the step has to ask for the group before a copy is made.

for a middle

Be able to multiply the two factors — the largest group's row count and the bytes of whatever was materialised — and say which of the two the code controls.

for a senior

Recognise the wide-table trap in review: a step expressed so that it receives the group whole drags every unread column along, and narrowing before the split is the cheap repair.

for a principal

The policy question is whether steps that receive a whole group are allowed at all in jobs whose input width is not under your control, and what review catches them.

## Bookkeeping is not copying In a **grouped operation** — split, apply, combine: one pass that gives every row a key, runs a computation once per key, and reassembles the answers — the split does not, by itself, build one table per group. What it has to establish is which rows carry which key value, and that can be recorded as positions against the original storage. A tool can then run a reduction straight over the rows it already holds, touching each once. The useful question is therefore not "does grouping copy?" but **what forces a copy** — and the answer is: something downstream that demands the group as an object it can be handed. ## What a held group is handed, and why it varies When the per-key work does need its group present, the amount of memory that costs depends entirely on the form the group arrives in. Three forms are common across this family of tools, and they differ by an order of magnitude on a wide table: | Form the group arrives in | What is resident for the biggest group | Typical cause | |---|---|---| | the values of one column for that group | rows in the group x the width of one value | a reduction the tool implements over a named column | | the values of several named columns | rows in the group x the width of those columns | a computation declared over a few columns at once | | a table of that group's rows, all columns | rows in the group x the full row width | a body handed the group as a whole table | The third row is where people are surprised. A table forty columns wide, reduced on one numeric column, can cost forty columns' worth of the largest group if the step was expressed as a body that receives the group whole. Nothing about the answer changed; only what was handed over did. ## What you can change The two factors multiply, and both are yours to influence: - **Rows in the largest group** — usually a property of the data, not the code, though a coarser or finer **grouping key** changes it. - **Bytes per row that were materialised** — squarely a property of how the step was written. Narrowing the table to the columns the computation reads *before* the split removes the unread columns from the bill entirely. Expressing the work as a reduction over a named column, rather than as a body that receives everything, does the same thing. There is a third, quieter factor: the representation of the columns that did come along. A column of long text values carried into a materialised group costs far more per row than a fixed-width numeric one, so a wide table's bill is rarely evenly spread across its columns. ## Do not generalise either way Two flat statements are both wrong, and each is true of some design in this family: - *"Grouping copies each group's rows into a table of its own."* True at the moment something demands a materialised group, and of designs that build sub-tables eagerly; false of designs that record positions and run the reduction over the original buffers. - *"A group is only ever handed the column being reduced."* True of the implemented reductions on most surfaces; false wherever a body receives the group as a table. The statement that survives both is the one about what forces the materialisation — and that is the thing to check in the code in front of you rather than assume from habit. ## Diagnosing it If a grouped step is consuming far more memory than the column it appears to be reducing could account for, work through this order: 1. Identify what the per-key step receives — one column's values, or the group's rows. 2. Multiply the largest group's row count by the width of what it receives; compare against what you observed. 3. If the second factor is the whole row, narrow the columns ahead of the split and measure again. 4. Only then look at whether the computation itself could be expressed in a folding form, which removes the residency question rather than shrinking it. ## What an interviewer is checking This question separates a candidate who knows that some grouped steps are expensive from one who can say *which bytes* are on the bill. The second candidate reaches for the cheap fix — drop the columns nobody reads before the split — and the first reaches for a bigger machine.

  • Why does narrowing the columns before the split help more than narrowing them after?
    After the split the expensive moment has already passed: whatever was materialised for the largest group was resident while the step ran. Removing columns earlier means they are never part of what is handed over, so the peak itself is smaller rather than the result being tidier.
  • Does the columns' representation matter, or only how many there are?
    Both. A fixed-width numeric column costs the same per row every time; a column of variable-length text can cost many times that, and one such column can dominate a materialised group. Counting columns is a rough proxy; counting bytes per row is the real figure.

saying these in an interview costs you the question

  • Assumes forming groups always builds a table per group.
  • Assumes only the reduced column is ever held.
  • Ignores the unread columns dragged into a materialised group.
  • Reaches for a bigger machine before narrowing the columns.
  • Thinks column representation makes no difference to what is held.