Unit content
MapReduce and shuffle-based data parallelism
MapReduce structures large data processing as independent mapping followed by grouping and reduction.
A map function processes each input record and emits zero or more key-value pairs. For word counting, a document containing red blue red can emit
(red, 1)
(blue, 1)
(red, 1)
The shuffle phase routes values so that all values for the same key are grouped together:
red -> [1, 1]
blue -> [1]
A reduce function then combines each group, producing
(red, 2)
(blue, 1)
Mapping is naturally data-parallel because input partitions can be processed independently. Reduction is parallel across different keys and can often use local partial reductions before the shuffle to reduce communication.
The expensive step is frequently data movement rather than the map function itself. Key skew can also overload one reducer even when input records were evenly partitioned.
MapReduce is therefore both a programming model and a decomposition strategy: expose independent record processing, make the required regrouping explicit, then aggregate grouped data.