skip to content

One reducer in a Hadoop MapReduce job runs for hours after every other reducer finishes — why?

level: seniorimportance: should knowfreq 45%

answer

  1. one task, all the rows
  2. the others finished; this one can't
  3. the partitioner had no choice
  4. a placeholder value nobody cleaned up
  5. more reducers will not split it

basics

~20 s

Almost always key skew. The default HashPartitioner sends every value for a given key to one reduce task, and a single key's group cannot be split across reducers, so one hot key pins one task. Speculative execution does not help, because the retry processes the same data.

solid answer

~50 s

The unit of parallelism on the reduce side is the **key group**, not the record. `HashPartitioner` assigns each key to `(key.hashCode() & Integer.MAX_VALUE) % numReduceTasks`, so all values for one key land on one reducer, and no amount of extra reducers splits them. When one key holds a large share of the data — a `NULL` or `"UNKNOWN"` sentinel, a default account id, a single enormous customer, a bot user — that reducer fetches, merges and reduces vastly more bytes than its peers. Confirm it by comparing per-task `REDUCE_INPUT_RECORDS` against `REDUCE_INPUT_GROUPS`: the slow task shows huge record counts over very few groups. Fixes, in order of preference: filter or special-case a sentinel key; add a combiner if the aggregation is algebraic; salt the hot key across N buckets and add a second aggregation pass; or replace the shuffle with a distributed-cache map-side join if the other side is small.

code

text · 5 lines
text
Task           REDUCE_INPUT_GROUPS   REDUCE_INPUT_RECORDS   REDUCE_SHUFFLE_BYTES   Elapsed
reduce_000000            41,208             1,905,442            212 MB     4m 11s
reduce_000001            39,884             1,877,003            209 MB     4m 02s
reduce_000002            40,551             1,893,770            211 MB     4m 08s
reduce_000003                 1           612,004,915             68 GB     3h 47m

go deeper

for a junior

Know that all values for one key go to one reduce task, so an unusually common key makes one task do far more work than the rest.

for a middle

Explain the mechanism precisely — HashPartitioner, the key group as the indivisible unit — and why adding reducers cannot split a single key. Be able to name the counters that confirm skew.

for a senior

Walk the full diagnosis: task-duration spread, per-task input records versus groups, identifying the hot key with a cheap profiling job, then choosing between filtering a sentinel, a combiner, salting with a second pass, or a distributed-cache map-side join.

for a principal

Treat recurring skew as a data-contract problem, not a job-tuning problem. Sentinel values, unbounded per-entity fan-out and unpartitioned backfills should be caught by upstream validation and key design, so pipelines are not repeatedly rescued by salting.

## Why one slow reducer is the default failure mode A MapReduce job's wall-clock time is the time of its slowest task. On the reduce side the work is partitioned by key hash, and the indivisible unit is a key group: `reduce()` receives one key and an iterator over *all* of its values, so those values must all arrive at the same task. That is a hard property of the model, not a tuning miss. Any input distribution where a few keys carry a large share of the records will therefore produce a straggling reducer, and the effect is superlinear — the hot task pays extra in fetching, in merging more sorted runs, in local disk pressure, and sometimes in garbage collection. ## Diagnosis Start in the ApplicationMaster's task list, sorted by elapsed time. If exactly one or two reduce tasks are outliers while the rest cluster tightly, you are looking at skew rather than a slow node. Confirm with per-task counters: - `REDUCE_INPUT_GROUPS` — how many distinct keys the task saw - `REDUCE_INPUT_RECORDS` — how many values it consumed - `REDUCE_SHUFFLE_BYTES` — how much it pulled across the network Skew looks like an enormous record count over a tiny group count. If instead the outlier has a normal record count but a long runtime, suspect the machine — a failing disk, a saturated NodeManager, a container squeezed on memory — and check whether a speculative attempt on a different node finished faster. To find *which* key, run a cheap profiling job: map the same key extraction, emit `(key, 1)`, aggregate, and look at the top of the distribution. In practice it is almost always a sentinel: an empty string, `NULL`, `"UNKNOWN"`, `-1`, `"default"`, an anonymous session id, or a date bucket that absorbed a backfill. ## Why the obvious fixes don't work **Raising `mapreduce.job.reduces` does nothing** for a single hot key. More reducers redistributes the *other* keys and leaves the hot one exactly where it was; you may even make things worse by producing more small output files. **Speculative execution** (`mapreduce.reduce.speculative`, on by default) also does not help: it launches a duplicate attempt of the slow task on another node, but that attempt must fetch and process the same enormous partition, so it will be slow too. Speculation is a cure for slow *machines*, not for uneven *data*. **Giving the reducer more memory** only postpones the problem, and it fails outright if the code materializes the value iterator into a list — a classic bug, since that iterator is a cursor over a merged run and was never meant to be collected. ## Fixes that work **Special-case or drop the sentinel.** If `"UNKNOWN"` is a placeholder with no business meaning, filter it in the mapper or route it to its own dedicated output. This is usually the cheapest and most honest fix, and it often reveals an upstream data-quality problem worth fixing at the source. **Add a combiner** when the aggregation is commutative and associative. Combining collapses each mapper's copies of the hot key to one record before the shuffle, which can turn a catastrophic skew into a harmless one. It does nothing for non-algebraic work such as collecting all values per key. **Salt the key.** Rewrite the hot key as `key#r` where `r` is a random value in `[0, N)`, so its rows spread across up to N reducers, then run a small second pass that strips the salt and merges the N partial results. This is a two-job pattern and the standard answer for aggregation skew. Salt only the known hot keys, so the second pass stays tiny. **Eliminate the shuffle for the skewed join.** If the skew comes from joining against a small dimension, ship the small side through the distributed cache, load it in the mapper's `setup()`, and join in `map()`. No key travels, so no reducer can be hot. This is the map-side join, and it is the same insight a broadcast join expresses in a modern engine. **Use a custom Partitioner** when skew is structural rather than single-key — for example when several known heavy keys hash into the same partition. A `Partitioner` that routes named heavy keys to dedicated reducers and hashes the rest evens out the load without a second pass. ## The honest framing for an interview Say what the model guarantees (all values for a key go to one task), show how you confirmed it from counters rather than guessing, explain why more reducers and speculation are not fixes, and then pick a remedy that matches the aggregation's algebra. That sequence — mechanism, evidence, ruled-out fixes, chosen fix — is what distinguishes a senior answer from a list of tricks.

  • Why doesn't speculative execution rescue a skewed reducer?
    Speculation launches a duplicate attempt of a slow task on another node and takes whichever finishes first. That helps when the *machine* is slow — a failing disk, a noisy neighbour — but a skewed task is slow because of its data, and the duplicate must fetch and reduce the same oversized partition. You end up doing the work twice at the same speed.
  • How does salting actually restore parallelism, and what does it cost?
    Appending a random suffix in `[0, N)` to the hot key turns one key group into up to N groups that hash to different reducers, so the heavy aggregation runs in parallel. The cost is a second job that strips the salt and merges the N partial results, plus the requirement that the aggregation be mergeable. Salt only the known hot keys so the second pass stays trivial.
  • What tells you the straggler is a bad node rather than skew?
    Check the counters. A skewed task shows a huge `REDUCE_INPUT_RECORDS` over a tiny `REDUCE_INPUT_GROUPS` count and a large `REDUCE_SHUFFLE_BYTES`. A bad node shows ordinary counters with a long elapsed time, and a speculative attempt elsewhere typically finishes quickly. Node symptoms also tend to affect map tasks on the same host.
  • Why is materializing the value iterator into a list dangerous?
    The `Iterable` handed to `reduce()` is a cursor over a merged, sorted run on disk, and Hadoop reuses the same value object as it advances. Copying every value into an `ArrayList` both breaks on object reuse unless you clone, and turns an unbounded key group into unbounded heap — so the hot key that merely made a task slow now makes it OOM.

Twenty supermarket checkouts open, nineteen empty, and one queue with a single customer whose trolley holds half the store's stock. Opening a twenty-first lane does not help, because that trolley cannot be split.

saying these in an interview costs you the question

  • Suggests raising the reduce count to split a single hot key
  • Says speculative execution will finish the slow task
  • Blames a slow node without checking per-task counters
  • Proposes only more reducer memory as the fix
  • Claims the framework splits a large key group across reducers

context