MapReduce is two ideas that are easy to state and hard to implement well. The first is a programming model: express a computation as a map function that turns input records into key-value pairs and a reduce function that combines all values sharing a key. The second is an execution engine that runs thousands of copies of those functions next to the data, moves the intermediate pairs across the network, and survives the failure of any machine along the way.

Many teams now run Spark or Tez instead, but MapReduce still runs large Hive, Sqoop and DistCp workloads, and its vocabulary of splits, spills, shuffles, combiners and committers underlies every engine that followed. This article follows one job from the client to the committed output on a Hadoop 2 or 3 cluster running YARN, with the configuration defaults that shape each step, a worked sizing example, and the failure modes operators actually meet.

Advertisement

The programming model

A map function takes one input record, such as a byte offset and a line of text, and emits zero or more intermediate pairs (k2, v2). The framework groups every intermediate pair by key and sorts the keys. A reduce function receives one key with an iterator over all its values and emits zero or more output pairs. That is the entire contract; everything else is the engine's job.

The contract buys fault tolerance. Because a map task depends only on its input split and a reduce task only on map outputs, the engine can rerun either one anywhere, as often as needed, and get the same result, provided your functions are deterministic and free of external side effects. A map that calls a remote service or writes to a database breaks that guarantee: a retried or speculative attempt performs the side effect twice. Keep side effects in the output format, where the commit protocol controls them.

The moving parts on YARN

Since Hadoop 2, MapReduce is an application running on YARN rather than a cluster service. A job involves five components. The client computes input splits and stages the job. The YARN ResourceManager accepts the application and allocates containers; the YARN overview explains its schedulers and queues. The MRAppMaster, one per job, runs in the first container, plans tasks, requests containers for them and tracks every attempt; the ApplicationMaster article covers its heartbeat protocol in depth. NodeManagers launch task containers and host the ShuffleHandler, an auxiliary service that serves map output to reducers. HDFS holds input, staging files and output, and the JobHistory server keeps records after the AM exits.

ClientJob.submit()1 stageHDFS stagingjar, conf, splits2 submitResourceManagerschedules containers3 launch AMMRAppMasterplans tasks, tracks attempts4 containersmap tasks read splits from HDFS, data-local where possibleMap tasksplit 0Map tasksplit 1Map tasksplit 2Local disk: sorted, partitioned map output + indexserved by the NodeManager ShuffleHandlerfetchfetchReduce task 0merge, reduce()Reduce task 1merge, reduce()5 commitHDFS outputpart-r-00000..., _SUCCESS
One MapReduce job on YARN. The client stages files and submits; the ResourceManager starts the MRAppMaster; the AM runs map tasks near their splits; reducers fetch sorted partitions from each node's ShuffleHandler; the committer publishes output and writes _SUCCESS.
Advertisement

A complete job in code

The example counts error lines per service in tab-separated application logs. It uses the org.apache.hadoop.mapreduce API, the one to use for new code.

public class ErrorCount {
  public static class ParseMapper extends Mapper<LongWritable, Text, Text, LongWritable> {
    private static final LongWritable ONE = new LongWritable(1);
    private final Text service = new Text();                       // reused, not reallocated
    @Override
    protected void map(LongWritable offset, Text line, Context ctx)
        throws IOException, InterruptedException {
      String[] f = line.toString().split("\t");                    // ts, service, level, msg
      if (f.length < 4) { ctx.getCounter("parse", "malformed").increment(1); return; }
      if (!"ERROR".equals(f[2])) return;
      service.set(f[1]);
      ctx.write(service, ONE);
    }
  }

  public static class SumReducer extends Reducer<Text, LongWritable, Text, LongWritable> {
    private final LongWritable total = new LongWritable();
    @Override
    protected void reduce(Text key, Iterable<LongWritable> values, Context ctx)
        throws IOException, InterruptedException {
      long sum = 0;
      for (LongWritable v : values) sum += v.get();                 // values are streamed, not a list
      total.set(sum);
      ctx.write(key, total);
    }
  }

  public static void main(String[] args) throws Exception {
    Job job = Job.getInstance(new Configuration(), "error-count");
    job.setJarByClass(ErrorCount.class);
    job.setMapperClass(ParseMapper.class);
    job.setCombinerClass(SumReducer.class);                         // safe: sum is associative
    job.setReducerClass(SumReducer.class);
    job.setOutputKeyClass(Text.class);
    job.setOutputValueClass(LongWritable.class);
    job.setNumReduceTasks(4);
    FileInputFormat.addInputPath(job, new Path(args[0]));
    FileOutputFormat.setOutputPath(job, new Path(args[1]));         // must not exist yet
    System.exit(job.waitForCompletion(true) ? 0 : 1);
  }
}

Three habits in it are worth copying. Writable objects are reused rather than allocated per record, because a map task may process tens of millions of records. Malformed input increments a counter instead of throwing, so one bad line does not kill a task after four attempts, and the counter tells you how many lines were skipped. The combiner is the reducer itself, which is only valid because addition is associative and commutative; a combiner that computed an average this way would be wrong.

Submission and input splits

waitForCompletion hands the job to JobSubmitter. It checks that the output directory does not exist, so a rerun cannot silently overwrite earlier results. It asks the InputFormat for splits, then copies the job jar, the configuration and the serialized split list into a staging directory on HDFS, and submits a YARN application. The ResourceManager starts the MRAppMaster, which reads the splits and creates one map task per split plus mapreduce.job.reduces reduce tasks (default 1, which is almost never what a large job wants).

# FileInputFormat split sizing (per file, uncompressed or splittable codec)
split_size = max(min_split_size, min(max_split_size, block_size))
# defaults: min ~1 byte, max Long.MAX_VALUE, block 128 MB  ->  split_size = 128 MB
maps = sum(ceil(file_size / split_size) for each input file)   # small files: >= 1 map each
# the last chunk of a file may be up to 10% larger than split_size (SPLIT_SLOP = 1.1)

A split is a logical range, not a copy of data. It carries the hosts that store its blocks, and the AM asks YARN for containers on those hosts so most maps read locally. Record boundaries are handled by the RecordReader: a line reader skips the partial first line of its split and reads past its end to finish the last line. Non-splittable compressed files, such as a gzip file, become one map task each regardless of size, which is a common cause of a job with one very slow map. Many small files create many tiny maps; see the HDFS overview for why blocks and files cost NameNode memory. Very small jobs can run in uber mode, inside the AM's own container, when mapreduce.job.ubertask.enable is true and the job has at most 9 maps and 1 reduce by default.

The map side: buffer, sort and spill

Each map output pair is assigned a partition by the Partitioner, by default the key's hash modulo the number of reducers, and written into an in-memory circular buffer of mapreduce.task.io.sort.mb (default 100 MB). When the buffer reaches mapreduce.map.sort.spill.percent (default 0.80), a background thread sorts its contents by partition and then by key, runs the combiner if one is set, and writes a spill file to local disk while the map keeps filling the remaining space.

When the map finishes, its spill files are merged, up to mapreduce.task.io.sort.factor (default 10) streams at a time, into a single file sorted by partition and key, with an index recording where each partition starts. The combiner may run again during this merge. That file stays on the node's local disk, not in HDFS, which is why a lost node means its completed maps must be rerun. A job with zero reducers skips all of this and writes map output directly through the output format.

Shuffle and reduce

Reducers start before all maps finish: the AM schedules them once mapreduce.job.reduce.slowstart.completedmaps (default 0.05) of maps are done, so fetching overlaps the map phase. Each reducer pulls its partition from every completed map over HTTP from the ShuffleHandler, using mapreduce.reduce.shuffle.parallelcopies (default 5) fetch threads. Fetched segments are held in memory or spilled to disk and merged into one sorted stream. The reduce phase article goes deep on the copy and merge machinery, and the shuffle service article covers the NodeManager side.

The reducer then walks the sorted stream, calling reduce once per key group. Grouping uses a comparator that can differ from the sort comparator, which is how secondary sort works: sort by (user, timestamp) but group by user, so each call sees one user's events in time order. Output goes through the OutputFormat's record writer into a task-attempt directory, never straight into the final location.

Commit, retries and fault tolerance

Output is published by an OutputCommitter. With the default FileOutputCommitter, each attempt writes under a _temporary directory; a task's output is committed only for the attempt the AM chooses, and job commit moves committed output into the destination and writes an empty _SUCCESS marker. Downstream jobs should wait for that marker rather than for files to appear. The committer has two algorithm versions: version 1 moves task output at job commit, which is slower but leaves no partial output on failure; version 2 moves at task commit, which is faster but can leave partial results if the job fails. Check which one your distribution defaults to, and prefer version 1 when readers may see the directory. Object stores need committers designed for them, since rename is not atomic there.

Failure handling is attempt-based. A failed map or reduce attempt is retried up to mapreduce.map.maxattempts and mapreduce.reduce.maxattempts (both 4) times, preferably on another node, and nodes with repeated failures are blacklisted for the job. Reducers that cannot fetch a map's output report it, and the AM reruns that map. Speculative execution launches a duplicate of an unusually slow attempt and keeps whichever finishes first. If the AM itself dies, YARN restarts it up to mapreduce.am.max-attempts (default 2), subject also to the ResourceManager's own cap, and the new AM can recover completed tasks from its history.

Worked example: sizing a 1 TB log job

Suppose the error-count job reads 1 TB of uncompressed logs on a cluster with 128 MB blocks. The split formula gives 128 MB splits, so about 8,192 map tasks. With 400 map containers available, the maps run in roughly 21 waves; if each map takes 60 seconds, the map phase takes about 21 minutes plus scheduling overhead. Doubling the split size with mapreduce.input.fileinputformat.split.minsize would halve the task count and the per-task startup cost, at the price of less parallelism and slower recovery per failed task.

Now the shuffle. Assume 2 percent of lines are errors and each emitted pair is about 30 bytes. Each map would emit roughly 128 MB × 0.02 ≈ 2.6 MB of pairs without a combiner, about 21 GB across the job, all sorted and fetched over the network. With the combiner, each map emits one pair per service it saw; with 200 services that is a few kilobytes per map and a few tens of megabytes in total. The combiner turns a network-heavy job into a scan-heavy one. These figures are illustrative; read the real ones from the job counters for map output bytes, combine input and output records, and reduce shuffle bytes.

Finally the reducers. With 200 keys, 4 reducers is plenty, and more would only add output files. The risk is skew: if one service produces most errors, its reducer receives most of the values. Summing is cheap, so skew barely matters here, but a job that builds per-key lists would make that reducer the long pole.

Failure modes

  • Key skew. One reducer runs for hours while others finish in minutes. Look at per-task shuffle bytes; salt hot keys or pre-aggregate.
  • Too many small maps. Thousands of tiny files each cost a task. Combine them with CombineFileInputFormat or compact upstream.
  • Non-splittable input. A large gzip file becomes one map. Use a splittable format.
  • Reducer out of memory. Code copies values into a list; stream the iterator instead, remembering that Hadoop reuses the value object between iterations.
  • Side effects in tasks. Retries and speculation duplicate external writes.
  • Counter explosion. Dynamic counter names hit mapreduce.job.counters.max (default 120) and fail the job.
  • Reading output early. Consumers that do not wait for _SUCCESS read partial results.

When MapReduce still fits

Spark and Tez keep intermediate data in memory and plan multi-stage pipelines as one DAG, so iterative and multi-step jobs run much faster there. MapReduce writes to disk between every stage, which is slow but predictable: it needs little memory per task, recovers from failure at fine granularity, and has almost no tuning surprises. It remains a sensible choice for very large single-pass batch jobs on memory-constrained clusters, for tools built on it such as DistCp, and wherever operational simplicity matters more than latency.

What to do next

  1. Run the error-count job on a sample of your own logs and read every counter the job history page shows.
  2. Compute your largest job's split count with the formula above and check it against the number of maps actually launched.
  3. Compare map output bytes with reduce shuffle bytes to see whether a combiner would help.
  4. Set the reducer count deliberately for each job instead of inheriting the default of 1.
  5. Find the slowest reduce task in a recent job and compare its shuffle bytes with the median to detect skew.
  6. Confirm your committer algorithm version and make every downstream consumer wait for _SUCCESS.
Key takeaway: A MapReduce job is a pipeline of splits, buffered and sorted map output, an HTTP shuffle, grouped reduces and a committed output, coordinated by a per-job ApplicationMaster on YARN. Its fault tolerance rests on deterministic tasks and attempt-scoped output, so keep side effects out of your functions. Size splits and reducers deliberately, use a combiner when the reduce function allows it, and read the counters: they show where the bytes and the time actually go.