skip to content

Narrow and Wide Steps

A step a worker finishes from the piece it already holds, against one that needs records from every other worker: telling them apart in your own code, and the rewrites that avoid the second.

on this pageshow

questions

4

A worker thread holds one piece of the input. What decides whether it can finish a step alone or needs records other workers hold?

level: juniorimportance: must knowfreq 78%

answer

  1. ask what one output depends on
  2. own piece, or every piece
  3. filter and derive against group and order
  4. shape, not line count
  5. tiny output can still be wide

basics

~20 s

Dependency shape decides: if each output depends only on records already in the thread's piece, it finishes alone — a narrow step. If one output needs records spread across every piece, the step is wide.

solid answer

~50 s

Ask one question about the step: to produce a single output record, which input records must be read? If the answer is "only records already sitting in the piece this thread holds", the step is **narrow** — the thread runs it to completion with no communication, so per-record filters, derived columns and one-to-many expansions all qualify. If the answer is "records that could be in any piece, on any machine", the step is **wide**: grouping by a key, dropping duplicates, joining on a key, ordering the whole dataset. Wide steps need the matching records brought together first, which is the redistribution usually called a shuffle. Note what the test ignores: how many lines the step is, how slow the per-record work is, and how small the output is. A global count is wide even though it emits one row.

go deeper

for a junior

Be able to state the test out loud: does one output need only the records in this thread's piece, or records that could be anywhere? Then sort the everyday operations — filter and derive on one side, group, distinct, join and global order on the other.

for a middle

Explain why the two categories differ in kind rather than degree: the wide one cannot produce correct output until records placed by the stored layout have been re-placed by your key. Show that a tiny output can still be wide.

for a senior

Use the test as the first pass on a slow run, before any measurement: name the key in each line and count the wide steps. Say plainly what varies between runtimes — chaining narrow work, whether a wide step is a scheduling split, and whether anything is rewritten at all.

for a principal

The point worth institutionalising is that cost tracks dependency shape, not code volume. Review standards, cost estimates and job templates that reason in rows processed will misprice work systematically; the count of wide steps is the cheaper predictor.

## The vocabulary this turns on A distributed run cuts its input into **pieces of the input** — one slice of the stored input that a single **worker thread** reads and processes from start to finish. A worker thread is one lane inside a **worker process**: one operating-system process on one machine, holding its own memory and several such lanes. A piece occupies one lane for its whole life, so the number of pieces is the ceiling on how much of the cluster can be busy at once. Given that arrangement, every step in your program is one of two things, and the difference is not a matter of degree. ## The test, in one question **To produce one output record, which input records must be read?** - *Only records already inside the piece this thread holds* → the step is **narrow**. The thread finishes it alone. Nothing crosses the network on its account. - *Records that could be sitting in any piece, on any machine* → the step is **wide**. Before the output for one key can exist, the records carrying that key have to be brought together. That bringing-together is a redistribution of records between workers; the term of art is a **shuffle**, meaning the step that sends each record from the worker holding it to the worker that will produce the output for its key. Nothing in the test mentions how much code the step contains, how expensive the per-record work is, or how large the output is. It is purely a statement about *where the inputs of one output live*. ## The everyday operations, sorted | step | one output depends on | shape | |---|---|---| | keep records matching a predicate | the record in hand | narrow | | add a derived or parsed column | the record in hand | narrow | | expand one record into several | the record in hand | narrow | | count or sum per key | every record carrying that key, wherever it sits | wide | | drop duplicate values of a column | every record with that value, wherever it sits | wide | | join two inputs on a key | records with that key from both inputs | wide | | order the whole dataset | every record, to know what precedes what | wide | | count all records | one partial result from every piece | wide, though tiny | The last row is the one people argue with. A global count emits a single number, so it feels free — but the single number depends on records in every piece, so per-piece partial results still have to meet somewhere. Its *shape* is wide; its *volume* is negligible. Shape and volume are independent axes, and only the shape is this question. ## Why the shape is what costs you Within a piece, a narrow step is a loop over records that are already in memory or already streaming off local storage. A wide step cannot start producing correct output until records that were placed by something else — usually the byte ranges of the stored files — have been re-placed by the key you asked about. That re-placement involves the network and, on most runtimes, some form of handover between producers and consumers. What that handover actually does, what it costs in bytes, and who waits for whom are subjects of their own; for reading your own code, the useful fact is simply that the categories differ in kind. Two corollaries that interviewers probe: 1. **Line count predicts nothing.** Fifty lines of per-record cleanup is fifty narrow steps. One line that groups by a key is one wide step, and it will usually dominate the run. 2. **Adding capacity does not change a shape.** More worker threads run more pieces at once; they do not turn a step whose output depends on records held elsewhere into one whose output does not. ## What genuinely varies between runtimes The narrow/wide distinction is common to this whole class of engine. What each does with it is not: - **Chaining the narrow work.** Some runtimes fuse a run of narrow steps into a single pass over the piece. Others are record-at-a-time and simply hand each record to the next operator as it is produced. The older two-phase disk-to-disk model chains narrow work inside one half of the job. In all three the practical effect is the same — the records are not re-placed — but the execution is not. - **What a wide step *is* to the scheduler.** On a runtime that schedules the work in units between redistributions, a wide step ends one unit and begins the next. On a runtime where every operator is live at once, it is a routing rule on an edge rather than a scheduling split. - **Whether anything is rewritten.** A declarative surface may reorder or eliminate work before running it; a program written as opaque per-record functions is largely executed as written, and the oldest model in this family rewrites nothing. ## Using the test on your own code Read each line and name the key. If a line mentions a key — group by it, join on it, distinct over it, rank within it — assume wide until you can show the records for that key are already together. If a line mentions no key and touches one record at a time, it is narrow no matter how slow it is. That reading takes a minute and predicts which line the run will spend its time in.

  • A step produces exactly one number for the whole dataset. Is it narrow?
    No. The one number depends on records in every piece, so each piece produces a partial result and those partials still have to meet on one worker. The volume that travels is trivial, but the dependency shape is the wide one, and that is what decides whether communication happens at all.
  • Two steps read the same number of records and one takes twenty times longer. What do you check first?
    Whether either is wide. Before looking at per-record cost, data types or machine health, name the key each step depends on: a step whose output for one key needs records held on other workers is in a different cost class from one that does not, and that difference usually explains a twentyfold gap on its own.
  • Does calling an expensive function on every record make a step wide?
    No. Expense per record and dependency shape are unrelated axes. A costly parse, a decompression or a remote lookup per record is still narrow, because each output depends only on the record in hand; it is slow, and the remedies for slow narrow work are different from the remedies for a redistribution.

Six people sit round a table, each holding a shuffled handful from one deck. "Turn your own red cards face up" is done by each of them alone, in parallel, without a word — narrow. "Tell me how many cards there are of each suit" cannot be answered by anyone from the handful they hold; the cards have to be passed around the table first, and only then counted — wide. The second request is shorter to say and far slower to satisfy.

saying these in an interview costs you the question

  • Says a step is expensive because it processes more rows, never mentioning dependencies.
  • Calls any step containing a user function wide.
  • Thinks adding worker threads turns a wide step into a narrow one.
  • Assumes a step producing a single number needs no records from other workers.
  • Believes every runtime rewrites wide steps away automatically.
  • Equates how much code a step contains with how much it costs.
open as a page

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%

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.

open as a page

A job groups records by account, then joins on account, then aggregates by account again. How can it be rewritten to regroup once?

level: middleimportance: should knowfreq 56%

basics

~20 s

Regroup once on the account key, then keep that division. After the first wide step every account's records already sit together, so later steps keyed on the same account finish locally — unless something in between re-divides or hides the key.

open as a page

Your program contains three wide steps. How does that number read differently on a runtime that schedules in units versus one where all operators run at once?

level: seniorimportance: should knowfreq 42%

basics

~20 s

On a runtime that schedules the work between regroupings, three wide steps mean four separately scheduled units of steps. On a runtime where every operator runs at once, they mean three routing hops each record may make in flight.

open as a page