A printed plan shows one step handling 2,000 pieces at once and the next handling 1 — what does that tell you?
answer
- pieces at once, not bytes
- widths come in blocks between crossings
- a width of one idles the cluster
- one thread, one worker's memory
- printed width is not the ran width everywhere
basics
~20 sThat the work has been funnelled through a single worker thread. Step width is how many pieces of the input a step processes at the same time, so a width of one means the cluster's size stops mattering from that step onward.
solid answer
~50 s**Step width** is how many pieces of the input a step processes at the same time — the plan's statement of how much of the cluster is working during that step. A drop from 2,000 to 1 means everything now passes through one worker thread, so that step is bounded by one thread's throughput and one worker's memory no matter how large the cluster is. The usual causes are an aggregate over the whole input with no grouping key, a globally ordered result assembled in one place, a write forced into a single output, a grouping whose key has one distinct value, or the result being pulled back to the coordinating process that assembled the graph. Note what the width does *not* say: it counts pieces, not records or bytes, so it cannot tell you whether those 2,000 pieces are evenly filled.
go deeper
Know that a step's width is how many pieces of the input it processes at the same time, and that a width of one means a single worker thread is doing all the remaining work.
Explain why widths change only across a point where records move between workers, name the usual causes of a collapse to one, and say that the width counts pieces rather than records or bytes.
From a collapse, identify which construct in the job produced it and whether the volume passing through the narrow point justifies concern. Say which of your statements about widths hold only for a job that runs once rather than continuously.
Weigh whether narrow tail steps are acceptable platform-wide given what the downstream consumers actually require, and how much a fixed-width continuous runtime constrains capacity planning across a fleet of jobs.
## What a width is **Step width** is how many pieces of the input a given step processes at the same time. It is the plan's answer to "how much of the cluster is busy here", and it is printed per step, not once per job, because it changes down the listing. It changes for a structural reason. Within a run of steps that needs no record movement, the width cannot change: each output piece is built from one input piece, so the count is inherited from whatever produced it. Across a point where records are sent between workers, the width is whatever that redistribution produced. So the widths in a plan come in blocks, one block per run between two movement points, and the interesting lines are the ones where a block's width differs sharply from its neighbour's. ## Reading a collapse to one A step of width 1 beneath a step of width 2,000 is the collapse worth recognising on sight. It says the entire result now passes through a single worker thread. From that step on: - the job's throughput is one thread's throughput, and adding machines changes nothing; - the whole of the data at that point must fit within one worker's working memory, or it spills to local disk — writing part of the working set out because it does not fit — which is slower again; - every other worker in the cluster is idle while it runs, which is exactly the shape of "why is this job slow when the cluster is mostly doing nothing". The common causes, all visible from the step just above the collapse: 1. **An aggregate over the whole input with no grouping key** — a single total has one destination by definition. 2. **A globally ordered result assembled in one place**, rather than by giving each worker a disjoint range of values. 3. **A write forced into a single output**, where the job was asked for one file or one destination handle. 4. **A grouping whose key has one dominant or single distinct value** — the extreme of *skew*, where one piece holds far more records than the rest. 5. **The result being pulled back to the coordinating process** — the single process that assembled the graph and receives whatever the demanding call asked for. Not every collapse is a defect. A final aggregate that emits one row, or a deliberately single output file that a downstream consumer requires, is a narrow step by design; what matters is whether the *volume* passing through the narrow point is small. ## What a width does not tell you This is the half of the question that separates a middle answer from a junior one. - It counts **pieces, not records or bytes**. Two thousand pieces can be two thousand equal pieces or one piece holding most of the data with 1,999 nearly empty ones, and the width line reads identically in both cases. - It says how many pieces are processed **at the same time**, not how many there are in total. Where the count of pieces exceeds the worker threads available, the rest are queued and the step runs in successive rounds. - It is not a statement about **duration**. A wide step over huge pieces can dominate a job whose narrow steps are instant. ## Where the models disagree about widths | Engine model | How the width behaves | |---|---| | Finite job, one run per submission | Widths are derived per run: from how the input was cut for the read, and from the redistribution's piece count after each movement point | | One fixed graph kept running, record at a time | The width is usually fixed for the life of the job; changing it typically means stopping and restarting, redistributing any retained state to the new number of workers | | Continuous work as a succession of small finite jobs | Each small job re-runs the same graph at the same widths, so a bad width is paid once per small job, repeatedly | On top of that, some engines can revise the not-yet-run part of a plan using statistics measured from work already finished, which can change a width after the text was first printed. That mechanism — replanning after measurement — is a subject of its own with its own owner; the point for plan reading is only that on such an engine the printed widths may not be the widths that ran, so you have to know which text you are looking at. ## The line this reading stops at Reading a width off the plan and interpreting a collapse is plan reading. **Choosing** the number — deciding how many pieces the input should be cut into, arguing about the arithmetic of bytes against threads — is a different subject and belongs with how the work is split in the first place. In an interview, say what the width shows and what the collapse costs; if the follow-up asks what number it should be, you have crossed into that other subject and should say so.
- Why can a width not change inside a run of steps with no movement between workers?Because each step there builds one output piece from one input piece. With no redistribution, there is nothing to renumber the pieces, so the count is inherited from the step that produced them and carries unchanged to the end of the run.
- The plan shows width 2,000 but the cluster offers 200 worker threads. What is actually happening?The 2,000 pieces are processed in successive rounds of roughly 200 at a time. The width is a property of the work, not of the machines; how many run concurrently is bounded by the threads available, which is why a width far above the thread count mainly buys finer scheduling, not more parallelism.
- Is a final step of width one always a problem?No. A single total, or a deliberately single output that a consumer requires, is narrow by design, and if only a few rows pass through it the cost is negligible. It is a problem when the volume crossing the narrow point is large, because one thread and one worker's memory then bound the whole job.
saying these in an interview costs you the question
- Reads a width as how much data each piece holds
- Thinks a width of one is always the engine misbehaving
- Believes widths can change inside a run with no movement
- Assumes a width is a number the author edits in the plan
- Says printed widths are always the widths that ran
- Treats width above the thread count as extra parallelism