skip to content

How does appending an artificial suffix to a heavy grouping key spread its records, and what second pass does that force?

level: juniorimportance: must knowfreq 60%

answer

  1. one key value, one destination
  2. change the key, not the machines
  3. small integer suffix, several destinations
  4. second pass drops the suffix
  5. combine partial results per original key

basics

~20 s

Appending a small integer to the grouping key turns one key value into several, so its records reach several destinations instead of one. A second pass then re-groups those partial results under the original key to produce the real answer.

solid answer

~50 s

A grouping cannot be computed from what one worker already holds, so it first sends every record to the destination that owns its key - and the rule mapping a key to a destination is a function of the key alone. One value covering a large share of the data therefore lands whole on one destination, and extra machines cannot divide it. The rewrite makes several keys out of one. Pass one attaches a small integer suffix, random or round-robin, and groups on `(key, suffix)`, so the heavy value reaches up to `W` destinations, each holding about a `W`th of its records and emitting a partial result. Pass two drops the suffix and combines the `W` partials under the original key. Pass two is cheap because it reads partial results rather than records - which is also the condition for the rewrite to be worth anything.

code

sql · 18 lines
sql
-- pass one: group on the original key plus an artificial suffix
WITH partials AS (
  SELECT
    customer_key,
    spread_suffix,                -- integer 0..7 attached to each record upstream
    SUM(amount) AS part_sum,
    COUNT(*)    AS part_rows
  FROM events
  GROUP BY customer_key, spread_suffix
)
-- pass two: drop the suffix, combine the eight partials per original key
SELECT
  customer_key,
  SUM(part_sum)                             AS total_amount,
  SUM(part_rows)                            AS total_rows,
  SUM(part_sum) / NULLIF(SUM(part_rows), 0) AS mean_amount
FROM partials
GROUP BY customer_key;

go deeper

for a junior

Recall the shape: one key value always lands on one destination, so the fix is to make several key values out of it and then combine the partial results in a second pass.

for a middle

Explain that the composite key is the only thing that changed, that the second pass reads partial results rather than records, and that the job now redistributes twice where it redistributed once.

for a senior

Show judgment about where the suffix goes: only on values measured to be heavy, at a width taken from their record counts, and stable enough that a recomputed piece places records the same way.

for a principal

Weigh a hand-written rewrite that has to be re-tuned whenever the key distribution moves against simply paying for the slow run, and say which jobs earn that ongoing maintenance.

## Why one key cannot be divided A distributed job cuts each step's input into **pieces** - one share of the input that a single **worker** (one process on one machine, with its own memory and its own local disk) handles independently of the others. A step costs what its slowest piece costs, so the pieces are meant to be comparable. Some steps cannot be computed from what a worker already holds. Grouping by a key is the obvious one: the records for a key may start out anywhere, so the step begins with a **redistributing step** - every record is sent to the destination that owns its key. Which destination that is comes from the **key-to-destination rule**, a function of the key value alone. That is not an implementation detail; it is what makes the grouping correct, because equal keys have to meet somewhere. It is also the trap. A **heavy key** - one value covering a large share of all the records - resolves to exactly one destination, and that destination receives the whole share by itself. This is **data skew**: an uneven number of records per piece. It is a different thing from a **straggler**, where a piece of ordinary size runs long because the machine underneath it is unhealthy. Skew of this kind is immune to the obvious levers: more workers, more pieces and a larger cluster all leave the heavy value's records on one destination, and the run still costs what that destination costs. ## The rewrite Since the key cannot be split, split the key value: 1. **Attach a suffix.** Before the redistributing step, give each record of the heavy value a small integer `s` in `0..W-1` - drawn at random or handed out round-robin - and group on the composite `(original key, s)`. 2. **Aggregate the spread shares.** The key-to-destination rule now sees up to `W` distinct values where it saw one and scatters them, so each destination receives roughly a `W`th of the heavy value's records and emits a **partial result** for its share. 3. **Combine.** A second redistributing step drops the suffix, groups the partial results by the original key, and folds the `W` partials into one answer per key. Nothing about the engine changed; only the values it was asked to group by. The technique therefore reads the same on runtimes whose redistribution writes to local disk and is fetched afterwards and on those that push records across the network as they are produced. ## What the two passes cost | | pass one | pass two | |---|---|---| | reads | every input record | one partial per `(key, suffix)` | | the heavy value contributes | its full share of records | `W` rows | | groups maintained | distinct keys x `W`, where the suffix applies | distinct keys | | the busiest destination | about 1/`W` of the original load | a few rows per key | Pass two is cheap for one reason only: its input is partial results rather than records. That is also its precondition. If the partial is as large as the group it summarises - the list of every record under the key, say - then pass two rebuilds the heavy group on one worker and the imbalance comes straight back. ## The four things it does not do - **It does not remove the movement.** The job now performs two redistributing steps where it performed one. What it buys is balance, not less traffic: the first movement's destinations are comparable, so the step costs about what its median piece costs instead of what its worst one costs. - **It does not help a straggler.** A unit that read a normal amount of input and still ran ten times longer than its peers has a machine problem. Spreading a key does nothing for it, which is why the diagnosis has to come before the rewrite. - **It does not fit every aggregate unexamined.** Pass two has to rebuild the answer from partials, so the combining operation must be associative and have an identity - the general law of any parallel reduction. A sum, a count, a minimum and a maximum satisfy it directly; a mean does not, and is recovered by carrying a sum and a count and dividing once at the end, never by averaging the partial averages. - **It is not free for the keys that were fine.** Applied to every key, the intermediate result becomes distinct keys x `W` rows. On a job with tens of millions of distinct keys and a width in the hundreds that is billions of partials for pass two to move - considerably worse than the hour you set out to remove. In practice the suffix is attached only to the values measured to be heavy, and every other key carries a constant suffix so that it produces exactly one partial. ## Random or derived A suffix drawn at random per record is the simplest thing to write, but it makes pass one's output non-reproducible. On runtimes that recover a lost piece by recomputing it from its inputs, the recomputed piece can place records differently from the one it replaces. Deriving the suffix from a secondary field of the record, reduced into `0..W-1`, keeps the assignment stable across recomputation while still spreading the value - provided that secondary field is itself varied.

  • Why is the second pass usually far cheaper than the first?
    Because it reads one partial result per `(original key, suffix)` instead of raw records. Where the suffix is attached only to the values measured heavy, the heavy value arrives as `W` rows and every other key as one, so the second redistribution moves a small fraction of the bytes the first one moved.
  • Does the rewrite reduce the total amount of record movement?
    No - it usually increases it a little, because the job now performs two redistributing steps where it performed one. What it buys is balance: the destinations of the first movement hold comparable amounts, so the step costs roughly what its median piece costs rather than what its worst piece costs.
  • What goes wrong if the suffix is never dropped before the second grouping?
    You publish up to `W` rows per original key, each carrying a partial answer. A consumer expecting one row per key either sees the key duplicated or silently reads one partial as the total, so figures come out as a fraction of the truth and nothing in the job reports an error.

A supermarket sends every shopper whose surname starts with S to one till, because the till is chosen from the surname. Opening more tills changes nothing for them. Hand the S shoppers numbered tickets one to eight and open eight S tills, and the queue divides - but at closing time someone still has to add up the eight tills to get the day's S total. The tickets are the artificial suffix; the closing total is the second pass.

saying these in an interview costs you the question

  • Cutting the input into more pieces will split the heavy key
  • Spreading the key removes the redistributing step entirely
  • One pass is enough; group by the suffixed key and publish that
  • Average the partial averages to get the key's mean
  • Attach the suffix to every key; it costs nothing
  • Spreading a key also fixes a unit slowed by its machine