skip to content

Which facts in the text an engine prints about a job exist only because the work is spread over machines?

level: juniorimportance: must knowfreq 58%

answer

  1. half the plan is about machines
  2. crossings, phases, widths, pruned read
  3. count network hops, not steps
  4. the read's column and condition list
  5. intent, never measurement

basics

~20 s

Four: the points where records must be sent between workers, where the listing breaks into runs of steps between those points, how many pieces each step processes at once, and what the read was told to skip.

solid answer

~50 s

Most engines can print the sequence of steps they intend to run — the step graph, the ordered set of steps derived from your program, each naming the steps whose output it reads. Half of that text would look the same on one machine: a read, a condition, a grouping, a write. The other half is distributed-only, and it is four things. First, the points where records must be sent between workers because the next step needs records another worker holds — what the industry calls a shuffle. Second, the way those points cut the listing into runs of steps that need no record movement. Third, each step's width: how many pieces of the input it handles at the same time. Fourth, what the read was instructed to produce — which columns, and which conditions were pushed down to it. Those four are what you actually read a distributed job's plan for.

go deeper

for a junior

Be able to name the four distributed facts: where records cross the network, where the listing breaks into runs between those points, how many pieces each step handles at once, and what the read was told to produce.

for a middle

Explain why each of the four exists only in a distributed setting, and which program constructs put a crossing there. Say plainly that the text states intent and carries no measurement of what will actually flow.

for a senior

Use the four as a checklist against the job you wrote: expected crossings against key changes, widths against cluster size, the read's instruction against the columns you use. Know which of these facts this engine's text does not carry.

for a principal

Judge how much a team should invest in reading printed text at all when the engines in use disagree about what they print, and decide where the platform should instead capture measurements that the text can never contain.

## What the printed text is An engine of this class does not run your program as written. It derives a **step graph** from it — the ordered set of steps the engine will run, each step naming the steps whose output it reads — and most such engines can render that graph as text. That rendering is the **printed plan**: the engine's own description of what it intends to run. About half of what it contains is ordinary query-plan material that would look much the same for the same computation on a single machine: a read at the bottom, some conditions, a grouping, a join, a write. Reading that half — walking the tree, comparing what an operator was predicted to produce against what it produced, judging per-operator cost — is a general skill that belongs with database execution plans and is not what this subject is about. The other half exists **only because the work is spread over many machines**, and it is the half a data-engineering interviewer is pointing at. ## The four distributed facts 1. **Network crossings.** A point in the plan where records must be sent between workers, because the next step needs records some other worker is holding. The industry word for this movement is a *shuffle*. Grouping by a key, joining two inputs on a key, producing one globally ordered result, and deduplicating across the whole input all tend to force one: records that have to meet must end up in one place. The exception worth knowing is an input that already arrives arranged by that key, in which case the engine may need no movement at all. 2. **Phases.** The crossings cut the listing into runs of steps that can be executed without moving any records. Most engines print the listing already grouped that way, and many call such a run a *stage*. A long run that never splits is telling you something concrete: every step inside it builds each output piece from exactly one input piece — a *narrow* dependency — so the whole run can be walked over each piece in one pass. 3. **Step width.** How many pieces of the input a step processes at the same time. This is the number that says how much of the cluster is working during that step, and it is not constant down the listing: a step after a crossing usually has the width the redistribution produced, while a step before it has the width the input was cut into. 4. **What the read was told to skip.** Two rewrites end up at the bottom of the plan and both are visible there. **Read-time column narrowing** tells the read to produce only the columns some later step actually uses. **A filter at the read** moves a condition down to the point of reading, so the rows it rejects are never produced at all. The printed text normally shows the read's column list and the conditions handed to it. ## Why those four and nothing else - They are the only facts in the text whose existence depends on there being more than one machine; everything else in the plan has a single-machine twin. - Each of the four is **checkable against the program you wrote**: you know how many times your job changes grouping key, so you know how many crossings to expect; you know which columns you use, so you can see whether the read was told about them. - They are all statements of **intent**. The text says what the engine means to do, not what will happen: it carries no record counts it has observed, no measure of how unevenly records will land, and no duration. ## What varies between engine models | Engine model | What its printed text looks like | |---|---| | A model that runs one grouping step at a time and writes every intermediate to disk before the next begins | Almost nothing is rewritten, so the printed shape is essentially the shape that was submitted, with one crossing per grouping step | | A model that keeps one fixed graph running and passes each record through it as it arrives | The text describes a graph that is already running; crossings are connections records are pushed across as they are produced, and widths are usually fixed for the life of the job | | A model that runs continuous work as a fast succession of small finite jobs | The text describes one of those small jobs; the same crossings and widths recur on every one of them | So "the engine prints a rewritten plan with pushed-down filters" is true of some engines here and false of others — in particular, a program written as per-record functions, where the author hands the engine a body to call once per record rather than a named operation, gives the engine nothing to push down and nothing interesting to print. ## How to read it in practice 1. Count the crossings, and compare that count with the number of times your program genuinely changes grouping key. More crossings than key changes is the interesting case. 2. Read the widths down the listing and look for a collapse — a step much narrower than its neighbours is where the cluster stops being used. 3. Look at the read's column list and its conditions, and check they match what you thought you asked for. 4. Stop. The text will not tell you how many records flow, how evenly they are spread, or how long anything takes; those come from measuring the run.

  • The read line lists three columns out of forty. Does that mean only three columns leave the storage layer?
    It means the read was instructed to produce three. How much is genuinely skipped depends on the source: a layout that stores columns separately can avoid touching the rest, while a row-oriented or opaque source may read everything and discard it after decoding. The plan records the instruction, not the effect.
  • Why would the printed text show no rewriting at all for a job that clearly could be narrowed?
    Because the engine can only rewrite what it understands. Where the author hands over a body of ordinary code to call once per record, the engine knows the body runs and nothing more, so it cannot move a condition through it. One engine model in this class also rewrites essentially nothing by design, running the shape it was given.
  • Does every engine of this class print a plan worth reading?
    No. What is printed ranges from a rewritten graph with the read's instructions shown, through a running graph annotated with per-step widths, to little more than the submitted shape. Before drawing conclusions, establish which of the four distributed facts this engine's text actually carries.

saying these in an interview costs you the question

  • Thinks the printed text reports how long each step took
  • Counts steps rather than points where records cross the network
  • Believes a width line reports how many bytes a piece holds
  • Assumes the read skipped columns merely because the program used few
  • Treats the text as identical across every engine of this class
  • Reads intent as observation and quotes it as run evidence