The reduce phase is where a MapReduce job pays for its shuffle. Every map task leaves a sorted, partitioned file on local disk; every reduce task must pull its partition from every one of those maps, merge thousands of sorted runs into one stream, call your reduce method once per key group, and commit its output atomically so that a failed or speculative attempt never leaves partial files behind. Most slow or failing MapReduce jobs are slow or failing here.
This article walks the reduce task from the inside: how fetchers copy map output, how the merge manager decides between memory and disk, what the final merge does, how the Reducer API iterates values, how secondary sort works, and how output is committed. It uses the property names and defaults from the current Hadoop mapred-default.xml, and ends with a sizing worked example and a checklist. The serving side of the shuffle, the NodeManager's shuffle service, is covered in external shuffle service.
The three sub-phases of a reduce task
A reduce task reports progress in three parts, and the job history UI shows time spent in each: copy (also called shuffle), sort (really a merge, because inputs are already sorted), and reduce. They overlap more than the names suggest. Merging starts during the copy as memory fills, and the last merge pass streams directly into your reducer rather than writing a fully sorted file first.
Reducers do not wait for all maps to finish before starting. mapreduce.job.reduce.slowstart.completedmaps (default 0.05) lets the ApplicationMaster launch reducers once 5% of maps are done, so copying overlaps with the map phase. That is a good default for a dedicated cluster and a poor one for a busy shared cluster, where early reducers hold containers while doing nothing but waiting for slow maps. Raising it to 0.8 or higher on shared clusters is common.
Copy: fetchers and the shuffle handler
The default shuffle consumer is org.apache.hadoop.mapreduce.task.reduce.Shuffle, configurable through mapreduce.job.reduce.shuffle.consumer.plugin.class. It runs mapreduce.reduce.shuffle.parallelcopies fetcher threads (default 5). The reducer learns about completed maps through task-completion events from the ApplicationMaster, and each fetcher asks a host's shuffle handler, the mapreduce_shuffle auxiliary service on the NodeManager (port mapreduce.shuffle.port, default 13562), for this reducer's partition of one or more map outputs on that host. Fetching is grouped by host, so a node with many finished maps is served in one connection.
Connect and read timeouts both default to 180,000 ms. When a fetch fails, the reducer retries, and reports the failure to the ApplicationMaster. If enough reducers report that a given map's output cannot be fetched, the ApplicationMaster declares that map attempt failed and re-runs it elsewhere, which is why a single bad disk or dead NodeManager late in a job can re-run maps that finished hours ago. If the reducer itself cannot make progress, it fails, and it is retried up to mapreduce.reduce.maxattempts (default 4) times before the job fails. When NodeManager recovery is enabled, mapreduce.reduce.shuffle.fetch.retry.enabled follows it by default, letting fetchers ride out a NodeManager restart instead of failing.
Merge: memory, disk and the MergeManager
Each fetched segment goes either into memory or straight to disk. The rules are fractions of the reducer's heap:
| Property | Default | Meaning |
|---|---|---|
mapreduce.reduce.shuffle.input.buffer.percent | 0.70 | Share of heap for holding fetched map outputs in memory |
mapreduce.reduce.shuffle.memory.limit.percent | 0.25 | Largest single segment kept in memory, as a share of that buffer |
mapreduce.reduce.shuffle.merge.percent | 0.66 | Buffer fill level that triggers an in-memory merge to disk |
mapreduce.reduce.merge.inmem.threshold | 1000 | Segment count that also triggers that merge |
mapreduce.task.io.sort.factor | 10 | Maximum segments merged at once on disk |
mapreduce.reduce.input.buffer.percent | 0.0 | Share of heap that may still hold map outputs when reduce() starts |
A segment larger than the single-segment limit bypasses memory and is written directly to disk. Smaller segments accumulate in the buffer until it passes the merge threshold or the segment count passes 1,000, at which point an in-memory merger writes them out as one sorted file. Meanwhile an on-disk merger combines disk files, at most io.sort.factor at a time, whenever too many accumulate. If the job has a combiner, it also runs during the in-memory merge, shrinking what reaches disk.
When copying ends, the final merge builds a single sorted iterator over the remaining in-memory segments and disk files. With mapreduce.reduce.input.buffer.percent at its default of 0.0, all in-memory data is flushed to disk before reduce() starts, freeing the whole heap for your code. If your reducer uses little memory, raising it (to 0.5, say) lets some map output stay in memory and saves a disk round-trip; if your reducer builds large in-memory structures, leave it at zero.
The Reducer API and the iterator trap
The reducer itself is small. The framework calls setup once, reduce once per key group with an iterable over that group's values, and cleanup once. The default run method ties them together, and you can override it for custom control, for example to stop early.
public class MaxTempReducer extends Reducer<Text, IntWritable, Text, IntWritable> {
private final IntWritable out = new IntWritable();
@Override
protected void reduce(Text station, Iterable<IntWritable> temps, Context ctx)
throws IOException, InterruptedException {
int max = Integer.MIN_VALUE;
for (IntWritable t : temps) {
max = Math.max(max, t.get()); // read the value, do not keep the object
}
out.set(max);
ctx.write(station, out);
}
}The trap is object reuse. The framework deserialises every value into the same Writable instance, and the key object is reused too. Code that stores the objects, for example by adding each IntWritable to a list, ends up with a list of references to one object holding the last value. Copy what you keep: new IntWritable(t.get()), or WritableUtils.clone. The iterable can also be traversed only once; a second loop sees nothing. If you need two passes, buffer the values yourself and watch the heap, or restructure the job so the values arrive in the order you need.
Grouping, sorting and secondary sort
Two comparators control the reduce side. The sort comparator (job.setSortComparatorClass) decides the order of keys in the merged stream. The grouping comparator (job.setGroupingComparatorClass) decides which adjacent keys belong to one reduce call. By default both are the key's own ordering.
Secondary sort uses the gap between them. To see each station's readings in time order without buffering, make the key a composite of station and timestamp, sort on both, partition and group on station alone. Each reduce call then receives one station's values already ordered by timestamp, and the current key object changes as you iterate, exposing each value's timestamp. The partitioner must agree with the grouping comparator; if it hashes the full composite key, one station's records land on different reducers and the grouping silently splits.
job.setMapOutputKeyClass(StationTime.class); // (station, timestamp)
job.setPartitionerClass(StationPartitioner.class); // hash(station) only
job.setSortComparatorClass(StationTimeComparator.class); // station, then timestamp
job.setGroupingComparatorClass(StationComparator.class); // station only
job.setNumReduceTasks(200);
Output: RecordWriter and the commit protocol
Your context.write calls go to a RecordWriter from the job's OutputFormat. With file output formats the writer does not write to the final directory. It writes into a task-attempt directory under _temporary, and the OutputCommitter decides what becomes visible. When the attempt succeeds, commitTask promotes its files; when the job succeeds, commitJob finishes the promotion and writes the _SUCCESS marker. Attempts that fail or lose a speculative race are aborted and their directories deleted.
FileOutputCommitter has two algorithms, selected by mapreduce.fileoutputcommitter.algorithm.version. Version 1 renames task output into a job-level temporary directory and renames everything again in commitJob, which is serial and slow for jobs with many files but leaves the destination untouched if the job fails before commit. Version 2 renames task output straight into the destination at task commit, which makes job commit fast but means a failed job can leave partial output visible. The current mapred-default.xml lists 2; older releases and several distributions ship 1, so check your cluster rather than assuming. Downstream readers should key on _SUCCESS, not on the existence of part files. On object stores where rename is a copy, use the store-specific committers instead of either algorithm.
Worked example: sizing a reduce stage
A job emits 2 TB of map output after the combiner, and each reducer container has a 4 GB heap. How many reducers?
Start from the data each reducer should handle. A common target is one to a few gigabytes of shuffled input per reducer, so the task finishes in minutes and a retry is cheap. At 2 GB each, 2 TB needs about 1,000 reducers. The shuffle buffer is 0.70 of 4 GB, about 2.8 GB, and the largest in-memory segment is 0.25 of that, about 700 MB. With 5,000 maps, each map sends this reducer about 2 GB / 5,000, roughly 400 KB, so every segment fits in memory; the buffer passes its 0.66 merge threshold after about 1.85 GB, so expect one or two in-memory merges to disk and a final merge over a handful of files, well under the sort factor of 10.
Now inspect the counters after a run. REDUCE_SHUFFLE_BYTES per task tells you whether the partition is even; REDUCE_INPUT_GROUPS against REDUCE_INPUT_RECORDS shows the values per key; SPILLED_RECORDS far above input records means multiple merge passes; FAILED_SHUFFLE above zero points to a sick node. If one reducer receives 60 GB while the median is 2 GB, adding reducers will not help: a single hot key is sent to a single reducer. Salt the key, pre-aggregate with a combiner, or handle the hot key in a separate pass. Spark faces the same skew problem with a different mechanism, described in Spark shuffle architecture.
Failure modes
- OutOfMemoryError during shuffle. Too much heap promised to the shuffle buffer, or many fetchers each holding a segment. Lower
input.buffer.percentormemory.limit.percent, or raise the container heap. - OutOfMemoryError in reduce(). Your code buffers a whole key group. Stream the values, or use secondary sort so you do not need to.
- Map re-runs late in the job. Fetch failures against one host; check that node's disks and NodeManager logs.
- One reducer runs for hours. Key skew. Speculative execution (
mapreduce.reduce.speculative, default true) cannot help, because the duplicate gets the same data; see speculative execution. - Wrong aggregates, no errors. Stored references to reused Writables, or a partitioner that disagrees with the grouping comparator.
- Idle reducers blocking other jobs. Slow start set too low on a shared cluster.
- Partial output visible. Algorithm version 2 and a job failure; readers ignored
_SUCCESS.
Tuning and trade-offs
Tuning the reduce phase is mostly about moving work away from it. A combiner and compressed map output shrink shuffled bytes, which cuts both network and merge time. Raising parallelcopies helps when there are many small maps on many hosts, at the cost of more concurrent connections per shuffle handler. Raising io.sort.factor to 50 or 100 reduces merge passes, at the cost of more open files and smaller read buffers per file. mapreduce.reduce.memory.mb defaults to -1, meaning it is derived from the heap setting; set container size and -Xmx together so the heap stays inside the container limit that YARN enforces.
More reducers mean smaller tasks, cheaper retries and more output files; fewer reducers mean larger files and longer tails. If downstream consumers suffer from small files, fix the reducer count rather than adding a compaction job later.
What to do next
- Open the job history for your slowest job and compare copy, sort and reduce times per reducer.
- Check
REDUCE_SHUFFLE_BYTESacross reducers; if the maximum is several times the median, fix the key skew before touching memory settings. - Set the reducer count from shuffled bytes, aiming for one to a few gigabytes per reducer.
- Audit reducers for stored Writable references and double iteration.
- Raise slow start on shared clusters, and enable map output compression and a combiner where the reduce function allows it.
- Confirm which committer algorithm your cluster uses and make downstream readers wait for
_SUCCESS.