skip to content

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%

answer

  1. same answer, different shape
  2. too many small, or one too large
  3. the consumer is reshaped, not the producer
  4. one side turns out to be small
  5. grouping and join boundaries only

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.

solid answer

~50 s

A redistributing step - one that cannot be computed from what a worker already holds, so every record is sent to the worker owning its key - leaves behind a measured size per destination. Three corrections follow from that measurement. **Merging**: if the destinations came out far below the size a unit should handle, a group of neighbouring ones is given to a single worker, so the consumer runs with fewer, larger units and pays less per-unit overhead. **Dividing**: a destination measured far above its peers is split into sub-pieces that run in parallel, which for a join also means copying the other side's matching rows to each sub-piece. **Copying the small side**: if one input of an unrun join measures small, the whole of it is sent to every worker and the large side never moves. All three leave the rows unchanged; which of them an engine implements, and for which operators, differs.

go deeper

for a junior

Recall that these corrections change how the work is arranged, never what the job outputs; the rows are the same either way.

for a middle

Name the three moves and the measurement that triggers each, and be clear that it is the step about to run that is reshaped, not the one that just finished.

for a senior

Show where the moves stop: they land at grouping and join boundaries, coverage differs between engines, and a job whose input never ends offers no finished step to measure.

for a principal

Decide how much of your uneven-work strategy can be delegated to a correction whose availability varies by engine and operator, and what the fallback is when it does not fire.

## The evidence, and the window in which it is usable A **redistributing step** is a step that cannot be computed from what one worker already holds, so every record is first sent to the worker that owns its key. When one completes, the runtime knows something it could only guess at before: how many bytes and records each destination actually received. **Runtime replanning** - the engine re-deciding a part of the plan it has not run yet from the sizes a finished step actually produced - spends that evidence on the step about to consume those destinations. The window is narrow and one-directional: the producing side is done, the consuming side has not started, and only the consuming side is open. ## Move one - fewer, larger units for the consumer A step is planned with some number of destinations. If the data turns out much smaller than assumed, most of them hold very little. Each unit costs something regardless of how much data it holds: it is scheduled onto a worker, it opens and closes its inputs and outputs, it reports when it is done. Thousands of units holding a few megabytes each spend more of the run on that overhead than on the work. The correction is to give a group of neighbouring destinations to one worker, so the consumer runs with fewer, larger units. Two things are worth being precise about: - it is the **consumer** that is reshaped. The producing side already ran with the count it ran with, and that is not undone; - it is **not** the same decision as choosing how many pieces to cut the input into before the job starts. That choice is made from guesses at submission time; this one is made from a measurement. Engines that do this generally merge neighbouring destinations rather than arbitrary ones, and aim at a target size rather than a target count. The target exists in order to keep a unit large enough to amortise its own overhead and small enough to fit a worker's memory; its value differs between engines and is not worth memorising. ## Move two - dividing one that came out far too large The mirror case: one destination measures many times its peers, so the consumer would run for as long as that one unit takes while everything else idles. Where the division is legal, the runtime cuts that destination into sub-pieces that run in parallel. For a join, the sub-pieces only produce the right answer if the other side's matching rows are copied to each of them, so the move buys parallelism with extra reads and extra network. The limits on when this is legal are a subject of their own. ## Move three - copying the small side The plan assumed two large inputs and therefore planned to redistribute both, sending records from each across the network so equal keys meet. The finished branch then reports that one side is genuinely small. The unrun join can be switched to **copying the small side**: the whole of the smaller input goes to every worker, and the larger input is neither grouped by key nor moved. This is usually the largest single win available, because it deletes an entire redistribution of the big input rather than reshaping one. It can only be decided after a measurement when the small side is itself computed. A branch that filters, joins and aggregates before reaching this point has no trustworthy declared size; only its finished output does. ## The three moves side by side | move | what triggers it | what changes | what it costs | |---|---|---|---| | merging destinations | most destinations far below a useful unit size | the consumer runs fewer, larger units | less parallelism if overdone | | dividing a destination | one destination far above its peers | that unit becomes several running in parallel | the other side of a join is read and sent more than once | | copying the small side | one unrun join input measures small | the large input is never grouped by key and never moves | the small input is held in every worker's memory | ## What it reaches, and what it does not The corrections land where the next step's shape depends on the sizes just produced, which in practice means the redistribution boundaries of groupings and joins. A chain of per-record steps has nothing to re-decide. Which operators are actually covered differs between engines, and so does whether a given correction is applied without being asked for. And there is a whole class of job that gets almost none of it: a job whose input never ends has no step that finishes, so there is no completed output to measure and nothing on which to base a correction. ## What never changes All three moves are meant to be result-preserving - the rows the job emits are the same rows. What changes is the number of units, the volume crossing the network, and the runtime. If a proposed correction would change the answer, it is not this mechanism.

  • Why bother merging undersized destinations - they are small, so surely they are cheap?
    Each unit costs scheduling, input and output handling and reporting whether it holds four megabytes or four hundred. With thousands of tiny units, that fixed cost dominates the run, and the units also become the granularity of everything downstream. Merging converts fixed overhead into useful work without touching the result.
  • Why can the switch to copying the small side only be made after a step finishes?
    Because 'small' is a measured fact about a computed input. When one side of the join is the output of filters, joins or an aggregation, nothing available at submission time predicts its size well enough to risk the choice - copying an input that turns out to be large would put it in every worker's memory at once.
  • Does the runtime re-plan every kind of step this way?
    No. It acts where the next step's shape depends on sizes that were just produced, which is the redistribution boundary of a grouping or a join. Per-record steps offer nothing to re-decide, and engines differ in which operators they cover at all.

saying these in an interview costs you the question

  • Claims merging destinations changes the job's result
  • Says the runtime re-runs the step that already finished
  • Believes copying the small side is chosen from the declared schema, not a measured size
  • Treats merging destinations as the same decision as choosing the count up front
  • Thinks dividing a destination is free of extra work on the other side