skip to content

A job groups records by account, then joins on account, then aggregates by account again. How can it be rewritten to regroup once?

level: middleimportance: should knowfreq 56%

answer

  1. one regrouping, then reuse its placement
  2. same key, same rule, same count
  3. opaque functions hide the key
  4. the unmatched join side still moves
  5. fewer records is not fewer regroupings

basics

~20 s

Regroup once on the account key, then keep that division. After the first wide step every account's records already sit together, so later steps keyed on the same account finish locally — unless something in between re-divides or hides the key.

solid answer

~60 s

The three steps name the same key, so in principle only the first has to move records. A **wide step** — one a worker cannot finish from the records it already holds — leaves behind a useful side effect: the data is now placed by the **division rule** you asked for, the function from a record's key to the number of the piece it belongs in. Any later step keyed on that same account, with the same rule and the same piece count, can then be satisfied from the piece in hand. The rewrite is therefore: order the keyed steps so they share one key, avoid re-dividing between them, avoid steps that may change the key, and remember that the second input to the join was never divided by account and still has to be. Whether the runtime *proves* the first division and skips the regrouping varies — on a declarative surface it often can, on opaque per-record code it cannot, and in the two-phase disk-to-disk model you fold the steps into one job yourself.

go deeper

for a junior

Remember the shape of the fix rather than its mechanics: steps that name the same key can often share one regrouping, so keeping them next to each other in the pipeline is worth doing before any tuning.

for a middle

State the four conditions — same key, same rule, same piece count, nothing invalidating in between — and name what breaks each. Be able to say why an opaque per-record function is the usual culprit.

for a senior

Separate removing a regrouping from shrinking one, and say which runtimes can prove a division and which cannot. Judge when merging pieces in place is the right trade against the imbalance it preserves.

for a principal

Across many jobs the win is conventional key choice: if teams key their datasets and their joins on the same identifier, regroupings are shared by default. That is a standards decision with a migration cost, not a per-job tweak.

## What a regrouping leaves behind A **wide step** is one a worker thread cannot finish from the piece of the input it already holds, because producing one output needs records currently sitting on every other worker. Executing one has a cost and also a product: afterwards, the records are no longer placed by the stored byte ranges they came from, they are placed by the **division rule** — the function from a record's key to the number of the piece that record belongs in. Everything with the same key is now in one place. That product is reusable, and reusing it is the rewrite. ## The conditions for reuse A later step is narrow on an already-divided input when all of these hold: 1. **Same key.** Grouping by account then joining on account qualifies; joining on region does not. 2. **Same division rule.** Two rules that both spread accounts evenly but differently put the same account in different places. 3. **Same piece count.** The rule maps a key to a piece number out of *n*; change *n* and every mapping changes. 4. **Nothing in between broke it.** This is the condition that actually fails in practice. ## What breaks a division - **An explicit re-division.** Asking for a different piece count between the steps throws the placement away and buys a second regrouping. - **A step that may change the key.** An arbitrary per-record function can rewrite the account field, and nothing outside it can see whether it did. A runtime that tracked the division conservatively must assume it is no longer valid. - **The other side of a join.** The second input was never placed by account. Only the side that already is can stay put; the other has to be brought into the same division. Some runtimes move only that side; others, unable to prove the first side's placement, redistribute both. - **A different key in the middle.** Grouping by account, then by region, then by account again means two of the three steps move records however you write them — unless you can reorder so the account-keyed steps are adjacent. ## The rewrites this buys you 1. **Reorder keyed steps so they are adjacent and share one key.** Three account-keyed steps in a row can share one regrouping; the same three separated by a region-keyed step cannot. 2. **Delete a re-division you added out of habit.** An explicit request to change the piece count before a step that was already narrow, or already placed by that key, is a pure cost. 3. **Merge pieces without moving records** when the only goal is fewer, larger pieces — glue neighbouring pieces together in place rather than re-dividing. It avoids a wide step, at the price of each surviving piece being correspondingly larger and of losing any evenness the re-division would have given. 4. **Scope a deduplication to a key you already have.** A global distinct is wide; distinct within an account, once records are placed by account, is local work. 5. **Keep the key visible.** Where the surface offers a way to express "this transformation preserves the key", use it; where it does not, leave the key-touching logic out of the opaque function. ## What this rewrite does *not* do It removes regroupings; it does not make the remaining one cheaper. Two changes are commonly confused with it and are different work: | change | effect on the wide step | why | |---|---|---| | moving a filter earlier | still wide, fewer records travel | the dependency shape is untouched | | raising the piece count | still wide | more, smaller pieces, same re-placement | | sharing one regrouping across three keyed steps | two wide steps become narrow | the placement already satisfies them | | merging pieces in place instead of re-dividing | the wide step disappears | no record changes worker | Writing both inputs to a join with the same key, rule and count *before* the job ever runs is a related but separate technique, and so is sending a small input to every worker rather than moving the large one; each belongs to its own subject. ## What varies by runtime This is where a confident-sounding claim goes wrong. Whether a second keyed step is *actually* executed as narrow depends on whether the runtime tracks the current division at all: - A declarative surface that models the division as a property of the data can often prove the placement and schedule the next keyed step without moving anything. - A pipeline of opaque per-record functions gives the runtime nothing to prove with; it will typically re-place the records. - In a continuous job the placement is established by the declared operator width, and keeping several keyed operators on the same key is how you avoid a second routing hop for every record. - In the two-phase disk-to-disk model there is one regrouping per job by construction, so "reuse" means folding the keyed logic into a single job rather than chaining jobs. So the advice is stable — same key, same rule, same count, nothing in between — while the payoff is something you confirm on the runtime in front of you rather than assume.

  • Why can an arbitrary per-record function between the two keyed steps cost you the second regrouping?
    Because nothing outside the function can see whether it rewrote the key. A runtime tracking the division has to assume the placement may no longer hold and re-establish it. Keeping key-touching logic out of opaque per-record code, or using whatever the surface offers to declare the key preserved, is what protects the reuse.
  • Does moving a filter before the grouped aggregate remove a wide step?
    No. It shrinks what the wide step has to carry, which is often worth doing, but the dependency shape is unchanged: the output for one key still needs records that may sit anywhere. Only a rewrite that lets a step be satisfied from the piece in hand removes the regrouping itself.
  • When is merging pieces in place a poor substitute for re-dividing?
    When evenness matters. Gluing neighbours together avoids a wide step but inherits whatever imbalance the original pieces had, and can concentrate it — one oversized piece then runs long after the rest finish. Re-dividing costs a regrouping and buys an even spread; which is right depends on how skewed the pieces already are.

saying these in an interview costs you the question

  • Says filtering earlier removes the regrouping rather than shrinking it.
  • Assumes any second step on the same key is automatically free.
  • Forgets the other join input was never placed by that key.
  • Ignores that a different piece count invalidates the placement.
  • Claims every runtime remembers how the data is currently divided.
  • Treats merging pieces in place and re-dividing as interchangeable.