skip to content

In a pipeline that filters rows, derives a column, counts distinct users per country and orders the result, which steps regroup records?

level: middleimportance: must knowfreq 70%

answer

  1. look for a key, not for complexity
  2. predicate and derivation stay local
  3. distinct, group, join, total order
  4. slow is not the same as wide
  5. concatenation is narrow; dedup is not

basics

~20 s

The filter and the derived column are narrow: each output depends only on the record in hand. The per-country distinct count is wide, and so is the global ordering — their outputs depend on records that may sit in any piece.

solid answer

~50 s

Two of the four steps regroup. The predicate and the derived column are **narrow**: a worker thread decides each record from the record itself, inside the piece of the input it already holds. `count(distinct user_id) per country` is **wide** twice over in principle — identifying one user's records as one user, and bringing one country's records together — although a runtime may satisfy both with a single redistribution. The final global ordering is **wide** as well: you cannot know what precedes what without seeing every record. The giveaway to train yourself on is a key: group by it, join on it, deduplicate over it, rank within it, or order across all of it, and records have to be re-placed. No key mentioned and one record at a time means narrow, however slow the per-record work is.

code

sql · 7 lines
sql
select
  lower(country_code) as country,
  count(distinct user_id) as users
from events
where event_date = date '2026-09-01'
group by lower(country_code)
order by users desc

go deeper

for a junior

Learn the giveaway: if a step names a key — grouping, joining, deduplicating, ranking, total ordering — records have to be brought together. If it touches one record at a time and names no key, it is narrow.

for a middle

Walk a real pipeline line by line and label each step before measuring anything. Be able to explain why a distinct count is wide despite a tiny output, and why a slow remote call per record is narrow despite dominating the clock.

for a senior

Separate the two repair families out loud: slow narrow work wants batching, concurrency or caching; a wide step wants a rewrite that removes or shares the regrouping. Say which runtimes can collapse two regroupings into one and which leave that to the author.

for a principal

Make labelling shape part of pipeline review. When every author can name the wide steps in their own code, capacity conversations stop being about rows processed and start being about how many times the data is re-placed.

## Reading a pipeline for its regroupings A **piece of the input** is one slice of the stored input that a single **worker thread** — one lane inside a worker process — reads from start to finish. A step is **narrow** when a thread can finish it from the piece already in its hands, and **wide** when producing one output needs records currently sitting on other workers. Reading your own pipeline for wide steps is the cheapest performance work available, and it is done by eye. Take the four steps in order. 1. **The predicate.** Each record is kept or dropped on its own contents. Narrow. It removes records, which shrinks everything downstream, but it does not change any later step's shape. 2. **The derived column.** A parse, a cast, a lower-casing, a lookup inside the record. Narrow, and narrow no matter how expensive the derivation is. 3. **The distinct count per country.** Two dependencies at once: the records of a single user may be scattered across pieces and must be recognised as one user, and the users of a single country must be counted in one place. Wide. 4. **The global ordering.** No record's position is knowable from the piece holding it. Wide. ## The giveaway is a key | what the step mentions | shape | why | |---|---|---| | nothing but the record | narrow | the inputs of one output are already present | | a grouping key | wide | records carrying the key may be anywhere | | a deduplication column | wide | equal values may be anywhere | | a join key | wide | matching records come from two places | | a total ordering or a global rank | wide | position depends on every other record | | an explicit re-division | wide | the re-placement *is* the step | | set intersection or difference | wide | membership is decided against the whole other side | | concatenating two inputs | narrow | no record has to meet another | That last row surprises people. Appending one dataset to another needs no record to meet any other, so the pieces of both inputs simply stand side by side; only a *deduplicating* union is wide, and it is wide because of the deduplication. ## Three traps in your own code - **Slow is not wide.** A per-record call to a remote service can dominate the wall clock while remaining narrow. Its remedies are batching and bounded concurrency, not anything to do with re-placing records — a completely different repair from the one a wide step wants. - **Small output is not narrow.** A distinct count emits one row per country and may still move a great deal of data to get there, because the shape of the dependency, not the size of the answer, decides whether records travel. - **A key you did not write is still a key.** Ranking within a group, a windowed running total over a partitioning column, a semi-join used as a filter and a de-duplicating union all name a key even when the word never appears in your code. ## What differs between runtimes, and what does not The classification above is a property of the computation, so it holds wherever you write it. What varies is the treatment: - A declarative surface may **collapse two regroupings into one** when both are keyed the same way, so the distinct-then-count pair may cost a single redistribution rather than two. A pipeline written as opaque per-record functions gets no such collapse, and the oldest execution model in this family — the two-phase disk-to-disk one — gives you exactly one redistribution per job and makes you fold your logic into it. - On a runtime that schedules work in units between redistributions, each wide step splits the program into separately scheduled units. On a record-at-a-time runtime everything runs at once and the wide step is a routing decision on every record. - Whether the ordering at the end is even attempted as a total order varies with the runtime and the output contract; what is universal is that a total order cannot be produced from one piece alone. ## The habit to build Read the pipeline once and write the shape beside each line before touching any measurement. A pipeline of twelve steps where two are wide has two places where its cost can live, and knowing which two lines those are turns an open-ended tuning exercise into a bounded one. If a wide step turns out to be avoidable — because the records it needs are already together, or because the operation can be reordered to share an earlier regrouping — that is the rewrite worth making; if it is not avoidable, at least you know what the run is paying for.

  • Is appending one dataset to another narrow or wide?
    Narrow, if duplicates are kept: no record has to meet another, so the pieces of both inputs simply stand side by side and the combined run has the pieces of both. A union that removes duplicates is wide, and it is wide because of the deduplication — equal records may sit anywhere — not because two inputs are involved.
  • A step calls a slow remote service once per record. Does that change its shape?
    No. Each output still depends only on the record in hand, so the step is narrow and stays narrow however long the call takes. That matters because the repairs differ: a slow narrow step wants batching, bounded concurrency or caching, while a wide step wants a rewrite that removes or shares the regrouping.
  • Why might the distinct count and the grouping cost one regrouping rather than two?
    Because both are keyed on values that can be redistributed together: once records are placed by country, deduplicating users within a country is local work. Whether that collapse happens depends on the surface — a declarative plan can spot it, a pipeline of opaque per-record functions cannot, and the two-phase disk-to-disk model expects you to fold it in by hand.

saying these in an interview costs you the question

  • Calls a step wide because it is slow or because it calls a remote service.
  • Assumes a small result means no records were regrouped.
  • Misses that a rank or running total within a group names a key.
  • Says appending two datasets requires bringing them into one division.
  • Treats the filter as the expensive step because it reads every record.