skip to content

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%

answer

  1. not every filter may move down
  2. grouping keys against aggregate results
  3. dropping rows changes the aggregate
  4. a condition on a total stays above

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.

solid answer

~50 s

The rule is equivalence, not cheapness. Moving a filter below a grouping step is allowed only where the rewritten graph still produces the same answer. A condition on the grouping key satisfies that: a record's key decides which group it lands in, so rejecting records by key removes whole groups and leaves every surviving group's result untouched - and doing it early means fewer records ever move between workers. A condition on an aggregate cannot move: the total or count does not exist until the grouping has run, and removing records before it would change the very number being tested. A condition on a non-key column is the interesting middle case - it usually cannot move, because dropping some of a group's records alters that group's aggregate. Engines differ in how hard they try; some attempt none of this, and a step whose body the rewriter cannot read blocks it outright.

go deeper

for a junior

Remember that filtering earlier is usually cheaper but is not always allowed. A test on the value you grouped by can happen first; a test on a total that grouping produces obviously cannot.

for a middle

State the rule as equivalence: a move is permitted only where the answer is unchanged. Then work the three cases - grouping key, aggregate result, other column - and say why the middle one is the interesting one.

for a senior

Demonstrate that you write the early key condition yourself rather than trusting the rewriter, that you check where the condition ended up, and that you know a body the rewriter cannot read will pin the condition in place.

for a principal

The angle is how much correctness your platform delegates to a rewriter. Aggressive reordering buys performance and creates a class of subtle wrong-answer risk; deciding which surfaces your teams write on largely decides how much of that risk you carry.

## The rule behind the rewrite A **plan rewriter** - the engine component that edits the declared graph into an equivalent, cheaper graph before running it - may reorder steps only where the reorder is *equivalence-preserving*: the rewritten graph produces the same answer for every input. Cheapness is the motive; equivalence is the permission. Moving a filter below a grouping step is the clearest place to see the two come apart, because the cheap move and the correct move are not always the same move. Why the engine wants the move at all: a grouping step needs records that other workers are holding, so records must be sent across the network first. Every record removed before that point is a record that is never moved, never held in the grouping's working memory, and never combined. Filtering earlier is one of the largest single wins available to a rewriter. ## The three cases | the condition names | may it move below the grouping? | why | |---|---|---| | only grouping keys | yes | a record's key decides its group, so rejecting records by key removes entire groups and changes nothing about the groups that survive | | an aggregate result (a total, a count, a maximum) | no | the value does not exist before the grouping runs, and removing its input records would change the number being tested | | a non-key column of the record | usually no | it would drop part of a group, so every aggregate over that group would be computed from a different set of records | The middle row is the one candidates get wrong in both directions. Someone who says 'filters always move down' publishes wrong totals. Someone who says 'filters can never cross a grouping step' gives up the single best rewrite available, and also mis-explains why: keys can cross precisely because they cannot split a group. The third row has exceptions worth naming rather than glossing: if the condition on a non-key column happens to be implied by the aggregate being computed, or if the grouping is over records the condition cannot reject, the move is still equivalent. A rewriter that cannot prove the implication will not attempt it, which is the normal outcome. ## What blocks the move even when it is valid - **A step the rewriter cannot look inside.** If the condition calls an author-supplied body, the rewriter can call it but cannot reason about it - it cannot establish that the body behaves the same on every record regardless of position, so it leaves the condition where it was written. - **Missing values.** A condition that rejects records whose key is absent behaves differently above and below a grouping that treats absent keys as one group of their own. A careful rewriter refuses the move rather than guess. - **An engine that does not attempt it.** A model that runs one grouping step at a time and writes every intermediate result to disk before the next begins rewrites essentially nothing, so the author's written order is the executed order. There, the early filter is something the author writes, not something the engine grants. ## What an author should do about it The practical answer, and the one a good candidate volunteers: 1. **Write key conditions early yourself.** It costs nothing if the rewriter would have moved them, and it is the whole win if it would not. 2. **Accept that an aggregate condition runs last**, and reduce what reaches the grouping by other means - narrowing columns, rejecting records on keys, folding partial results per worker before the records move. 3. **Do not assume a rewrite happened.** The engine's own rendering of what it intends to run will show where the condition ended up; reading it is a separate skill and belongs with the plan-reading material, but the habit of checking is part of this one. ## The trap in the vocabulary 'The filter moved down' is said about two different edits. One is this one: past a grouping step, into the part of the graph that runs before records are redistributed. The other is the condition reaching the read itself, so rejected rows are never produced. They travel together in the common case - a key condition moves below the grouping and then continues to the read - but they are separate moves with separate conditions, and only the second can be blocked by an input that has no way to accept a condition.

  • Why is moving a key condition below the grouping worth so much?
    Because the grouping needs records that other workers hold, so records are sent across the network before it runs. Every record rejected before that point is never moved, never held in the grouping's working memory and never combined. Filtering below the grouping usually continues down to the read, where the rejected rows are not even produced.
  • A condition names a non-key column. Can it ever move below the grouping?
    Only where the rewriter can prove that dropping those records leaves every surviving group's result unchanged - for example where the aggregate is defined over exactly the records the condition keeps. Otherwise the move would compute each aggregate from a different set of records, so it is refused. Being unable to prove it is the usual outcome, and refusing is the correct behaviour.
  • How do missing values interfere with the move?
    A condition on a key that some records lack can behave differently on either side of the grouping, because a grouping may collect absent keys into a group of their own. If the rewriter cannot establish that both placements reject the same records, it leaves the condition where the author wrote it. Engines differ in how much of this they attempt.

saying these in an interview costs you the question

  • Moves a condition on a computed total below the grouping step
  • Says no condition may ever cross a step that combines records
  • Assumes a condition on any column is safe to apply before grouping
  • Judges the move by cost rather than by whether the answer is preserved
  • Assumes every engine attempts this rewrite at all