MapReduce and Distributed Computation

Topics Covered

The MapReduce Model

The Map Phase

The Reduce Phase

Partitioning: Deciding Which Reducer Gets Which Key

The Shuffle: The Expensive Phase

Fault Tolerance Through Determinism

Distributed Joins

Reduce-Side Join (Sort-Merge Join)

Map-Side Join (Broadcast Join)

Partitioned Join

Join Skew: The Hidden Performance Killer

Choosing a Join Strategy

Beyond MapReduce

No Iteration Without Chaining Jobs

Excessive Disk I/O

The Materialization Problem

No Real-Time Processing

The Path to Modern Frameworks

Apache Spark Architecture

RDDs: Resilient Distributed Datasets

Lazy Evaluation and the DAG

In-Memory Processing

DataFrames and Spark SQL

Wide vs. Narrow Dependencies

Spark Fault Tolerance: Lineage

Checkpointing: Bounding Recovery Cost

Choosing a Processing Framework

Batch vs. Stream vs. Micro-Batch

Decision Framework

Resource Considerations

The Convergence of Frameworks

Common Interview Pitfalls

MapReduce is a programming model for processing large datasets across a cluster of machines. Google introduced it in a 2004 paper that described how they processed 20+ petabytes of data daily using this model. The core insight is that many data processing tasks decompose into two phases: a map phase that transforms individual records in parallel, and a reduce phase that aggregates results by key. By constraining your computation to these two phases, the framework handles distribution, fault tolerance, and scheduling automatically.

That constraint is the whole idea, and it is worth holding onto as the course moves on. The same restriction, applied to an unbounded input, produces the operators in stream processing patterns, and the argument over whether one engine should serve both is lambda and kappa architectures. The practical patterns built on top of the model are in batch processing patterns.

Two numbers set the character of every job in this model, and both come from the same place: a job finishes when its slowest task finishes, not when its average task finishes.

Tjob=max⁡itinotE[ti]T_{\text{job}} = \max_{i} t_i \quad \text{not} \quad \mathbb{E}[t_i]

Simulating 20,000 jobs whose tasks average 60 seconds, with 1 percent of tasks running 10 times slow, the mean job time is 120 s at 10 tasks, 417 s at 100, 691 s at 1,000, and 751 s at 10,000: a job of ten thousand 60-second tasks takes 12.5 times the length of its average task. Adding more machines does not help, because the problem is one machine, and this is the entire justification for speculative execution. Launching a backup copy of any task running past 1.5 times the median caps the same simulated job near 150 s.

The second number is the shuffle:

C=M⋅RC = M \cdot R

Every mapper must deliver a partition to every reducer, so 10,000 mappers and 1,000 reducers is 10 million transfers, most of them small. That product, not the volume of data, is why the shuffle is the phase that fails, and why merging map output and reducing the reducer count are the first tuning moves.

The power of MapReduce is not in the individual operations (map and reduce are basic functional programming primitives). The power is in the contract: you write two simple functions, and the framework handles distributing them across thousands of machines, moving data between phases, retrying failed tasks, and producing correct output even when individual machines fail. This separation of application logic from distributed systems concerns is what made MapReduce transformative.

A large input split into blocks, mapped in parallel, shuffled by key, and reduced, with the data moving at each stage.

The Map Phase

The map function takes one input record and produces zero or more key-value pairs. Each mapper processes its input split independently, with no communication between mappers. This is what makes the map phase embarrassingly parallel: doubling the number of mappers halves the processing time with no coordination overhead.

Consider counting word frequencies across a 1TB text corpus. Each mapper receives a 64MB block, tokenizes the text, and emits (word, 1) for every word encountered. A mapper processing the sentence "the cat sat on the mat" emits six pairs: (the, 1), (cat, 1), (sat, 1), (on, 1), (the, 1), (mat, 1). No mapper needs to know what any other mapper is doing.

 
map(document):
  for each word in document:
    emit(word, 1)

The Reduce Phase

The reduce function receives a key and all values associated with that key across every mapper's output. It aggregates those values into a final result. For word count, the reducer for key "the" receives the list [1, 1, 1, ...] from every mapper that encountered "the," and sums them.

 
reduce(word, counts):
  emit(word, sum(counts))

Reducers are also parallel, but the degree of parallelism depends on the number of distinct keys, not the input size. If your data has 10,000 unique keys and you run 100 reducers, each reducer handles roughly 100 keys.

Partitioning: Deciding Which Reducer Gets Which Key

The framework uses a partitioner to determine which reducer receives which keys. The default partitioner hashes the key and takes the modulo of the number of reducers: reducer_id = hash(key) % num_reducers. This distributes keys roughly evenly across reducers, assuming the hash function produces a uniform distribution.

Custom partitioners matter when the default hash produces skewed assignments. If your keys are dates and most of your data is from the current month, the hash partitioner sends a disproportionate volume to whichever reducer gets the current month's date. A range partitioner that assigns date ranges to reducers can balance load more effectively.

The number of reducers is a configuration choice with real consequences. Too few reducers means each one processes too much data and becomes a bottleneck. Too many reducers means excessive overhead from task scheduling and too many small output files that downstream jobs must merge. A common heuristic is to size each reducer's input to roughly 1-2GB.

The Shuffle: The Expensive Phase

Between map and reduce sits the shuffle. This is where the framework redistributes all mapper outputs so that every value for a given key arrives at the same reducer. If 1,000 mappers each emitted a pair with key "error," all 1,000 values must travel across the network to a single reducer.

Map output partitioned, sorted, spilled and fetched, so that every pair reaches the reducer that owns its key.

The shuffle is the bottleneck in most MapReduce jobs. It involves three expensive steps:

  1. Sort: Each mapper sorts its output by key and writes sorted runs to local disk.
  2. Transfer: The framework sends each key's data to the assigned reducer across the network.
  3. Merge: Each reducer merge-sorts the incoming sorted runs from all mappers.

For a job processing 1TB of input that produces 500GB of intermediate key-value pairs, the shuffle must move 500GB across the cluster network.

Optimizations like combiners (local pre-aggregation on the mapper side) reduce shuffle volume. A combiner for word count sums local counts before sending, turning 1,000 copies of (the, 1) into a single (the, 1000) per mapper.

A combiner reducing on the mapper before the shuffle, and the bytes that never cross the network because of it.
Interview Tip

When analyzing MapReduce performance, focus on the shuffle. The map and reduce phases scale linearly with more machines, but the shuffle is bounded by network bandwidth and disk I/O. Reducing shuffle volume through combiners, better key design, or filtering early in the map phase is usually the highest-leverage optimization.

Fault Tolerance Through Determinism

MapReduce achieves fault tolerance by treating tasks as deterministic functions over immutable input. If a mapper fails, the framework simply reruns it on the same input split on another machine. The output is identical because the input has not changed and the function is deterministic. Reducer failures work the same way: rerun the reducer with the same set of intermediate key-value pairs.

This design avoids the complexity of checkpointing or distributed transactions. The cost is that failed tasks must redo all their work from scratch, but for batch processing where jobs run for minutes to hours, rerunning a single task is cheap compared to the engineering cost of incremental checkpointing.

Level Expectations

Mid-level engineers should be able to explain the map-shuffle-reduce pipeline and identify the shuffle as the bottleneck. Senior engineers should reason about combiner placement, partition skew (one reducer getting disproportionate traffic), and when MapReduce is the wrong model entirely. Staff engineers evaluate whether the batch processing model fits the business latency requirements or whether a streaming approach is needed.