Learning path

Full curriculum

Full curriculum

Arrows go from each prerequisite to the units that depend on it. Hover or focus a unit to highlight its path.

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.