Jeffrey Dean@Google and Sanjay Ghemawat@Google
The need
Google's workload means the data to compute is usually huge. You have to spread it across hundreds or thousands of machines to finish in a reasonable time. How to compute in parallel, distribute data, and handle faults becomes the problem.
To solve these, the authors abstract out a simple computation, and hide the details of parallelism, fault tolerance, data distribution, and load balancing.
The programming model they designed is inspired by the map and reduce primitives in Lisp.
Programming model
A set of key/value pairs as input to the computation, producing a set of key/value pairs as output.
The MapReduce library expresses this through a Map function and a Reduce function.
Map function
The Map function, written by the user, takes a key/value pair as input and produces a set of intermediate key/value pairs. The MapReduce library groups all intermediate values with the same key and passes them to the Reduce function.
Reduce function
The Reduce function, also written by the user, takes an intermediate key and a set of values for that key. It merges those values into a possibly smaller set. Typically each Reduce call writes 0 or 1 output value. Intermediate values are handed to reduce via an iterator. That lets us handle a list of values too large to fit in memory.

Implementation
Map calls are spread across many machines. Input is automatically split into M pieces, which different machines can process in parallel. Intermediate keys are partitioned into R pieces, then Reduce is called. R (the number of partitions) is chosen by the user.
- The MapReduce library first splits the input files into M pieces, typically 16MB or 64MB each. Then it starts many copies of the user program on a set of machines.
- One copy is special — the master. The rest are workers; the master assigns them tasks. There are M map tasks and R reduce tasks to assign. The master picks idle workers and gives them a map task or a reduce task.
- A worker given a map task reads the corresponding input split, parses key/value pairs out of the data, and passes them to the user-defined map function. Intermediate key/value pairs from map are buffered in memory.
- Periodically, buffered pairs are written to local disk, partitioned into R regions by the partitioning function. The locations of these buffered pairs on local disk are sent back to the master, which forwards them to workers running reduce.
- When a reduce worker is told these locations by the master, it uses RPC to read the buffered data from the map workers' local disks. After reading all intermediate data, it sorts by intermediate key so that all occurrences of the same key are grouped. Sorting is needed because many different keys typically map to the same reduce task. If the intermediate data is too large for memory, an external sort is used.
- The reduce worker walks the sorted intermediate data. For each unique intermediate key, it passes the key and the corresponding set of values to the user's Reduce function. Reduce output is appended to the final output file for this reduce partition.
- When all map and reduce tasks are done, the master wakes the user program. At that point the MapReduce call in the user program returns to user code.
After a successful run, MapReduce output is available in R output files (one per reduce task; names specified by the user). Users usually don't merge these R files into one — they typically pass them as input to another MapReduce call, or use them from another distributed application that can handle input split across files.
Master's data structures
The master keeps several data structures. For each map and reduce task, it stores the state (idle, in-progress, or completed) and the identity of the worker machine (for non-idle tasks).
The master is the conduit that propagates the locations of intermediate file regions from map tasks to reduce tasks. So for each completed map task, the master stores the locations and sizes of the R intermediate file regions that map produced. After a map task finishes, it gets updates of this location and size info. The info is pushed incrementally to workers that are in the middle of reduce tasks.
Fault tolerance
Because the MapReduce library is meant to help process huge data on hundreds or thousands of machines, it has to tolerate machine failures gracefully.
Worker failure
The master periodically pings each worker. If it gets no response within a certain time, it marks that worker as failed. Any map tasks completed by that worker are reset to idle, so they can be scheduled on other workers. Likewise, any map or reduce tasks in progress on the failed worker are reset to idle and eligible to be rescheduled.
Completed map tasks are re-executed on failure because their output lives on the failed machine's local disk and is no longer reachable. Completed reduce tasks don't need re-execution; their output is in the global file system.
When a map task is first run by worker A and then by worker B (because A failed), all reduce workers are notified of the re-execution. Any reduce task that hasn't yet read from A will read from B.
MapReduce is resilient to large-scale worker failure. For example, during one MapReduce operation, network maintenance on the running cluster made a group of 80 machines unreachable for a few minutes. The MapReduce master simply re-executed the work those unreachable workers had done, kept going, and finished the operation.
Master failure
It's easy to have the master write periodic checkpoints of the data structures above. If the master task dies, a new copy can start from the last checkpoint. Given there's only one master, failure is unlikely; so if the master fails, the current implementation aborts the MapReduce computation. The client can detect this and retry as needed.
Semantics in the presence of failures
When the user-supplied map and reduce operators are deterministic functions of their inputs, the distributed implementation produces the same output as a sequential, fault-free execution of the whole program.
We rely on atomic commit of map and reduce task output for this. Each in-progress task writes its output to private temporary files. A reduce task produces one such file; a map task produces R (one per reduce task). When a map task finishes, the worker sends the master a message with the names of the R temp files. If the master gets a completion message for a map task already done, it ignores it. Otherwise it records the R file names in the master data structures.
When a reduce task finishes, the reduce worker atomically renames its temp output file to the final output file. If the same reduce task runs on multiple machines, there are multiple rename calls on the same final output file. We rely on atomic rename from the underlying file system so the final FS state contains the data from only one execution of that reduce task.
The vast majority of our map and reduce operators are deterministic.
In this case, our semantics are equivalent to sequential execution, which makes it easy for programmers to reason about program behavior. When map and/or reduce operators are non-deterministic, we provide weaker but still reasonable semantics. With non-deterministic operators, the output R1 of a particular reduce task is equivalent to the output of R produced by some sequential execution of the non-deterministic program. Outputs of different reduce tasks R, though, may correspond to different sequential executions of that non-deterministic program.
Locality
Network bandwidth is a relatively scarce resource in our compute environment. We save bandwidth by using the fact that input data (managed by GFS [8]) is stored on local disks of the machines in the cluster. GFS divides each file into 64MB chunks and stores several replicas of each chunk (typically 3) on different machines. The MapReduce master considers input-file location info and tries to schedule a map task on a machine that already holds a replica of the corresponding input. If that fails, it tries to schedule the map task near a replica (e.g. a worker on the same network switch as the machine that has the data). When a large MapReduce runs on a large fraction of the cluster's workers, most input is read locally and doesn't consume network bandwidth.
Task granularity
There are practical limits on M and R. We often use 2,000 worker machines to run MapReduce computations with M = 200,000 and R = 5,000.
Backup tasks
One common reason a MapReduce operation takes a long time overall is "stragglers": a machine taking an unusually long time to finish one of the last few map or reduce tasks. Stragglers can appear for many reasons. For example, a machine with a bad disk may hit correctable errors often, dropping read performance from 30 MB/s to 1 MB/s. The cluster scheduler may have put other tasks on the machine, so CPU, memory, local disk, or network contention slows the MapReduce code. One recent bug in machine-init code disabled the processor cache: computation on affected machines slowed by more than a hundred times.
We have a general mechanism to mitigate stragglers. When a MapReduce operation is close to done, the master schedules backup executions of the remaining in-progress tasks. As soon as the primary or the backup finishes, the task is marked complete. We've tuned this so it typically adds only a few percent of extra compute resources. We found it significantly reduces the time to finish large MapReduce operations.
Refinements
Partitioning function
MapReduce users specify how many reduce tasks / output files they want (R). Intermediate data is partitioned across those tasks by a partitioning function on the intermediate key. A default hash partition is provided (e.g. hash(key) mod R). That tends to give fairly balanced partitions. Sometimes it's useful to partition by some other function of the key. For example, sometimes the output key is a URL, and we want all entries for a single host in the same output file. To support that, users can supply a special partitioning function. Using hash(Hostname(urlkey)) mod R puts all URLs from the same host in the same output file.
Ordering guarantees
We guarantee that within a given partition, intermediate key/value pairs are processed in increasing key order. That makes it easy to produce a sorted output file per partition, which is useful when the output format needs efficient random access by key, or when users of the output find sorting convenient.
Combiner function
Sometimes there is significant repetition of intermediate keys produced by each map task, and the user-specified Reduce is commutative and associative. Word counts often follow a Zipf distribution, so each map task produces hundreds or thousands of records of that form. All those counts would be sent over the network to a single reduce task, then added by Reduce into one number. We let the user specify an optional combiner that partially merges this data before it goes over the network.
The combiner runs on every machine that executes a map task. Usually the same code implements combiner and reduce. The only difference is how the MapReduce library handles the function's output. Reduce output is written to the final output file. Combiner output is written to an intermediate file that will be sent to a reduce task.
Partial combining significantly speeds up some classes of MapReduce operations.
Input and output types
The MapReduce library supports reading input data in many different formats.
Side effects
Sometimes MapReduce users find it convenient to produce auxiliary files as extra output from their map and/or reduce operators. We rely on the application writer to make such side effects atomic and idempotent. Typically the application writes a temp file and atomically renames it once the file is fully generated.
We don't support atomic two-phase commit of multiple output files from a single task. So tasks that produce multiple output files with cross-file consistency requirements should be deterministic. This restriction has never been a problem in practice.
Skipping bad records
Sometimes there's a bug in user code that makes Map or Reduce deterministically crash on certain records. That can keep a MapReduce operation from finishing. The usual fix is to fix the bug, but sometimes that's not feasible — maybe it's in a third-party library without source. And sometimes ignoring a few records is acceptable, e.g. statistical analysis of a large dataset. We provide an optional execution mode where the MapReduce library detects which records cause deterministic crashes and skips them so progress can continue.
Each worker process installs a signal handler for segmentation violations and bus errors. Before invoking user Map or Reduce, the library stores the argument sequence number in a global variable. If user code raises a signal, the handler sends a "last gasp" UDP packet with the sequence number to the MapReduce master. When the master sees multiple failures on a particular record, it says that record should be skipped on the next re-execution of that Map or Reduce task.
Local execution
Debugging problems in Map or Reduce can be tricky, because the real computation happens in a distributed system, often on thousands of machines, with work assignment decided dynamically by the master. To help debugging, profiling, and small-scale testing, we built an alternative MapReduce library implementation that runs all the work of a MapReduce operation sequentially on the local machine. Users get controls to restrict the computation to particular map tasks. They invoke their program with a special flag, then can easily use whatever debugging or testing tools they find useful (e.g. gdb).
Status information
The master runs an internal HTTP server and exports a set of status pages for humans. The pages show progress: how many tasks done, how many in progress, input bytes, intermediate bytes, output bytes, processing rate, etc. They also have links to each task's stderr and stdout files. Users can use this to predict how long the computation will take, and whether more resources should be added. The pages also help tell when a computation is much slower than expected.
The top-level status page also shows which workers failed, and which map and reduce tasks they were working on at the time. That's useful when diagnosing bugs in user code.