skip to content

A parallel aggregation groups records by key and assigns each key to a worker with hash(key) modulo the worker count. Most workers finish quickly while one runs for hours. What is happening, what options do you have, and what guarantees about output ordering does this style of processing give you?

level: seniorimportance: should knowfreq 42%

answer

  1. hash balances keys, not records
  2. a key is atomic - cannot be split across workers
  3. adding workers does not fix skew
  4. salt hot key -> aggregate -> merge partials (two stages)
  5. no order: between groups, within a group, or across runs

basics

~20 s

Key skew: one key (or a few hashing together) holds a huge share of the records, and a key cannot be split across workers, so its worker becomes the straggler. Options: pre-aggregate locally, salt the hot key into sub-keys and aggregate twice, or isolate hot keys. No ordering is guaranteed - not of values within a group, nor of outputs.

solid answer

~60 s

The job is **key-skewed**. Hash partitioning balances *keys*, not *records*: if one key holds 40 percent of the rows, its worker does 40 percent of the reduce work regardless of how many workers exist. Adding workers does not help, because the unit of assignment is a key. A secondary variant is hash collision, where several heavy keys land in the same partition. Remedies, in order: 1. **Pre-aggregate locally** before the exchange, so each worker ships one merged value per key. This alone often removes the problem. 2. **Salt** the hot key: emit `(key, random 0..n)` in a first pass, aggregate the sub-keys in parallel, then a second, tiny pass merges the n partials per key. Requires an associative merge. 3. **Isolate** known hot keys onto their own tasks, or use a range/custom partitioner informed by sampling. 4. **Re-key** so the grouping is finer-grained if the business question allows. Ordering: none. Values arrive at a group in arbitrary order, and groups are emitted in arbitrary order. If you need order, sort explicitly - do not rely on the partitioner or on the input order surviving.

code

text · 10 lines
text
// stage 1: spread the hot key over n sub-keys
map(record):
    k = record.key
    salt = (k in HOT) ? random(0, n-1) : 0
    emit((k, salt), lift(record.value))

reduce((k, salt), values) -> emit(k, merge_all(values))

// stage 2: merge the at-most-n partials per key
reduce(k, partials) -> emit(k, finish(merge_all(partials)))

go deeper

for a junior

Recognize that one key holding most of the records makes its worker the bottleneck, and that outputs come back in no particular order.

for a middle

Explain that hash partitioning balances keys not records, that a key cannot be split, and name pre-aggregation and salting as the fixes.

for a senior

Diagnose from per-task metrics, distinguish skew from hash collision and slow hosts, walk through two-stage salted aggregation and its associativity requirement, and be explicit about the three places ordering is not guaranteed.

for a principal

Treat skew as a data-distribution property to design for: sampled or custom partitioning, isolating sentinel keys as a data-quality issue, choosing keys and aggregations that stay mergeable, and setting explicit contracts about ordering with downstream consumers.

## What the partitioner actually balances Hash partitioning assigns key k to worker `hash(k) mod R`. Over many distinct keys this distributes *keys* very evenly. It says nothing about the number of *records* behind each key. Real data is almost never uniform in that respect: a `null` or empty-string key catches all the missing values, a default account absorbs unattributed events, one celebrity user generates a million rows, one country dominates the traffic. Record counts per key typically follow a heavy-tailed distribution, so a handful of keys hold a large fraction of the data. Because the model's contract is that one worker sees all values for a key, a key is **atomic** with respect to placement. The worker owning the hot key must process all of it. Job runtime is set by the slowest task, so the job takes as long as the hottest key - and adding workers changes nothing, which is the diagnostic signature. ## How to confirm it Look at the per-task record counts or bytes-processed distribution rather than averages. Skew looks like a long tail: median task processes a few million records, one processes a billion. Complementary evidence: the straggler's input is huge but its CPU profile looks normal; the job's runtime is insensitive to added parallelism; a `count(*) group by key order by count desc` on a sample immediately names the culprit. Distinguish it from two look-alikes: an unlucky **hash collision** (several distinct heavy keys mapping to the same partition - fixable by changing R or the hash), and a **slow machine** (same input size, longer time - fixable by speculative re-execution). ## Remedies **1. Local pre-aggregation.** If the aggregation has an associative, commutative merge, fold each worker's own records per key before the exchange. The hot key's payload shrinks from 'all its records' to 'one value per source worker'. This is the cheapest fix and often complete on its own. It fails when the operation is not mergeable - for instance, when the reducer must see raw records to sort or to emit them all. **2. Salting (two-stage aggregation).** Split the hot key artificially: ``` stage 1: emit ((key, rand(0..n-1)), value) -> n sub-groups, aggregated in parallel stage 2: emit (key, partial) from each sub-group -> merge n partials per key ``` Stage 1 spreads the load across n workers; stage 2 is tiny because it handles only n records per hot key. This requires the aggregation to be associative so partials can be merged. Salt only the hot keys if you can identify them - salting everything multiplies stage-2 volume for no benefit on cold keys. **3. Custom or sampled partitioning.** Instead of a plain hash, sample the input to estimate key weights and build a partition map that places heavy keys alone and packs light ones together. This is what range partitioners with sampled splitters do for sorting, and it is the general answer when skew is chronic rather than incidental. **4. Isolation and separate handling.** Route known hot keys (null, default, sentinel) to a dedicated path, or filter them out and handle them with a purpose-built query. Frequently the hot key is a data-quality artifact - a null placeholder - that should not be aggregated with real keys at all. **5. Re-keying.** If the grouping key is coarser than the question requires, refine it: group by `(country, day)` rather than `country`, and roll up afterwards. This turns one huge group into many medium ones. What does *not* work: adding workers, increasing memory (it defers the failure from a timeout to a spill, or from a spill to an out-of-memory), or shuffling harder. Those treat the symptom. ## Ordering: what you do and do not get This style of processing gives up ordering in three distinct places, and candidates routinely assume otherwise: - **Between groups.** Which key's output appears first depends on partition assignment and task completion. There is no global order unless you sort. - **Within a group.** Values reach a reducer in whatever order the exchange delivered them; the framework may group by sorting on the key, which orders keys within a partition but not the values behind a key. If the aggregation needs the values in order - a last-value-wins, a session reconstruction, a first-event timestamp - you must include the ordering field in the value and sort explicitly, or use a secondary-sort mechanism that sorts on a composite key. - **Across runs.** Because the grouping and merge order vary with worker count and scheduling, an order-sensitive or non-associative computation gives different answers on different runs. This is also why parallel floating-point sums are not bit-reproducible. The safe posture: assume the only ordering guarantee is the one you implement yourself, and make every aggregation order-independent (associative, commutative) so it cannot depend on something the framework never promised.

  • Why does doubling the number of workers not help a skewed aggregation?
    Because the unit of placement is a key, and one key's values must all be processed in one place. Doubling workers halves the load of the many light keys, which were never the bottleneck, while the hot key's task is unchanged. Runtime is set by the slowest task, so the job time barely moves - and that insensitivity is the clearest diagnostic that you are looking at skew rather than under-provisioning.
  • What does salting require of the aggregation, and what does it cost?
    It requires an associative merge, because the per-salt partial results must be combined in a second stage to reconstruct the key's true aggregate. The cost is a second pass over a small amount of data plus extra job complexity, and if you salt every key rather than only the hot ones you multiply the second stage's input by the salt factor for no benefit. Aggregations that need to see all raw records at once, such as an exact median over unbounded data, cannot be salted without changing the algorithm to a mergeable sketch.
  • A downstream consumer depends on results arriving grouped and in key order. What do you tell them?
    That this processing model guarantees no ordering - not between groups, not among the values within a group, and not stably across runs - so their dependency is on an accident of the current configuration. The correct fix is an explicit sort, either as a final ordered stage or by having the consumer sort or use an ordered store. Encoding the ordering field into the key or value and sorting explicitly is the only durable answer.

Sorting parcels by destination city across many desks: if half the parcels go to one city, that desk works all night no matter how many other desks you open.

saying these in an interview costs you the question

  • Fixing a skewed reduce by adding more workers or memory.
  • Assuming hash partitioning balances the number of records rather than the number of keys.
  • Believing a single key's values can be split across two reducers and still produce the right total.
  • Salting every key indiscriminately instead of only the identified hot ones.
  • Relying on values arriving at a group in input order, or on groups being emitted in key order.

context