skip to content

MapReduce Programming

The map/shuffle/sort/reduce model: deliberately restrictive, and the ancestor of every distributed engine that followed. Interviewers use it to test that you can decompose a problem into map and reduce steps and can explain what a combiner or a custom partitioner actually buys you.

questions

6

What phases does a Hadoop MapReduce job move through from input split to final output?

level: juniorimportance: must knowfreq 55%

answer

  1. four phases, one pipeline
  2. one map task per input split
  3. the sort happens before the network hop
  4. a partitioner picks the reducer
  5. output lands as part-r files

basics

~20 s

A Hadoop MapReduce job runs map, then shuffle and sort, then reduce. One map task handles each InputSplit; its output is partitioned and sorted on local disk, fetched by reducers over the network, merged by key, and written out by the OutputFormat.

solid answer

~40 s

The `InputFormat` divides the input into **InputSplits** and the framework starts one map task per split; a `RecordReader` turns the split into key/value pairs for `map()`. Each map task buffers its output in memory and, on spill, assigns every record a partition (by default `HashPartitioner` on the key modulo the reduce count) and sorts by key, writing sorted spill files to the mapper's **local** disk that are then merged into one partitioned file. In the **shuffle** (copy) phase every reduce task fetches its own partition from every map task, then merge-sorts the pieces so all values for a key are adjacent. `reduce()` is called once per key with an iterator over its values, and the `OutputFormat` writes results to HDFS as `part-r-00000`, one file per reduce task, plus a `_SUCCESS` marker.

code

java · 25 lines
java
public class WordCount {

  public static class TokenizerMapper extends Mapper<Object, Text, Text, IntWritable> {
    private final IntWritable one = new IntWritable(1);
    private final Text word = new Text();

    public void map(Object key, Text value, Context context) throws IOException, InterruptedException {
      for (String token : value.toString().split("\\s+")) {
        word.set(token);
        context.write(word, one);
      }
    }
  }

  public static class SumReducer extends Reducer<Text, IntWritable, Text, IntWritable> {
    public void reduce(Text key, Iterable<IntWritable> values, Context context)
        throws IOException, InterruptedException {
      int sum = 0;
      for (IntWritable v : values) {
        sum += v.get();
      }
      context.write(key, new IntWritable(sum));
    }
  }
}

go deeper

for a junior

Be ready to name the phases in order and say what each one does in a sentence, using word count as the worked example. Knowing that one map task handles one input split is the single most-expected fact here.

for a middle

Explain the mechanics between the phases: the in-memory buffer and its spills, why records are sorted before they cross the network, and what the partitioner decides. Be precise that intermediate output lives on the mapper's local disk, not HDFS.

for a senior

Show where real jobs spend their time. Interviewers expect you to point at the copy and merge phases as the bottleneck, connect that to output volume per mapper, and explain why the mandatory sort is what makes task retries and speculative attempts safe.

for a principal

Own the framing question: what the rigid map/shuffle/sort/reduce contract buys in fault tolerance, and what it costs in latency and expressiveness once a workload becomes multi-pass. That tradeoff is the argument every later engine was designed around.

## What a MapReduce job is made of A MapReduce job is one distributed batch computation submitted to a YARN cluster. YARN starts an ApplicationMaster for the job, which requests containers and runs two kinds of **tasks** inside them: map tasks and reduce tasks. The vocabulary matters, because it is shared across engines and means different things: in MapReduce a *task* is a whole mapper or a whole reducer, whereas in Spark a task is one partition's work inside a stage and in Flink a subtask is one parallel instance of an operator. Do not carry one engine's meaning into another. ## Phase 1 — input and map The job's `InputFormat` computes **InputSplits**: logical byte ranges over the input, each carrying hints about which hosts hold that data so YARN can place the container near it. The framework starts exactly one map task per split, and a `RecordReader` turns the split's bytes into key/value pairs. Your `map(K1, V1)` is then called once per record and emits zero or more intermediate `(K2, V2)` pairs. Map tasks are independent — no map task ever sees another's data. ## Phase 2 — partition, sort and spill (still on the map side) Intermediate output does not go straight anywhere. Each map task writes into an in-memory circular buffer sized by `mapreduce.task.io.sort.mb` (100 MB by default). When the buffer passes `mapreduce.map.sort.spill.percent` (0.80), a background thread spills it. Before writing, each record is assigned a partition by the `Partitioner` — the default `HashPartitioner` computes `(key.hashCode() & Integer.MAX_VALUE) % numReduceTasks` — and the buffer is sorted by partition, then by key. Each spill is therefore a sorted, partitioned file on the mapper's **local disk**, never in HDFS. At the end of the map task the spills are merged into one partitioned sorted file plus an index. If a combiner is configured, it may run over the spills and over the merge. ## Phase 3 — shuffle (the copy) This is the only phase that moves data across the network, and it is where most jobs spend most of their time. Each reduce task fetches its own partition from every map task's local output, served over HTTP by an auxiliary service inside the NodeManager. Copying can begin before all mappers finish — `mapreduce.job.reduce.slowstart.completedmaps` defaults to 0.05, so fetchers start once 5% of maps have completed — and `mapreduce.reduce.shuffle.parallelcopies` (default 5) controls the fetcher threads per reducer. ## Phase 4 — merge and sort on the reduce side The fetched segments are merged, partly in memory and partly on disk, into a single stream sorted by key so that all values for one key are adjacent. The sort is not an optional nicety: grouping is *implemented* by sorting, which is why MapReduce is often described as map/shuffle/sort/reduce rather than just map/reduce. Crucially, `reduce()` cannot begin until every map task has finished, because a straggling mapper may still be holding values for a key the reducer is about to close. ## Phase 5 — reduce and output `reduce(K2, Iterable<V2>)` is invoked once per key, in key order, with an iterator over that key's values. That iterator is a cursor over a merged run, not a materialized list — you cannot iterate it twice, and holding all of its values in memory is how reducers OOM on a hot key. Output goes through the `OutputFormat`'s `RecordWriter`; an `OutputCommitter` promotes each successful task attempt's temporary directory to the final path, which is what makes a speculative or retried attempt safe. The result is one file per reduce task — `part-r-00000`, `part-r-00001`, … — plus an empty `_SUCCESS` marker. ## What is ordered and what is not Each reducer's output file is sorted by key, but there is no global order across output files: keys are distributed by hash, so `part-r-00000` holds an arbitrary subset. Total order requires `TotalOrderPartitioner` driven by a sampled partition file, which is how a terasort-style job produces globally sorted output. ## Why the shape is so rigid The restrictions — no communication between mappers, a mandatory sort, materialized intermediate output on local disk — are what buy fault tolerance. Any task can be killed and re-run from its input, and a failed reducer simply re-fetches its partitions, because the mapper outputs are still sitting on disk. Later engines relaxed the rigidity and kept the idea; understanding these five phases is what makes a Spark stage boundary or a Flink keyBy legible.

  • Where does the sorting actually happen — on the map side or the reduce side?
    Both. Each map task sorts its buffer by partition and key before every spill, and merges its spills into one sorted, partitioned file. Each reduce task then merge-sorts the sorted runs it fetched from all mappers. Sorting on the map side is what makes the reduce-side merge cheap, and it is why grouping and ordering by key are the same operation in MapReduce.
  • Can reduce tasks start before all map tasks finish?
    The copy phase can — fetchers start once `mapreduce.job.reduce.slowstart.completedmaps` (default 0.05) of maps are done, so the network transfer overlaps with mapping. But the user's `reduce()` cannot run until every map task has completed, because an unfinished mapper might still produce values for a key the reducer would otherwise close early.
  • How many output files does a MapReduce job produce, and why?
    One per reduce task, named `part-r-00000` upward, plus a zero-byte `_SUCCESS` marker written by the OutputCommitter. With `mapreduce.job.reduces` set to 0 the job is map-only and the files are named `part-m-NNNNN` instead. The count is a property of the reduce parallelism, not of the input size or the number of HDFS blocks.

Think of a national postal sort: every local office (mapper) sorts its own sacks into per-destination bins, the trucks carry each destination's bin from every office (shuffle), and the destination office merges the arriving sacks into one ordered pile before anyone opens a letter.

saying these in an interview costs you the question

  • Says reduce() begins as soon as the first mapper output arrives
  • Claims the shuffle produces one globally sorted output across all reducers
  • Describes a MapReduce task as one partition's work inside a stage
  • Thinks intermediate map output is written to HDFS with replication
  • Says the client or ApplicationMaster merges partial results

context

open as a page

Why can adding a combiner to a Hadoop MapReduce job make an average come out wrong?

level: middleimportance: must knowfreq 48%

basics

~20 s

A combiner is a reduce-style function the framework may run zero, one or many times over a mapper's own output. An arithmetic mean is not decomposable that way, so averaging partial groups and then averaging those averages gives a wrong result.

open as a page

In Hadoop MapReduce, what is an InputSplit and what determines its size?

level: middleimportance: should knowfreq 42%

basics

~20 s

An InputSplit is the logical byte range one map task will process, plus host hints for locality — not a physical copy of data. With FileInputFormat its size is max(minSize, min(maxSize, blockSize)), so by default it equals the HDFS block size of 128 MB.

open as a page

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

level: seniorimportance: should knowfreq 45%

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.

open as a page

How do you decide which legacy Hadoop MapReduce pipelines to rewrite on Spark and which to leave alone?

level: principalimportance: should knowfreq 35%

basics

~20 s

Rank jobs by how much the MapReduce model actually costs them. Multi-job chains and iterative algorithms pay an HDFS round trip between every stage and repay a rewrite; stable single-pass map-only jobs are already IO-bound and rarely justify the risk.

open as a page

What does a Hadoop MapReduce job do differently when mapreduce.job.reduces is set to 0?

level: middleimportance: nice to knowfreq 30%

basics

~10 s

It becomes a map-only job: no partitioning, no sort, no shuffle. Each map task writes its records straight through the OutputFormat to HDFS as part-m-NNNNN files, and any configured combiner never runs.

open as a page