skip to content

Grouping 50 million rows by a key returns immediately — what has actually been built at that point?

level: middleimportance: must knowfreq 55%

answer

  1. a fast return means little happened
  2. positions, not payload
  3. the copy is demanded, not automatic
  4. a whole-group argument forces materialisation

basics

~20 s

Normally just bookkeeping: a record of which row positions carry which key value. No per-key table has been built and no computation has run. Rows are copied only when something later demands a materialised group.

solid answer

~50 s

Forming the groups normally produces a **grouped handle** — what a grouping call hands back before any per-group computation has run. It records which rows belong to which key and nothing else, which is why it returns fast on a large table: the key expression has been evaluated and positions recorded, but no per-key table exists and no reduction has run. Rows get copied when something demands a materialised group, most often a hand-written body whose argument is a whole sub-table; a reduction the library implements itself can usually fold straight over the original storage, gathering by position. Designs also differ in how long the handle lives — in some, grouping is a state the table carries until dropped; in others it is a one-shot handle consumed by the next call — so check which model you are in before chaining a second step onto it.

code

pseudocode · 16 lines
pseudocode
# forming the groups: bookkeeping only
positions = empty map from key value to list of row positions
for i in 0 .. rowCount - 1:
    k = evaluate(keyExpression, row i)
    append i to positions[k]

# a reduction the library implements itself:
# fold the original storage at those positions, build nothing per key
for each (k, idx) in positions:
    emit k, foldOver(column, idx)

# a body the library cannot look inside, handed a whole group:
# the rows at idx must be gathered into a table first
for each (k, idx) in positions:
    groupTable = copyRows(table, idx)   # <- the copy happens HERE
    emit k, myBody(groupTable)

go deeper

for a junior

Know that a grouping call returning instantly on a huge table has not computed anything: it has only worked out which rows share each key value.

for a middle

Explain the difference between recording positions and copying rows, and name what forces a copy — a supplied body whose argument is a whole materialised group.

for a senior

When a grouped step allocates more than expected, inspect what the apply is handed before you blame data volume, and know whether your tool keeps grouping as table state.

for a principal

Decide as a team whether supplied per-group bodies are acceptable in shared pipelines, since they are what turns a streaming grouped step into one that materialises.

## What a grouping call hands back A **grouped handle** is what a grouping call returns before any per-group computation has run: it records which rows belong to which key value and nothing else. That is the whole reason a grouping call over fifty million rows can return in a moment. The expensive things — reading each group's values, folding them, assembling a result — have not happened. Concretely the split does two things: 1. Evaluate the **grouping key** — the column, columns or expression whose value decides which rows belong together — once for every row. 2. Organise those values so the rows carrying each distinct value can be found again, by recording row **positions** against each key value. ## Bookkeeping against copying The distinction worth carrying out of this question is between *forming* the groups and *materialising* them. - **Forming the groups** is bookkeeping: a mapping from key value to a set of row positions. It holds no copy of the data. - **Materialising a group** means gathering that key's rows into a table of its own. That is real work and real memory, and it happens only when something asks for it. Some designs do materialise a sub-table per group, and any design materialises at the moment something demands a whole table as an argument. So the useful question is never "does grouping copy?" but "what in my step forces the copy?" ## What forces the copy | what the apply is | what it needs | does it materialise | | --- | --- | --- | | a reduction the library implements (sum, count, minimum) | the column's values at the recorded positions | normally not | | a body handed the group's column | that one column, gathered | one column at most | | a body handed the whole group as a table | every column of the group, side by side | yes | | anything returning rows rather than a value | somewhere to assemble the rows | yes, on the output side | The practical consequence is that a **hand-written per-group body** — a function you supply that the library cannot look inside — is the usual reason a grouped step that should have streamed suddenly allocates. It is not that the body is yours; it is what the surface hands it. ## What the bookkeeping itself costs It is not free. One key evaluation happens per row, and the structure holds one entry per distinct key plus one position per row. Where the number of distinct keys approaches the number of rows, that structure is comparable in size to the column you grouped on. "Grouping returned instantly, so grouping is cheap" is therefore only half true: it returned instantly because the work was deferred, not because the work is small. ## How long the handle lives Designs disagree here, and the disagreement produces a real bug class: - In some, grouping is a **state the table carries** until explicitly removed, so subsequent operations honour it and a step written as though it were over the whole table is silently computed per group. - In others, grouping produces a **transient handle** consumed by the very next call, after which the result is an ordinary table with no memory of the split. Neither is wrong; they are different contracts. The hazard — a summary computed one level less grouped than the author believed — is live only in the first model, so it is worth knowing which one you are in before chaining a second step. ## Sorted runs against hashed buckets There are two common routes to the bookkeeping: - A **sorted split** orders the key column so equal keys form a contiguous run. It pays an ordering cost up front and makes the subsequent reads sequential. - A **hashed split** drops each key value into a bucket by its hash. It avoids the ordering cost, and the reads it feeds are scattered. Both produce the same grouping. They differ in cost, and in whether they leave an ordering behind as a by-product. ## What this changes about how you write the step - Read a fast-returning grouping call as **nothing having happened yet**, never as evidence the operation is cheap. - If the step allocates unexpectedly, look at what the apply is handed before you look at the data volume. - Prefer a reduction the library implements when one exists, because that is the path that can fold over the original storage. - When you do need your own body, prefer a surface that hands it a column over one that hands it a whole table, where the design offers both.

  • If forming the groups does not copy rows, how does a reduction read each group?
    By gathering from the original storage at the positions the split recorded. A sum reads that column's values at those positions and folds them into one accumulator, so nothing per key needs to exist beyond the accumulator and the position list.
  • Is the grouped handle still usable after you have computed one thing from it?
    That varies by design. In some, grouping is a state the table carries until explicitly dropped, and later operations still act per group — with the hazard of a step being grouped when you thought it was not. In others the handle is consumed by the next call and the result is an ordinary table.
  • How do you find out how many groups a split produced without paying for a heavy computation?
    Ask for the number of distinct key values the split recorded, or run the cheapest reduction available over the handle. The bookkeeping already knows how many keys it holds, so that count is nearly free compared with any per-group work.

A cloakroom writes you a numbered ticket; it does not carry your coat into a room with all the other blue coats. Forming the groups is writing the tickets — cheap, and nothing has moved. Only someone who insists on seeing every blue coat laid out together makes anyone actually carry them.

saying these in an interview costs you the question

  • Says grouping immediately builds one sub-table per key value.
  • Assumes a fast return means the computation has already run.
  • Believes a reduction must materialise each group before it can fold it.
  • Treats the grouped handle as a finished result rather than deferred work.
  • Thinks grouping is always a state the table keeps until you drop it.