§ 2.4Module 2

Combiners and Details of MapReduce Execution

On this page

2.4 Combiners and MapReduce Execution

Recall first

A mapper emits 1,000 (word,1) pairs but only 20 distinct words. Where can those 1,000 values be reduced before crossing the network, and why can the framework skip that step?

First principles

Execution mechanics

A Hadoop MapReduce job normally follows this path (official tutorial):

  1. InputFormat divides input into InputSplits and a RecordReader turns bytes into input pairs.
  2. A map task processes its split and serializes intermediate pairs into an in-memory buffer.
  3. When thresholds are reached, the mapper spills records to local disk, sorting by key and partitioning them for reducers. Multiple spill files may be merged.
  4. An optional combiner performs local aggregation on map output.
  5. Each reducer fetches its partition from every mapper: this is the shuffle.
  6. The reducer merges/sorts fetched segments, groups equal keys, and calls reduce(key, values).
  7. Reducer output is committed to the final filesystem path.

The shuffle and sort overlap: reducers can fetch map output while maps are still finishing. Intermediate output is often compressed to reduce network and disk traffic, trading CPU for I/O savings.

Combiner

A combiner is an optional local pre-aggregation function. For word count, it can turn:

(cat,1) (cat,1) (dog,1) -> (cat,2) (dog,1)

before the shuffle. It reduces bytes transferred and values processed, but it is not a miniature reducer that the program may rely on. Hadoop may invoke it zero, one, or multiple times, and it may run on only some map outputs. Therefore the final reducer must remain correct without it.

A safe combiner generally uses an associative and commutative aggregation, such as sum, count, min, or max. Do not use a combiner for non-associative operations such as “take the first global record,” or for average unless you combine sufficient statistics (sum,count) rather than local averages.

Worked average example

Values for a key are 10, 20, and 30. Two maps see [10,20] and [30]. A wrong combiner emits local averages 15 and 30; averaging these gives 22.5, not the true 20. A correct combiner emits (sum,count) as (30,2) and (30,1); the reducer sums to (60,3) and divides to 20.

Execution trade-offs

Exercise — revealed answer

Exercise: Can a combiner be required for correctness in a word-count job?

Answer: No. Without a combiner, all (word,1) values still reach the reducer and sums remain correct. The combiner is an optimization. It is safe here because addition is associative and commutative, and the reducer accepts partial sums.

Exam lens

State three facts: optional, local, and not guaranteed to run. Then trace map buffer → spill/sort/partition → combiner → shuffle → merge/group → reduce. For average, mention (sum,count) to demonstrate understanding rather than memorization.

Rapid revision checklist

Key takeaways

  1. The combiner is a local optimization, not a semantic phase guaranteed by the framework.
  2. MapReduce execution is dominated by serialization, local spill/merge, network shuffle, and reducer grouping.
  3. Correctness must survive arbitrary combiner invocation and task retry.

Sources