skip to content

Job as a Dataflow

A written program becomes a graph of steps before any of it runs: what the engine is allowed to rewrite, when execution actually starts, and why two identical runs can disagree.

on this pageshow

explore

questions

25

What does a cluster engine know about a named operator that it does not know about a function you hand it to run per record?

level: juniorimportance: must knowfreq 70%

answer

  1. meaning against a bare instruction
  2. one is read, one is only called
  3. named operator has known semantics
  4. a body says only 'call me per record'
  5. visibility is what buys the freedom

basics

~20 s

A named operator carries its meaning - the engine knows it keeps rows, or groups by a key - so it may reorder, narrow, fuse or skip it. A handed-over body means only 'call this per record', so it runs exactly where it stands.

solid answer

~50 s

A cluster engine turns your program into a **step graph**: the ordered set of steps it derives from the program, each naming the steps whose output it reads. On a **declared-operator surface** you name operations the engine already understands - keep these rows, produce these fields, group by this key, join on that key - so the step that lands in the graph carries a meaning the engine can act on. On a **per-record function surface** you hand over a function the engine calls once per record; the step that lands in the graph says only `call this body on every record`. That difference is the whole trade. Where an engine acts on meaning at all, only the named form can be reordered, narrowed or dropped; the body can only be called, in the position you put it. What you buy by handing over a body is expressiveness - anything the language can do.

go deeper

for a junior

Recall the one-line difference: a named operation tells the engine what you want, a handed-over function tells it only that something must be called once per record. That is the whole first-screen answer.

for a middle

Explain why visibility is what buys freedom: a step the engine can read may be reordered, narrowed, fused or skipped where that is provably safe, and a step it can only call may not. Say that engines differ in how much they rewrite at all.

for a senior

Show that you decide where the seam falls in a real job rather than preaching one surface. Name work that genuinely has no named equivalent, and be clear about what you gave up by handing it over.

for a principal

The platform-level angle is which surface a whole team standardises on: the named vocabulary is portable and legible to tooling, the function surface is expressive and hard to reason about, and the right default depends on what your engine actually does with a declared graph.

## The step graph and what goes into it A distributed processing engine does not run your program line by line. It first derives a **step graph** - the ordered set of steps it takes from your program, each step naming the steps whose output it reads - and only then decides how to run it. What the engine is able to do with that graph depends entirely on what each step *says*. Two authoring surfaces put very different things into it. - A **declared-operator surface** is one where the author names operations the engine already understands - keep the rows matching a condition, produce these fields, group by this key, join on that key, sum this column - over a **distributed collection**, the engine's handle on a set of records spread across the cluster. What lands in the graph is *discard rows whose amount is below one hundred*. - A **per-record function surface** is one where the author hands the engine a function and the engine calls that function once per record. What lands in the graph is *call this body on each record*. The engine knows the body will run. It knows nothing about what the body does. ## Meaning is what buys freedom Because a named operator carries meaning, an engine that has a **plan rewriter** - the component that edits the declared graph into an equivalent, cheaper graph before running it - can reason about that step rather than merely schedule it. At the level that matters here, four freedoms follow from being able to read a step: - **reorder** it against another step, where the result is provably the same; - **narrow** what the read produces, because the fields the named steps reference are visible; - **fuse** it with adjacent per-record work into a single pass, so no intermediate collection exists between them; - **skip** it, where it provably changes nothing. Each of those rewrites has conditions and limits of its own, which is a separate subject. The point on this axis is earlier and simpler: **none of them is available for a step whose content the engine cannot read.** A handed-over body is a box with exactly one contract - *call me once per record, in this position* - so the only plan the engine may legally make for it is to call it once per record, in that position. | What the engine holds | Named operator | Handed-over body | |---|---|---| | In the graph | the operation's meaning | "call this per record" | | Fields touched | visible | unknown; it must assume any | | Can be reordered | yes, on engines that rewrite at all, where equivalence is provable | no | | Can be dropped | yes, where it provably does nothing | no | | Checked before the run | usually, against the known record shape | only as it runs, on a worker | | Expressiveness | limited to the named vocabulary | anything the language can do | ## What varies between engines Do not carry one engine's behaviour into this as if it were the model of the class. - **Not every engine rewrites anything.** A **two-phase disk-handoff model** - one that runs a single grouping step at a time and writes every intermediate result to storage before the next begins - leaves a rewriter almost nothing to work with. There, a declared surface buys checking and readability rather than a cheaper plan. - **Which surface is idiomatic differs.** In a **continuous record-at-a-time model**, where one fixed graph of steps stays running and each record passes through it as it arrives, a per-record body is the normal way to express logic that remembers something between records for a key; often no named operator exists for it. In a **repeated-small-batch model**, which runs continuous work as a fast succession of small finite jobs, the declared surface is usually the default one. - **The freedoms are not all-or-nothing.** Some engines can still place, connect or fuse a body with its neighbours without reading it. What none of them can do is change what the body *means*. ## The mixed job is the normal job Real jobs use both surfaces. The named vocabulary is deliberately small - it is small precisely so the engine can reason about it - and plenty of honest work has no name in it: parsing an awkward text format, a bespoke scoring rule, a call out to something else. Handing over a body for that is correct, not a defect. The question is never "named good, body bad"; it is **where the seam between them falls**, and how much of the job you leave on the side the engine cannot see. ## What an interviewer is listening for "Named operators are faster" is a memorised slogan and invites an immediate follow-up the candidate then cannot answer. The model is: a named operator tells the engine *what* you want, so it stays free to choose *how*; a handed-over body tells it only *that* something happens, so it must do exactly what you wrote, where you wrote it. A candidate holding that can reason about a job, and an engine, they have never seen.

  • If the engine cannot read a handed-over body, how does it know where to run it?
    Position in the graph is all it needs. The author placed the step between two others, so the engine calls the body once per record on whatever records reach that point, on whichever worker holds them. It schedules the call; it does not interpret it. That is exactly why the step stays where it was put.
  • Does writing on the declared surface make a job faster by itself?
    No. It makes the job *legible*, which is what lets an engine that rewrites do so. Where the engine rewrites little - notably a model that writes every intermediate result to storage before the next grouping step - the same program on either surface can run much the same way, and the declared form still buys earlier checking and clearer code.
  • Can a single job use both surfaces at once?
    Normally it does, and that is not a compromise. Engines in this class generally let a declared graph contain a step whose body is ordinary code, and let a named operation be applied to the output of one. The design decision is which work goes on which side, not choosing one surface for the whole job.

Telling a travel agent "cheapest return to the coast on Monday, aisle seat" states a goal, so they stay free to combine flights, switch carriers or rebook you when something is cancelled. Handing them a sealed itinerary marked "execute in order" states a procedure: they can carry it out, but they cannot improve it, because they were never told what you were actually trying to achieve.

saying these in an interview costs you the question

  • Claims the engine reads and improves any code you write, whatever surface it is on
  • Says a handed-over per-record body is always slower than the named form
  • Believes every engine of this class rewrites the declared graph before running it
  • Thinks a job must be written entirely on one surface or the other
  • Treats the difference as a matter of syntax or personal style rather than visibility
  • Assumes the engine can work out what a body does by inspecting the code
open as a page

Two runs of one job over the same input write the same rows in a different order — why?

level: juniorimportance: must knowfreq 65%

basics

~20 s

Nothing in the job asked for an order, so none was kept. The work is cut into pieces processed by many workers at once, and rows appear in whatever sequence those pieces finished and were collected. Only a declared ordering step makes row order stable.

open as a page

In an engine that runs no declared step until something demands an answer, what distinguishes a describing call from a demanding one?

level: juniorimportance: must knowfreq 74%

basics

~20 s

A describing call only adds a step to the graph the engine will run and computes nothing. A demanding call asks for an answer or for output to be written, and that demand is what makes the recorded steps execute.

open as a page

A program reads a hundred-column input, keeps two columns and filters on a date - which two edits does the engine make before running it?

level: juniorimportance: must knowfreq 68%

basics

~20 s

Two: read-time column narrowing, which tells the read to produce only the two columns later steps use, and filter at the read, which moves the date condition down so rejected rows are never produced. Both shrink what enters the job.

open as a page

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%

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.

open as a page

A job declares three per-record steps in a row with no redistribution between them - what does an engine that fuses them actually run?

level: middleimportance: must knowfreq 54%

basics

~20 s

One pass over each record that applies all three operations in sequence. Step fusion removes the intermediate collections between the steps, not the work each operation does. The fused chain must end wherever records have to move between workers.

open as a page

A project and a filter are written after a step whose body the engine cannot read. Why does neither reach the read?

level: middleimportance: must knowfreq 70%

basics

~20 s

Both rewrites need to know what the unreadable body reads and produces. A condition on a value the body computed cannot be evaluated before the body has run, and a read cannot be narrowed to fields nobody can prove the body ignores, so whole records still arrive.

open as a page

In a job's printed plan, how do you locate the points where records cross the network, and what does a long unbroken chain mean?

level: middleimportance: must knowfreq 64%

basics

~20 s

Look for where the listing breaks: a crossing is a point where records must be sent between workers, and most engines cut the printed steps into runs around it. An unbroken chain means every step builds each output piece from one input piece.

open as a page

A job spends most of its time in a step whose body the engine cannot read. Which structural changes reduce that cost without rewriting the body?

level: seniorimportance: must knowfreq 58%

basics

~20 s

Do by hand what the engine cannot deduce: apply the conditions and select the fields before the step, hand the body named fields instead of the whole record, move the parts expressible as named operators out of the body, and call it after a reduction rather than on every raw record.

open as a page

A job's middle step is ordinary author-written code the engine cannot look inside. What does the engine still know about that step?

level: juniorimportance: should knowfreq 60%

basics

~20 s

The engine knows only that this step must be called once per record, where the author put it, along with what it was handed and the shape of its result. It cannot see which fields the body reads, whether repeating it is safe, or what it costs.

open as a page

When a step names the two record fields it uses instead of reading them inside a function body, what does the engine gain?

level: middleimportance: should knowfreq 46%

basics

~20 s

Two things become visible: which fields the step needs, and whether those names and types exist. The engine can then check the step before any record moves, and can tell the read to produce only the fields some step actually references.

open as a page

A colleague says a cluster engine will optimise whatever program you give it - when is that true, and when is it false?

level: middleimportance: should knowfreq 55%

basics

~20 s

Only a program whose steps the engine understands can be improved. Work handed over as a function body is executed close to literally, and some engine models rewrite next to nothing on either surface, so the claim is conditional on both the surface and the engine.

open as a page

A total over a billion amounts differs in its last digits between two runs over the same input — what explains it?

level: middleimportance: should knowfreq 46%

basics

~20 s

Partial results were combined in a different order, and each addition in the usual binary floating-point format rounds. Nothing ever promised a combination order, so the two totals are both correctly rounded and simply not the same number.

open as a page

A per-record step body reads the wall clock and a random draw — what does that cost the job's repeatability?

level: middleimportance: should knowfreq 52%

basics

~20 s

The output becomes a function of when and where the body ran, not of the input. Two executions disagree, and one execution can disagree with itself when a lost piece is recomputed and stamped with later values than its first attempt produced.

open as a page

Why is a job's parse error reported at the call that writes output rather than at the parse step, and how do you localise it?

level: middleimportance: should knowfreq 56%

basics

~20 s

Because the parse step only recorded itself; no record reached it until the write demanded an answer, so the failure surfaces there, carried back from a worker. Localise it by demanding answers from shorter prefixes of the graph over a small sample.

open as a page

A stopwatch around a job's declared steps reports four milliseconds while the run takes forty minutes — what did it measure?

level: middleimportance: should knowfreq 41%

basics

~10 s

It measured graph construction: recording the steps, and perhaps resolving names and types. No records were read inside the timed region, so the number describes the program's description of the work, not the work.

open as a page

In a job's step graph, which filters may the engine move below a grouping step, and which must stay above it?

level: middleimportance: should knowfreq 44%

basics

~20 s

A condition naming only grouping keys may move below, because rejecting whole records cannot change a surviving group's result. A condition on an aggregate must stay above: the value does not exist yet, and dropping its inputs would change it.

open as a page

A printed plan shows one step handling 2,000 pieces at once and the next handling 1 — what does that tell you?

level: middleimportance: should knowfreq 47%

basics

~20 s

That 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.

open as a page

A job is written entirely as per-record function bodies - which parts would you re-express as named operators, and which would you leave?

level: seniorimportance: should knowfreq 42%

basics

~20 s

Re-express anything with a named equivalent - selecting, discarding, combining, grouping, joining, aggregating, sorting - and keep as a body only work with no name: bespoke parsing, custom scoring, per-key logic that remembers. Split fat bodies so the nameable half becomes visible.

open as a page

Before calling two runs of one job equal, what has to be pinned and what equality should you actually assert?

level: seniorimportance: should knowfreq 40%

basics

~20 s

Pin the input records, the code and the settings that shape the plan, the step width, and any clock or random read inside a step body. Then assert a chosen equality — same rows ignoring order, same order, exact values, or values within a stated bound — because some variation survives pinning.

open as a page

Why is an input read twice when one distributed collection feeds two separate calls that each demand an answer, and which engine models never re-read it?

level: seniorimportance: should knowfreq 47%

basics

~20 s

Because a declared step holds no result of its own: each demand runs the graph again from its sources. Models that write every intermediate to storage, and models that keep one graph running and fan records out to both consumers, do not re-read.

open as a page

Some engine models rewrite nothing before running - what does the written order of steps then decide, and what falls to the author?

level: seniorimportance: should knowfreq 36%

basics

~20 s

Written order becomes executed order: no columns are dropped at the read, no condition moves down, and nothing is fused across a grouping step. The author has to narrow the read, filter first and fold partial results before records move.

open as a page

A filter written above an author-supplied per-record body never reaches the read - which edges stop a rewriter, and how do you move them?

level: seniorimportance: should knowfreq 48%

basics

~20 s

Rewriting stops at three edges: a body the rewriter can only call and never read, a read with no way to accept a condition, and a step whose output something else depends on. Authors move the first by declaring the readable part of the work in named operators instead.

open as a page

A per-record body written for a second language runtime dominates a job's time. What does each record pay, and what changes if records cross in batches?

level: seniorimportance: should knowfreq 54%

basics

~20 s

Each record leaves the engine's own runtime and comes back: it is converted into a form both sides read, crosses a process boundary twice, and is converted back, on top of the body's own work. Handing many records over at once pays the crossing cost once for the batch instead of once per record.

open as a page

A job's plan shows few points where records cross the network and even widths, yet the run takes hours — what can the printed text not tell you?

level: seniorimportance: should knowfreq 43%

basics

~20 s

Everything about volume and evenness. The printed text states the shape the engine intends to run: how many records flow along each edge, how unevenly they land across pieces, how much a crossing moves and how long anything takes are all measurements it does not contain.

open as a page