Between the last map() call and the first reduce() call, MapReduce does the work that gives the model its guarantee: every value for a given key ends up at one reducer, and that reducer sees its keys in sorted order with their values grouped. That phase is called shuffle and sort. It is usually the most expensive part of a job, because every intermediate byte is serialized, sorted, written to local disk at least once, sent over the network and merged again.
This article follows those bytes. It starts inside the map task, where records land in an in-memory collection buffer and are sorted and spilled, continues through the per-task output file and its index, crosses the network through the NodeManager's ShuffleHandler, and ends where the reduce side takes over. A worked example shows how buffer size changes disk traffic, and a counter table turns a slow job into a specific diagnosis. The full job flow is in MapReduce overview, in depth and the reduce side's fetchers and merge manager are covered in detail in MapReduce reduce phase, in depth. Configuration names below are the Hadoop 2 and 3 names from mapred-default.xml.
What the phase guarantees, and what it costs
The contract has three parts. Partitioning: each map output record is assigned to exactly one of R reduce partitions by the partitioner, by default a hash of the key modulo R. Sorting: within a partition, records are ordered by the job's sort comparator. Grouping: the reducer is called once per group of consecutive keys that the grouping comparator considers equal, which is how secondary sort works.
The cost is proportional to map output bytes, not input bytes. A job that filters 1 TB down to 10 GB of map output has a cheap shuffle; a join that emits every input row with a new key has a shuffle as large as its input, sometimes larger. Before tuning, compare map output bytes to map input bytes: it separates a machinery problem from a data-volume problem.
The collection buffer
When a mapper calls context.write(key, value), the framework's MapOutputBuffer computes the partition and serializes key and value straight into a byte array called the collection buffer. Its size is mapreduce.task.io.sort.mb, 100 MB by default, and it must fit inside the map task's JVM heap.
The buffer is circular and holds two kinds of data growing in opposite directions from a moving boundary called the equator: serialized key and value bytes, and fixed-size accounting metadata. The metadata for each record is four integers, 16 bytes, recording the partition, the key start, the value start and the value length. Sorting moves only these 16-byte entries, never the serialized records, which is why sort cost depends on record count as much as on bytes.
When the used portion reaches mapreduce.map.sort.spill.percent, 0.80 by default, a background spill thread starts sorting and writing the filled region while the mapper keeps writing into the remaining 20 percent. If the mapper fills that remainder before the spill finishes, context.write blocks. Small records pay a large metadata overhead, 16 bytes against a 20-byte record is 80 percent, and a mapper that outruns the disk stalls on the buffer.
Sorting: partition first, then key
The spill thread sorts metadata entries by partition number and then by key, using org.apache.hadoop.util.QuickSort unless map.sort.class says otherwise. Key comparison is the hot loop. If the key class registers a raw comparator that compares serialized bytes, no objects are created during the sort. If it does not, every comparison deserializes two keys. For custom keys, write the raw comparator:
public class EventKey implements WritableComparable<EventKey> {
long userId; long ts;
public void write(DataOutput out) throws IOException { out.writeLong(userId); out.writeLong(ts); }
public void readFields(DataInput in) throws IOException { userId = in.readLong(); ts = in.readLong(); }
public int compareTo(EventKey o) {
int c = Long.compare(userId, o.userId);
return c != 0 ? c : Long.compare(ts, o.ts);
}
/** Compares serialized bytes directly: no object allocation during sort and merge. */
public static class Raw extends WritableComparator {
public Raw() { super(EventKey.class); }
@Override public int compare(byte[] b1, int s1, int l1, byte[] b2, int s2, int l2) {
int c = Long.compare(readLong(b1, s1), readLong(b2, s2));
return c != 0 ? c : Long.compare(readLong(b1, s1 + 8), readLong(b2, s2 + 8));
}
}
static { WritableComparator.define(EventKey.class, new Raw()); }
}The raw comparator must agree exactly with compareTo, or map-side sort and reduce-side merge will disagree and groups will split. Test the two against each other on random keys.
Spill files and index records
Each spill produces one file in the task's local directories, laid out as R contiguous segments, one per partition, in IFile format: a stream of length-prefixed key and value records, optionally compressed per segment, with a checksum. Alongside it the task keeps an index with one record per partition: start offset, raw length and compressed part length, three longs, 24 bytes. The index lets any reader jump straight to partition r without scanning.
Spill indexes are held in memory up to mapreduce.task.index.cache.limit.bytes, 1 MB by default, and written to disk beyond that. With 24 bytes per partition, 1 MB covers roughly 43,000 partition entries, for example 43 spills of a job with 1,000 reducers. Jobs with many thousands of reducers therefore start writing spill indexes to disk earlier than people expect.
Compression applies here, not in the buffer. Setting mapreduce.map.output.compress to true with a fast codec such as Snappy or LZ4 in mapreduce.map.output.compress.codec shrinks every later step: spill writes, merge reads, disk usage on the NodeManager and network transfer. For text-heavy intermediate data it is usually the single cheapest improvement available.
Merging spills and where the combiner runs
If a map task spilled more than once, it must merge its spills into one final output file, file.out, with one final index, file.out.index. The merge is a k-way merge of already-sorted segments, done partition by partition, reading at most mapreduce.task.io.sort.factor streams at once, 10 by default. With more spills than the factor, intermediate merge rounds write and re-read data again.
A combiner, if configured, runs on each spill as it is written, and runs again during the final merge when the number of spills is at least mapreduce.map.combine.minspills, 3 by default. Because it may run zero, one or several times, a combiner must be commutative and associative and must produce output of the same types as its input. Sums, counts, minimums and maximums qualify; averages do not unless you carry a sum and a count.
Serving map output: the ShuffleHandler
Once a map task finishes, its output stays on the local disk of the node that ran it. Reducers pull it from the ShuffleHandler, a Netty-based HTTP server that runs inside every NodeManager as an auxiliary service, not inside the task. That is why a map's output survives the exit of its container, and why a NodeManager that is lost takes its map outputs with it.
<!-- yarn-site.xml on every NodeManager -->
<property>
<name>yarn.nodemanager.aux-services</name>
<value>mapreduce_shuffle</value>
</property>
<property>
<name>yarn.nodemanager.aux-services.mapreduce_shuffle.class</name>
<value>org.apache.hadoop.mapred.ShuffleHandler</value>
</property>
<!-- mapred-site.xml -->
<property>
<name>mapreduce.shuffle.port</name>
<value>13562</value>
</property>A reducer fetcher requests one or more map outputs for its partition in a single HTTP request, authenticated with a job token. The handler looks up the partition's offset in the index, using an index cache of its own, and streams that byte range from file.out, with zero-copy transfer where the platform allows it. If the aux service is missing, jobs fail at launch with an error that the mapreduce_shuffle auxService does not exist. Spark's equivalent, the external shuffle service, solves the same problem of outliving executors; External shuffle service compares the two designs.
The reduce side in one paragraph
Reduce tasks start before all maps finish, once mapreduce.job.reduce.slowstart.completedmaps of maps are done, 0.05 by default, so they can fetch early outputs while later maps run. Each reducer runs mapreduce.reduce.shuffle.parallelcopies fetcher threads, 5 by default, into a shuffle buffer that is mapreduce.reduce.shuffle.input.buffer.percent of heap, 0.70 by default. A single map output larger than mapreduce.reduce.shuffle.memory.limit.percent of that buffer, 0.25, goes straight to disk. In-memory outputs are merged to disk when usage passes mapreduce.reduce.shuffle.merge.percent, 0.66. The detail is in the reduce-phase article linked above.
Worked example: sizing the sort buffer
A map task reads a 256 MB block and emits 4 million records averaging 100 serialized bytes, 400 MB of map output. Each record also costs 16 bytes of metadata, so the buffer needs 116 bytes per record, 464 MB for the whole output.
With the default 100 MB buffer, a spill starts at 80 MB, which holds about 690,000 records. The task spills 6 times, then merges 6 files in one pass because 6 is below the merge factor of 10. Data is written twice, about 800 MB, and read back once. The Spilled Records counter reads 8 million, twice Map output records.
Raising mapreduce.task.io.sort.mb to 600 MB gives 480 MB usable per spill, enough for the full 464 MB. The task spills once, skips the merge, and Spilled Records equals map output records. The buffer lives in the map heap, so mapreduce.map.java.opts must rise to around -Xmx900m or more to leave room for the mapper's own objects, and mapreduce.map.memory.mb must rise to cover heap plus JVM overhead, for example 1,280 to 1,536 MB. Cheaper alternatives are map output compression, which does not shrink the buffer but halves every disk and network byte, and a combiner, which reduces the records spilled.
Reading a slow shuffle from its counters
| Counter or signal | What it tells you | Action |
|---|---|---|
| Spilled Records / Map output records well above 1.0 | multiple spills and merge passes per map | raise io.sort.mb with heap, add a combiner, compress map output |
| Map output materialized bytes far below Map output bytes | compression is working | keep it; compare codecs on CPU cost |
| Reduce shuffle bytes much higher for a few reducers | key skew; one partition carries most of the data | salt hot keys, pre-aggregate, or use a custom partitioner |
| Shuffle Errors: CONNECTION, IO_ERROR, WRONG_LENGTH | fetch failures from specific nodes | check the NodeManager, its disks and network; failures reported to the AM cause maps to re-run |
| Long shuffle time with low network use | too many tiny fetches or too few fetchers | fewer reducers, larger map outputs, raise parallelcopies moderately |
| GC time elapsed high in reducers during shuffle | shuffle buffer too large for the heap | lower input.buffer.percent or raise reduce heap |
| Local disk full on NodeManagers | intermediate data larger than local-dirs capacity | compress, add disks to local-dirs, reduce concurrent containers |
Fetch failures deserve a specific note. A reducer that cannot fetch an output retries with back-off and reports the failure to the MapReduce application master. When enough reports accumulate against one map, the AM declares that map's output lost and re-runs the map elsewhere. A single bad disk can therefore make an entire job repeat its map phase piecemeal. Look at which host the failures name before tuning anything, and see speculative execution for how stragglers interact with this.
Failure modes
- Mapper stalls on write. The spill thread cannot keep up with output; the task log shows the collector waiting for the spill. Fix disk throughput, compress output or reduce record count with a combiner.
- Map task out of memory at start. io.sort.mb was raised without raising the heap. The buffer is allocated up front, so the failure appears immediately.
- Split groups at the reducer. A raw comparator disagrees with compareTo, or the grouping comparator is inconsistent with the sort comparator.
- Maps re-run late in the job. Fetch failures against a lost or unhealthy NodeManager; the outputs lived only on that node.
- One reducer runs for hours. Skew. No shuffle setting fixes it; the key distribution must change.
Trade-offs
A bigger sort buffer trades heap for disk passes and is worth it only when Spilled Records shows extra spills. More reducers make each partition smaller and each merge cheaper, but multiply the number of fetches, M maps times R reducers, and the number of output files. Compression trades CPU for disk and network, and with modern codecs the trade almost always favours compressing. A combiner is nearly free when the aggregation is associative. Beyond these, the lever that matters most is sending less data: filter early, project only needed fields, and choose keys that spread evenly.
What to do next
- Pull the counters of your slowest recurring job and compute Spilled Records divided by Map output records.
- Turn on mapreduce.map.output.compress with Snappy or LZ4 and compare Map output materialized bytes and shuffle time.
- If the ratio is above 1.0, size io.sort.mb from records times (average record size plus 16 bytes), then raise the map heap and container memory together.
- Add a combiner wherever the reduce function is associative and commutative.
- Register raw comparators for custom keys and test them against compareTo on random data.
- Compare Reduce shuffle bytes across reducers to detect skew before touching fetcher settings.
- Confirm mapreduce_shuffle is configured on every NodeManager and alert on Shuffle Errors per host.