Most MapReduce tuning advice arrives as a list of forty properties with recommended values. Copying such a list is how clusters end up with 8 GB heaps for tasks that need 600 MB, a thousand reducers for a job that emits a few gigabytes, and a job that got slower after it was "tuned". The values that are right for your job depend on how many bytes it reads, how many it emits from the map side, how those bytes are distributed across keys, and how much memory one record group needs.

This article gives you a method instead of a list. It explains where a MapReduce job spends its time, which few settings control each part, how to read the job's own counters to find the part that is actually slow, and how to change one thing at a time. The defaults quoted here come from Hadoop's current mapred-default.xml. The worked example takes one slow, failing aggregation job through four measured changes.

Tune from evidence, not from a property list

A MapReduce job is a pipeline of four stages: read input splits, run the mapper and sort its output locally, shuffle the sorted output across the network to reducers, then reduce and write. Each stage has a small number of settings that matter and many that do not. Wall-clock time is set by the slowest stage and, within a stage, by the slowest task, so the first question is never "what should the sort buffer be" but "which stage and which tasks are slow, and why".

The tools for answering that question already exist in every cluster: the job history server's task timeline, per-task counters and the container logs. A tuning session is a loop. Run the job on realistic input, record wall time and a handful of counters, form one hypothesis, change one setting, run again, and keep the change only if wall time or resource use improved without breaking anything else. Changing five settings at once makes the result impossible to explain, and an unexplained improvement tends to disappear when the data changes.

Input splitshow many map tasksMap + sort bufferspills to local diskShufflefetch over the networkReduce + writeone task per partitionsplit.minsize/maxsizeCombineTextInputFormatmap.memory.mbio.sort.mb, combinermap output codecparallelcopies, slowstartjob.reducesreduce.memory.mbcontrolscontrolscontrolscontrolsCounters tell you which box is slowSPILLED_RECORDS, MAP_OUTPUT_MATERIALIZED_BYTES, GC_TIME_MILLIS, REDUCE_INPUT_GROUPSChange one knob per run, compare counters and wall time, keep the change only if both improve.
The four stages of a MapReduce job, the settings that control each one, and the counters that show which stage is slow.

Containers and heaps are different limits

Every task runs in a YARN container with a memory limit, and inside the container runs a JVM with a heap limit. These are two different numbers, and confusing them causes the most common MapReduce failure. The container size is mapreduce.map.memory.mb or mapreduce.reduce.memory.mb. The heap is the -Xmx in mapreduce.map.java.opts or mapreduce.reduce.java.opts.

In Hadoop 3 both memory settings default to -1, which means "infer". If you give a heap size and no container size, the container is derived from the heap using mapreduce.job.heap.memory-mb.ratio (default 0.8, so the heap is 80% of the container). If you give a container size and no -Xmx, the heap is derived the other way. If you give neither, the container is 1024 MB. The safest habit is to set only the container size and let the ratio compute the heap, because the JVM always needs memory outside the heap for thread stacks, metaspace, the JIT, direct buffers and native compression libraries.

When a task's process tree exceeds its container, the NodeManager kills it and the attempt log says the container "is running beyond physical memory limits". When the heap fills, the JVM throws OutOfMemoryError instead. The two errors need opposite fixes: the first means the gap between heap and container is too small, the second means the heap itself is too small or the code holds too much. Also remember that YARN rounds every request up to a multiple of yarn.scheduler.minimum-allocation-mb, so asking for 1,100 MB on a cluster with a 1,024 MB minimum really costs 2,048 MB of capacity.

# Per-job settings: size the container, let the ratio derive the heap
hadoop jar agg.jar com.example.Agg \
  -D mapreduce.map.memory.mb=2048 \
  -D mapreduce.reduce.memory.mb=4096 \
  -D mapreduce.job.heap.memory-mb.ratio=0.8 \
  -D mapreduce.map.java.opts="-XX:+UseG1GC" \
  -D mapreduce.reduce.java.opts="-XX:+UseG1GC" \
  /data/events/2026/10 /out/agg/2026-10
# Map heap becomes about 1638 MB, reduce heap about 3277 MB.

How many map tasks: splits

The number of map tasks equals the number of input splits. For FileInputFormat the split size is max(minSize, min(maxSize, blockSize)), using mapreduce.input.fileinputformat.split.minsize (default 0) and mapreduce.input.fileinputformat.split.maxsize. With defaults, one split is one HDFS block, which is a sensible unit: a 128 MB block typically maps in tens of seconds, long enough to amortise container start-up but short enough that a retry is cheap.

Two situations break that. Many small files produce one split per file, so a directory of 200,000 files of 50 KB launches 200,000 tasks that spend most of their life starting JVMs. CombineTextInputFormat with a split.maxsize of a few hundred megabytes packs many files into each split. Non-splittable compression produces the opposite: a gzip file cannot be split, so a 20 GB gzip file is one map task no matter what you set. Store large inputs uncompressed, in a splittable container format, or as many moderately sized compressed files. The small files article covers the storage side, and the map phase article explains how splits become tasks.

How many reducers, and why skew ignores that number

The number of reducers is mapreduce.job.reduces, and its default is 1. A single reducer is correct for a tiny result and catastrophic for anything large, because every byte of map output funnels through one task. Too many reducers is also wasteful: each one is a container, each writes at least one output file, and thousands of tiny output files become the next job's small-files problem.

A workable starting point is to divide the job's total map output after any combiner (the MAP_OUTPUT_MATERIALIZED_BYTES counter from a run without map output compression, since the counter is measured after compression) by a target of roughly one to two gigabytes per reducer, then round to fit the reduce containers your queue can run at once. That target is a heuristic, not a rule; reducers that do heavy per-group work want less data each.

Reducer count cannot fix skew. The default HashPartitioner sends every record with the same key to the same reducer, so if one key carries 30% of the data, one reducer carries at least 30% of the work whatever the total count. You see it in the timeline as one reducer running long after the rest, with a REDUCE_INPUT_RECORDS counter far above its peers. The fixes are in the data design: salt hot keys and aggregate in two passes, or pre-aggregate with a combiner.

Moving fewer bytes

The cheapest byte is one you never shuffle. Three levers reduce map output. A combiner pre-aggregates on the map side when the reduce function is associative and commutative; compare COMBINE_INPUT_RECORDS with COMBINE_OUTPUT_RECORDS to see what it buys. Map-output compression (mapreduce.map.output.compress=true, default false) shrinks spill files and network transfer; choose a fast codec such as org.apache.hadoop.io.compress.SnappyCodec or Lz4Codec rather than the default DefaultCodec, which is zlib and costs much more CPU. Finally, emit less: drop unused fields in the mapper and use compact Writable types rather than Text holding JSON.

The sort buffer (mapreduce.task.io.sort.mb, default 100) decides how often a map task spills. The quick test is SPILLED_RECORDS divided by MAP_OUTPUT_RECORDS, both taken from the map tasks (the Map column on the history server's counter page), because the job total also counts reduce-side spills. Near 1 means one spill per record; above 1 means re-spilling in merge passes. Raise the buffer only if the heap has room, and read the shuffle and sort article for buffer sizing, the spill threshold and merge factors in detail.

Scheduling knobs

A few settings change when work runs rather than how much there is.

  • mapreduce.job.reduce.slowstart.completedmaps (default 0.05) starts reducers when 5% of maps finish so they can fetch while maps run. On a busy shared queue those early reducers sit in containers waiting for map output that other jobs could have used; 0.5 to 0.8 is often fairer.
  • mapreduce.reduce.shuffle.parallelcopies (default 5) is how many map outputs a reducer fetches at once. Raise it moderately for jobs with many thousands of maps and a fast network.
  • mapreduce.map.speculative and mapreduce.reduce.speculative (both default true) launch backup attempts for slow tasks. They help with a sick node and waste capacity when slowness comes from skew, because the backup has the same data. Disable speculation for tasks that write to external systems.
  • mapreduce.job.ubertask.enable (default false) runs a tiny job, up to 9 maps and one reducer by default, inside the ApplicationMaster's JVM, saving container launches for jobs that finish in seconds. MRv1's JVM reuse has no YARN equivalent; uber mode is the closest substitute.
  • mapreduce.task.timeout (default 600000 ms) kills a task that neither reads, writes nor reports progress for ten minutes. A reducer doing a long computation per group should call context.progress() rather than raising the timeout for every job.

Reading the counters

Counters are the measurement instrument. The job history web UI shows them per job and per task, and the command line can fetch any single one, which makes it easy to record a before and after.

JOB=job_1759480000000_4211
G=org.apache.hadoop.mapreduce.TaskCounter
for k in MAP_OUTPUT_RECORDS SPILLED_RECORDS MAP_OUTPUT_MATERIALIZED_BYTES \
         REDUCE_SHUFFLE_BYTES GC_TIME_MILLIS CPU_MILLISECONDS; do
  printf '%-32s %s\n' "$k" "$(mapred job -counter $JOB $G $k)"
done
SignalWhat it usually meansFirst lever
Spill ratio well above 1Sort buffer too small for the map outputCombiner, io.sort.mb, emit less
GC_TIME_MILLIS over ~10% of CPU_MILLISECONDSHeap too tight or too many live objectsContainer size, object reuse
Shuffle phase dominates reduce timeToo many bytes or too few fetchersMap output compression, combiner
One reducer far slower than the othersKey skewSalting, two-pass aggregation
Thousands of maps under 10 secondsSmall files or tiny splitsCombineTextInputFormat

Worked example: a nightly aggregation

Consider a daily job that counts events per user and per event type over 600 GB of uncompressed text in 128 MB blocks, about 4,800 splits. It was configured with 20 reducers, 2,048 MB reduce containers and -Xmx2048m in the reduce Java options. It runs for 3 hours 10 minutes, and several reduce attempts fail every night. The figures below illustrate the method; your ratios will differ.

  1. Failures first. The attempt logs say "running beyond physical memory limits". The heap equals the container, so any off-heap use crosses the line. Remove -Xmx, set mapreduce.reduce.memory.mb=3072 and let the 0.8 ratio give a heap of about 2,457 MB. The failures stop. Wall time barely moves, which is expected: this change fixed correctness, not speed.
  2. Bytes moved. Counters show 9.1 billion map output records, map-side SPILLED_RECORDS at 2.6 times that, and 610 GB of materialized map output: the job shuffles more than it reads. The reducer only sums counts, so add the reducer class as the combiner. Map output records stay the same, but materialized bytes fall to 140 GB because most (user, type) pairs repeat within a split.
  3. Compression. Enabling map output compression with Snappy reduces materialized bytes again to roughly 55 GB in this illustration. Map CPU rises slightly; shuffle time falls sharply.
  4. Reducers. Twenty reducers each now receive about 2.75 GB compressed. Using the heuristic on the 140 GB measured before compression suggests around 100 reducers, which also fits the queue. Raise mapreduce.job.reduces to 100 and set slowstart to 0.7 so reducers stop holding containers during the long map phase.

The job now finishes in about 55 minutes with no failed attempts and uses fewer container-hours than before. Notice what was not changed: the sort buffer, the merge factor, the number of fetchers. The counters never pointed at them, so changing them would have added risk without evidence.

Failure modes

  • Heap equal to container. Physical-memory kills that look random. Leave 20% or more for off-heap memory, more for jobs using native codecs or large direct buffers.
  • Huge heaps by default. Cluster-wide 8 GB task heaps cut parallelism in half and make GC pauses longer. Size per job from PHYSICAL_MEMORY_BYTES and GC time.
  • One giant gzip input. A single map task processes the whole file while the rest of the job waits. Recompress into splittable chunks.
  • Reducer count of 1 in production. The default survives from a test run. Always set mapreduce.job.reduces explicitly.
  • Combiner that is not associative. Averages computed in a combiner give wrong answers silently; the combiner article shows the correct pattern.
  • Heap pressure in the reducer. Code that buffers all values for a key in a list runs out of heap on hot keys. Stream through the values iterator or use a secondary sort.

Trade-offs

ChoiceYou gainYou pay
Bigger containersFewer OOMs, larger sort buffersFewer concurrent tasks, longer GC pauses
More reducersParallel reduce, smaller per-task memoryMore output files, more containers
Map output compressionLess disk and network I/OCPU on both sides of the shuffle
Late slowstartFairer use of a shared queueLess overlap of shuffle and map
Speculation onProtection from slow nodesDuplicate work, unsafe for external writes
Larger splitsFewer task launchesLonger retries, less parallelism

What to do next

  1. Pick your three most expensive jobs from the history server and record wall time, container-hours and the six counters in the script above for each.
  2. Search their configurations for -Xmx values equal to the container size and replace them with container-only settings.
  3. Compute the spill ratio and materialized bytes; add a combiner wherever the reduce function allows it, and enable Snappy or LZ4 map output compression.
  4. Set mapreduce.job.reduces from uncompressed materialized bytes, not habit, and check the slowest reducer's input records for skew.
  5. Look for jobs with thousands of very short maps and switch them to CombineTextInputFormat.
  6. Re-run with one change at a time, keep a short log of each change and its measured effect, and check the next quarter's runs to confirm the gains held. Pair this with daemon heap tuning for the services underneath.
Key takeaway: Tune MapReduce jobs by measuring which stage is slow and changing one setting at a time. Size containers and let mapreduce.job.heap.memory-mb.ratio derive the heap so off-heap memory has room, size splits to avoid both tiny tasks and unsplittable files, set the reducer count from materialized map output rather than the default of 1, and cut shuffled bytes with combiners and fast map output compression before touching sort-buffer internals. Counters such as SPILLED_RECORDS, MAP_OUTPUT_MATERIALIZED_BYTES and GC_TIME_MILLIS tell you which lever matters; skew needs a data-design fix, not more reducers.