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):
- InputFormat divides input into InputSplits and a RecordReader turns bytes into input pairs.
- A map task processes its split and serializes intermediate pairs into an in-memory buffer.
- 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.
- An optional combiner performs local aggregation on map output.
- Each reducer fetches its partition from every mapper: this is the shuffle.
- The reducer merges/sorts fetched segments, groups equal keys, and calls
reduce(key, values). - 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
- More map spills mean more disk I/O; a larger sort buffer may reduce spills but leaves less memory for the mapper.
- More reducers can improve parallelism and failure isolation but create overhead and many output files.
- A reducer with a highly frequent key can be a straggler even when other reducers finish.
- A map-only job writes map results directly without shuffle/sort when reduction is unnecessary.
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
- Define combiner and why it saves network traffic.
- State why a combiner cannot be required for correctness.
- Identify InputSplit, spill, partition, shuffle, merge, and commit.
- Give an associative/commutative safe example.
- Explain the average trap.
Key takeaways
- The combiner is a local optimization, not a semantic phase guaranteed by the framework.
- MapReduce execution is dominated by serialization, local spill/merge, network shuffle, and reducer grouping.
- Correctness must survive arbitrary combiner invocation and task retry.
Sources
- Apache Hadoop MapReduce Tutorial — mapper, combiner, shuffle, reducer.
- Apache Reducer API.
- Google MapReduce paper.
- Syllabus-aligned supplement: the spill-buffer and partial-statistics examples expand the syllabus topic using the official execution documentation; no textbook page numbers are asserted.