Some engine models rewrite nothing before running - what does the written order of steps then decide, and what falls to the author?
answer
- written order is executed order
- no rewriter, no narrowing at the read
- intermediates written out between grouping steps
- the author moves the condition down
basics
~20 sWritten 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.
solid answer
~50 sTwo situations produce this. One is a two-phase disk-handoff model, which runs one grouping step at a time and writes every intermediate result to disk before the next begins - it has almost nothing to rewrite, because each step is submitted and completed on its own. The other is a program written entirely as author-supplied per-record bodies, where even an engine with a plan rewriter can only call what it was handed. In both, the graph the author wrote is the graph that runs, in the order it was written. So the author does the work a rewriter would have done: select only the fields later steps need, apply the most selective conditions immediately after the read, fold partial results per worker before records are sent across the network, and avoid reading the same input twice. That habit costs nothing where a rewriter exists and is the whole difference where one does not.
go deeper
Take away one rule: put the field selection and the conditions as early in the program as the answer allows. Some engines would move them for you and some never will, and writing them early is right either way.
Explain why written order becomes executed order when a step's output is written out before the next step starts, and be able to list what a rewriter would have done that now falls to the author.
Show that you can tell rewriter-driven edits from runtime behaviour, and that you write narrow reads, early conditions and local folds as a matter of habit rather than trusting a rewrite you have not verified happened.
The judgment call is whether the engine model itself, not the plan, is the thing to change. A workload dominated by intermediates written out between steps is not a tuning problem, and the migration cost is weighed against what a whole-graph rewriter would return.
## The negative case, and why it is still asked Most discussion of plan rewriting assumes a rewriter. It is worth being able to describe the case where there is none, both because such systems are still in production and because the same shape appears inside modern engines the moment an author stops declaring and starts handing over code. Two distinct situations: - **A two-phase disk-handoff model**: an engine model that runs one grouping step at a time and writes every intermediate result to disk before the next begins. There is no whole-graph view to optimise. Each step is submitted, runs to completion, leaves its output in the cluster-visible storage, and the next step reads it back. - **A program written on a per-record function surface**, where the author hands the engine a function to call once per record and the engine knows only that the function runs. Even an engine with a capable rewriter has nothing to analyse. The runtime may still chain adjacent per-record functions into one pass over each record, but that is execution behaviour, not a plan rewrite, and it changes no ordering and drops no columns. ## What the written order now decides | what a rewriter would have done | what happens without one | |---|---| | narrow the read to the fields later steps use | the read produces whole records, and unused fields are decoded, carried and moved | | move a condition down to the read | the condition runs exactly where it was written, after everything above it has already processed the rejected records | | fuse adjacent per-record steps | adjacent per-record work may still be chained in one pass at run time, but nothing is folded across a grouping step | | reorder independent work | the declared order is the executed order | | fold partial results before records move | nothing is inserted; if the author did not write it, every record crosses the network | The practical effect is that cost becomes a property of the text. Two programs computing the same answer, differing only in where the author put the condition, differ in runtime by whatever the rejected records cost every step above the condition. ## What the author must do by hand 1. **Project early.** Drop fields that no later step reads, immediately after the read. This is the edit with the largest effect on everything downstream, because a record's width is paid at every step and again whenever records move. 2. **Filter first, most selective first.** Put conditions as close to the read as correctness allows. A condition on a grouping key belongs before the grouping; a condition on an aggregate cannot move and should not be forced. 3. **Fold locally before records move.** Combine what can be combined on the worker that already holds the records, so that fewer records are sent. Whether a partial fold is valid at all depends on the combining operation, which is the parallel-reduction subject and not this one - here the point is only that no rewriter is going to insert the fold for you. 4. **Read the input once.** Without a whole-graph view, two branches that read the same input read it twice. If both are needed, arrange the job so the read happens once and both branches proceed from it. 5. **Keep the step count honest.** Where every step writes its intermediate to storage, a step is not free the way a fused step is; combining work that can be done in one pass over a record genuinely reduces passes. ## Saying it without over-claiming The sentence to avoid is 'the engine optimises your program'. What is true across this class is narrower and more useful: - Only what is declared in operators the engine already understands can be rewritten. - Engines differ in how much they rewrite - some analyse the whole declared graph before running any of it, some rewrite a little, and the oldest model in this family rewrites essentially nothing. - Some engines also change part of the plan mid-run on statistics measured from work already finished. That is a different mechanism with a different owner, it applies to some operations only, and it does not rescue a program that never gave the rewriter anything to read. ## The habit that survives every model Everything in the hand-written list above is free where a rewriter exists: if the engine would have moved the condition down, writing it down changes nothing, and if it would not have, you have just made the largest available improvement. That asymmetry is why experienced authors write narrow reads and early conditions regardless of which engine they are on - and why an interviewer treats it as a marker of somebody who has worked on more than one.
- Records still flow one at a time through adjacent steps here. Does that mean something was fused?No. Chaining adjacent per-record work into one pass is execution behaviour a runtime can have with no rewriter involved; it drops no columns, moves no condition and changes no ordering. A plan rewrite edits the graph before the run. The two look alike from outside and have different causes, which is exactly what interviewers probe.
- Why is writing the narrow read and the early condition by hand recommended even on an engine that would do it for you?Because the downside is zero and the upside is large. If the rewriter would have made the edit, writing it changes nothing about the result or the cost. If it would not - a body it cannot read, an input that cannot take a condition, a surface with no named operators - you have made the largest single improvement available. The habit also survives moving between engines.
saying these in an interview costs you the question
- Says every processing engine rewrites the program before running it
- Treats one-pass record flow as evidence of a plan rewrite
- Expects unused columns to be dropped without the author selecting them
- Believes mid-run replanning compensates for a program with nothing to analyse
- Writes conditions at the end of the program and relies on the engine to move them