skip to content

Replanning at Runtime

The runtime's own remedy, and only one of them: it measures what a finished step actually produced and re-plans the rest on that evidence, which takes a model where steps finish.

on this pageshow

questions

4

A step finished and reported 200 MB where the plan assumed 200 GB - what can a runtime that re-plans mid-run change now?

level: middleimportance: must knowfreq 55%

answer

  1. the plan is not fixed at submission
  2. evidence arrives when a step finishes
  3. measurement replaces the opening estimate
  4. only steps that have not started
  5. no finished step, nothing measured

basics

~20 s

Runtime replanning - re-deciding unrun steps from the sizes a finished step actually produced - can change only what has not started: the shape of the next step's pieces, or a join switched to copying the small side. Finished work stands.

solid answer

~50 s

The plan as submitted is built from estimates. When a step finishes, its output is recorded - bytes and record counts, and the same numbers per destination - and a runtime that re-plans mid-run reads those actual figures and re-decides the part of the plan that has not started. With one input measured at 200 MB rather than 200 GB, the obvious change is the join that follows: instead of redistributing both sides so equal keys meet, it can send the whole small side to every worker so the large side never moves. It can also hand several undersized destinations to one worker, or divide one that came out far too large. What it cannot do is undo the step that already ran, and it cannot change the answer. It is also conditional: this needs a boundary where a step finishes and a description the engine is allowed to rewrite, and some engines in this class supply neither.

go deeper

for a junior

Recall that a job's shape is not necessarily fixed the moment you submit it: on some runtimes, what a finished step actually produced can change how the remaining steps run.

for a middle

Explain where the evidence comes from - the recorded output of a step that has completed, counted per destination - and be exact that only the unrun part of the plan is open to change.

for a senior

Say plainly that this is measurement after the fact rather than estimation before it, and that it is available per engine and per operator rather than everywhere, so it is a correction and not a guarantee.

for a principal

Judge how far a platform should lean on a correction that exists only at step boundaries, and decide what you standardise on for the jobs whose steps never finish.

## The correction this leaf is about **Runtime replanning** is the engine re-deciding a part of the plan it has not run yet, using the sizes a finished step actually produced instead of the estimate it started with. The situation in the question is its clearest case: the plan was built expecting 200 GB out of one branch, the branch has now completed, and the recorded output is 200 MB. Every decision downstream of that branch was made against a number that is wrong by three orders of magnitude, and some of those decisions have not been executed yet. Two words carry the whole mechanism. **Finished** - the evidence exists only once a step has completed and its output has been recorded somewhere the engine can count. **Actually** - this is measurement after the fact, not estimation before it, which is exactly what separates it from a query planner choosing a shape from statistics gathered about a table before anything ran. ## Why the opening estimate is wrong so often Something has to run first, so the initial plan is built from whatever the sources declare: file sizes, a schema, a recorded row count. Those describe the data as it sits in storage. They do not describe what a filter, an expansion or an earlier grouping will leave of it, and the error compounds with depth. A large input behind a selective predicate and a join arrives at the next step as a small fraction of itself, and nothing available at submission time knows the fraction. ## What a completed step hands back The evidence available at a step boundary is coarse but exact: - the total bytes and records the step emitted; - the same numbers **per destination** - the share of the output each unit of the next step will read. This is the figure that makes uneven work visible, as one destination far above the median of its peers; - in most designs, nothing about individual key values. The evidence is per destination piece rather than per key, which is why the available corrections are about the *shape of the pieces* and not about the key itself. ## What is open and what is closed | | can change after the measurement | cannot change | |---|---|---| | the step that finished | nothing | its records were produced and moved; that cost is paid | | the consumer's unit count | undersized destinations can be handed together to one worker, or an oversized one divided | the count the producing side already ran with | | an unrun join | may switch to copying the whole small side to every worker, so the large side never moves | a join whose redistribution has already happened | | the result | nothing | every move is meant to be result-preserving; the rows are the same | So the honest answer to *what can it change now* is: the join that has not started, and the number of units the next step runs with. Not the redistribution already paid for, and not the answer. ## Why a small measurement is the valuable one Of the corrections available, the one a 200 MB measurement unlocks is the largest. A join planned as a two-sided redistribution - both inputs sent across the network so records with equal keys meet on one worker - becomes **copying the small side**: the whole of the smaller input is sent to every worker, so the larger input is never grouped by key and never moves. The plan could not choose that at submission time, because it believed both sides were large. The measurement is what makes the choice both legal and safe. ## Where this does not exist at all This is the part that makes the subject answerable by an engineer who has used a different engine. The mechanism needs three things, and engines of this class supply different subsets: 1. **A boundary where a step finishes and its output can be measured.** A runtime that pushes records to their destination as they are produced, materialising nothing in between, never reaches one. 2. **A description the engine is allowed to rewrite.** Where the program describes what is wanted, the engine may re-decide how; where it is a chain of opaque per-record functions, there is nothing to re-decide even when the boundary exists. The oldest, two-phase disk-to-disk model in this family rewrites nothing. 3. **An implementation for the operator in front of it.** Even where the mechanism exists it typically covers groupings and joins rather than every step, and which corrections are enabled out of the box differs between engines. Treat it as a correction that *may* apply, never as a promise that the runtime will notice a problem and fix it. ## What an interviewer is listening for That you separate measurement after the fact from estimation before it; that you know only the unrun part of the plan is in play; and that you state the condition instead of the universal - 'on a runtime where steps finish and their output sizes can be measured'. A candidate who says 'the engine will just handle uneven pieces' has asserted something false on part of this market and only partly true on the rest.

  • Why wait for the whole step to finish rather than re-planning from the first few units that report?
    Partial evidence is biased: the units that report first are usually the small ones, and the one that matters is the one still running. Some runtimes do act on partial evidence for limited decisions, but the change that flips a join to copying the small side needs a total size, and a total exists only once every producing unit has reported.
  • If the measurement is so much better, why does the plan start from an estimate at all?
    Because something has to run first. The shape of the earliest steps must be chosen before any of them has produced anything, so it comes from what the sources declare - file sizes, a schema, a stored row count. Replanning is a correction applied at the first boundary where reality is known, not a replacement for planning.
  • Does re-planning change what the job outputs?
    No. Every move is meant to be result-preserving: merging destinations, dividing one, or copying the small side of a join all yield the same rows. What changes is how many units run, how much crosses the network, and how long the job takes.

A builder quotes the whole house from the drawings, then measures the first room once its walls are up. The rooms already built stay as built; the order for everything still unbuilt is redone from the measurement.

saying these in an interview costs you the question

  • Says the runtime re-plans from statistics collected before the job started
  • Thinks re-planning can undo a redistribution that has already run
  • Assumes every engine of this class re-plans mid-run
  • Believes the entire plan is rebuilt, not only the unrun part
  • Claims the correction can change the job's output
open as a page

After a redistributing step finishes and reports its actual output sizes, which changes can a runtime still make to the unrun grouping or join?

level: middleimportance: should knowfreq 45%

basics

~20 s

Commonly three result-preserving moves, where an engine offers them: hand several undersized destinations to one worker, divide one that came out far too large, and switch a join to copying the small side once that side measures small.

open as a page

Runtime re-planning divided an oversized destination piece in a join, yet left an equally oversized one in a grouping alone - what differs?

level: seniorimportance: should knowfreq 40%

basics

~20 s

A join can be cut inside one key, because the other side's matching rows are copied to each sub-piece and the union is the same result. A grouping cannot: every record for one key must meet in one place.

open as a page

Your batch jobs lean on runtime re-planning around uneven pieces, and the same logic must now run continuously over an endless input - what coverage is lost, and how do you plan around it?

level: principalimportance: nice to knowfreq 28%

basics

~20 s

Nearly all of it. Re-planning needs a finished step's measured output, and an endless input has none. Uneven-work decisions move to design time and to the author, and changing one means a restart rather than a mid-run correction.

open as a page