skip to content

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%

answer

  1. one key, one destination
  2. cutting between keys is always fine
  3. cutting inside a key needs reconstruction
  4. a join copies the other side
  5. an aggregate needs every record together

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.

solid answer

~60 s

Both pieces are oversized, but only one of them can be cut. The key-to-destination rule - the function turning a record's key into the number of the destination that receives it - sends every record for one key to one place, and that is not negotiable for either operator. The difference is what may be done *inside* a key. For a join, a heavy key's records can be spread over several sub-pieces provided the other side's matching rows for that key are copied to each; every pair still meets somewhere, and the union of the sub-results is the join. For a grouping, the aggregate has to see all of that key's records to produce its one row, so splitting them yields partial results that nobody combines - unless the aggregate is rewritten into two passes, which is not something a runtime will generally do for an arbitrary user-supplied aggregate. If the grouping's piece were oversized because it held many medium keys, it could be cut on key boundaries; the case that defeats it is bulk concentrated in a single key.

go deeper

for a junior

Recall the rule underneath all of this: the same key always goes to the same destination, so one key's records cannot be shared out over several workers just because there are many of them.

for a middle

Explain the difference between cutting between keys, which keeps each key whole and is always sound, and cutting inside one key, which needs a way to reconstruct the whole from the parts.

for a senior

Show the reconstruction argument for the join - the other side's matching rows copied to every sub-piece, so every pair still meets - and say why an aggregate has no equivalent without a combining pass.

for a principal

Set the boundary policy: which uneven-key cases you expect the runtime to absorb, which you accept will reach production untouched, and what the team does at that point.

## The rule both operators obey The **key-to-destination rule** is the function that turns a record's key into the number of the destination that receives it. The same key always yields the same destination, which is exactly why one key's records cannot simply be scattered: they arrive together by construction. A **heavy key** - one key value covering a large share of all records - therefore lands on one destination alone, and that destination is the oversized piece in both halves of the question. **Runtime replanning** - the engine re-deciding part of the plan it has not run yet from the sizes a finished step actually produced - can see that one destination is many times its peers. Seeing it and being allowed to cut it are two different things. ## Which cuts are legal at all There are exactly two kinds of cut, and they have different legality: 1. **Cutting between keys.** The oversized destination holds many keys and is simply the fattest range. Splitting it on a key boundary keeps every key whole, so each sub-piece computes complete results for its own keys. This is legal for a grouping and for a join alike, and a runtime that divides pieces at all will generally do it. 2. **Cutting inside one key.** The oversized destination is one key's records. Now a split separates records that have to be considered together, and whether that is recoverable depends on the operator. The scenario in the question is the second kind, and the two operators answer it differently. ## Why the join can be cut inside a key A join pairs each record on one side with each matching record on the other. If a heavy key's records on the large side are spread over four sub-pieces, and the matching rows for that key from the other side are **copied to all four**, then every pair that existed before still exists in exactly one sub-piece, and the union of the four outputs is the same set of rows as before. Nothing has to be combined afterwards; the correctness argument is that pairing is per-record, and each record still meets everything it matched. That is not free: - the other side's matching rows are read and sent once per sub-piece, so four sub-pieces means four copies of them crossing the network; - it only pays when the piece genuinely dominates the step's runtime - a piece twice its peers is not worth the replication; - if the heavy key is heavy on **both** sides, the copies multiply and the move stops being attractive. ## Why the grouping cannot An aggregate over a key produces one result from all of that key's records. Split them across four sub-pieces and you have four partial results and no step that combines them - the plan has no such step, because the plan was written on the assumption that one destination sees the whole key. In principle the runtime could insert the combining step, turning the grouping into two passes. In practice it usually will not, for two reasons: it must be able to prove that the aggregate decomposes that way, which it cannot do for an arbitrary function supplied by the author, and on some engines it is not permitted to rewrite that operator at all. Where the runtime declines, the two-pass rewrite becomes the author's job rather than the runtime's - a different subject with a different owner. ## The two cases side by side | | oversized piece holds many keys | oversized piece is one heavy key | |---|---|---| | grouping | can be cut on key boundaries; each key stays whole | cannot be cut; the aggregate needs every record for that key in one place | | join | can be cut on key boundaries | can be cut inside the key, if the other side's matching rows are copied to every sub-piece | | cost | none beyond scheduling more units | for the join, the copied side is read and sent once per sub-piece | ## What varies between engines Say this out loud in an interview rather than asserting a universal. Whether a runtime divides oversized pieces at all, whether it does so for joins only or for groupings too, whether it will insert a combining pass, and how far above its peers a piece must be before anything happens are all per-engine decisions - and in a job whose input never ends there is no finished step to measure, so none of it applies. The mechanism is a correction that may be available, not a property of distributed processing. ## The judgment this question is testing An interviewer asking this wants to hear that you know why one key resists division at all, rather than that you know a feature exists. The chain is: the key-to-destination rule puts one key on one destination, division inside a key is only sound when the operator lets you reconstruct the whole from the parts, a join lets you by copying the other side, and an aggregate does not without a combining pass nobody inserted. The corollary matters as much: when the runtime leaves your grouping alone, it has not failed - it has hit a boundary, and the remedy from there is the author's.

  • What if the oversized grouping piece holds a thousand medium keys rather than one heavy key?
    Then it can be divided. Different keys are independent, so re-cutting the range on key boundaries keeps every key whole and the results are complete per sub-piece. Whether a particular runtime bothers is another matter, but nothing about the operator forbids it. The case that defeats division is bulk concentrated in a single key.
  • What does dividing a join's piece actually cost?
    The other side's matching rows are read and transmitted once per sub-piece, so the replication factor is the number of sub-pieces. You are trading extra reads and network for parallelism, which pays only when that one piece dominates the step. If the key is heavy on both sides, the copies multiply and the trade stops being worth it.
  • Could the runtime rewrite the grouping into two passes by itself?
    Only for aggregates it can prove decompose, and only on engines permitted to rewrite that operator. An arbitrary author-supplied aggregate gives it nothing to reason about, so it leaves the grouping as planned. From that point the two-pass rewrite is the author's move, not the runtime's.

saying these in an interview costs you the question

  • Says the runtime spreads one key's records across two destinations for a grouping
  • Thinks copying the other side of a join changes the join's result
  • Assumes any oversized piece can be divided regardless of the operator
  • Assumes the runtime always knows the record count of each individual key
  • Believes dividing a join piece costs nothing extra