The map phase is where a MapReduce job touches its input. It decides how many tasks run, where they run, how each one turns raw bytes into records, and how much work each record costs. Many slow jobs are decided here: one giant gzip file that becomes a single straggler, or thousands of tiny files that become thousands of containers.

This article walks the phase in order: split computation at submission, task placement, record reading, the mapper's lifecycle, and what happens to its output. It stops at the map-side sort buffer, which the shuffle and sort article covers in depth. For the job as a whole, start with the MapReduce overview.

Advertisement

What the map phase is responsible for

A map task has one input, a split, which is a logical byte range of one or more files plus a list of hosts that hold it. It has one job: read every record in the split, call map() once per record, and hand each emitted key-value pair to the output collector. The split is logical and its boundaries need not align with records; the RecordReader turns the byte range into whole records.

The map phase: from files to key-value pairs in the sort bufferInput filesHDFS blocksInputFormat.getSplitsclient, at submitMR AppMasterone map task per splitYARN containernode-, rack- or off-localRecordReaderbyte range to recordsMapper.runsetup, map*, cleanupSort bufferpartition, sort, spillOutputFormatmap-only jobsCountersrecords, locality, GCblock locationssplits filerequest + hintskey, valuereducers > 0reducers = 0Everything to the left of the sort buffer is this article; the buffer, spills and shuffle are the next stage.
The client computes splits at submission; the ApplicationMaster requests one container per split with locality hints; each task's RecordReader feeds Mapper.run, whose output goes to the sort buffer when there are reducers, or straight to the OutputFormat when there are none.

From files to splits

Splits are computed on the client when the job is submitted, by the job's InputFormat. For every file-based format this is FileInputFormat.getSplits, and the rule is short: the split size is max(minSize, min(maxSize, blockSize)). With the defaults, a minimum of 1 byte and an unbounded maximum, the split size equals the HDFS block size, so each block becomes one map task and each task can read its block from a local disk.

// FileInputFormat, simplified: one list of splits for the whole job
long minSize = max(getFormatMinSplitSize() /* 1 */, conf "mapreduce.input.fileinputformat.split.minsize");
long maxSize = conf "mapreduce.input.fileinputformat.split.maxsize" (default Long.MAX_VALUE);

for (FileStatus file : listStatus(job)) {
    if (!isSplitable(job, file.getPath())) {                  // e.g. gzip
        splits.add(split(file, 0, file.getLen(), hostsOf(firstBlock)));
        continue;
    }
    long splitSize = Math.max(minSize, Math.min(maxSize, file.getBlockSize()));
    long remaining = file.getLen();
    while ((double) remaining / splitSize > 1.1) {           // SPLIT_SLOP: 10% slack
        long off = file.getLen() - remaining;
        splits.add(split(file, off, splitSize, hostsOf(blockAt(off))));
        remaining -= splitSize;
    }
    if (remaining != 0) {
        splits.add(split(file, file.getLen() - remaining, remaining, hostsOf(lastBlock)));
    }
}

The loop has one subtlety. It keeps cutting full-size splits only while the remaining bytes exceed 1.1 split sizes. That slop factor prevents a tiny trailing task: a file 5 percent larger than a block becomes one slightly larger split rather than a full split plus a sliver.

The two settings move the size in opposite directions. To get fewer, larger splits than blocks, raise mapreduce.input.fileinputformat.split.minsize above the block size. To get more, smaller splits, lower mapreduce.input.fileinputformat.split.maxsize below it. Setting the maximum above the block size alone does nothing, because the formula takes the smaller of maximum and block size. Splits never span files, which is why small files need a different format.

Advertisement

Worked example: counting splits

Take a 128 MiB block size and a job whose input directory holds four kinds of files:

InputRule appliedSplits
One 1 GiB text file1024 / 128 = 8 exactly8
One 135 MiB text file135 / 128 = 1.05, below 1.1, so one split1
One 300 MiB text file300/128 = 2.3, cut 128; 172/128 = 1.3, cut 128; 44 MiB left3
One 2 GiB gzip filegzip is not splittable1
10,000 files of 1 MiBone split per file10,000

The job gets 10,013 map tasks. Two lines dominate its runtime. The gzip file becomes one task that must read and decompress 2 GiB alone, most of it from remote nodes because only the first block's hosts are hinted, so it finishes long after everything else. The small files become ten thousand containers, each paying start-up cost to process 1 MiB, and the NameNode had to serve ten thousand block lookups at submission. Fixing both, by recompressing the gzip file with a splittable codec and packing the small files with CombineTextInputFormat at 256 MiB per split, brings the job to a few dozen tasks of similar size.

Record boundaries: never lose or duplicate a line

Because splits are cut at byte offsets, a line almost always straddles two splits. LineRecordReader resolves this with two symmetric rules. A reader whose split does not start at offset zero discards everything up to and including the first newline, because that partial line belongs to the previous split. And every reader keeps reading past its end offset until it finishes the line that started within its range. Each line is therefore read by exactly one task: the one whose range contains its first byte, or, for a line starting exactly at a split boundary, the task before it, which reads one extra line.

The cost is a short read into the next block. The same pattern is required of any custom record reader: define a synchronization point, such as a newline, a record marker or a sync block in a container format, skip to the first one after your start, and read until the first one after your end. Formats without one, such as a multi-line JSON document, must return false from isSplitable.

One operational limit belongs here: a corrupt file with no newlines for gigabytes makes one record enormous. Setting mapreduce.input.linerecordreader.line.maxlength makes the reader skip overlong lines instead of running out of memory.

Compression and small files

Whether a compressed file can be split depends on whether a reader can start decompressing in the middle. Gzip and plain Snappy or LZ4 streams cannot, so a file in those formats is one split regardless of size. Bzip2 can be split, at a high CPU cost. Container formats such as SequenceFile, Avro, ORC and Parquet compress blocks internally and carry sync points or footers, so they split cleanly whatever codec is inside. Prefer container formats, and keep any unsplittable file within one block.

Small files are the opposite problem: too many splits. CombineFileInputFormat and its text subclass CombineTextInputFormat pack many files, preferring files on the same node and then the same rack, into splits up to the configured maximum size. That fixes the task count but not the NameNode load of millions of files, which is a storage design problem covered in the HDFS small files problem.

Locality: where map tasks run

Each split carries the hosts holding its data, which come from HDFS block locations; see blocks and replication for how replicas are placed. The ApplicationMaster passes these hosts and their racks to YARN as preferences. The scheduler tries the node first, then the rack, then anywhere, trading a short wait for locality against leaving containers idle.

The job counters report the result: DATA_LOCAL_MAPS, RACK_LOCAL_MAPS and OTHER_LOCAL_MAPS, shown in the UI as data-local, rack-local and other local map tasks. On a healthy cluster reading HDFS, most maps should be data-local. A low ratio means a busy cluster, input concentrated on a few nodes, or input in object storage, where locality does not exist.

The Mapper lifecycle and object reuse

Every map task runs one instance of your mapper through a fixed loop:

// org.apache.hadoop.mapreduce.Mapper.run: the whole map loop
public void run(Context context) throws IOException, InterruptedException {
    setup(context);
    try {
        while (context.nextKeyValue()) {
            map(context.getCurrentKey(), context.getCurrentValue(), context);
        }
    } finally {
        cleanup(context);
    }
}

That loop explains the right place for each kind of work. Expensive initialization, such as loading a lookup table or opening a connection, belongs in setup, once per task. Per-record work belongs in map. Flushing aggregates kept in memory, as in the in-mapper combining pattern, belongs in cleanup. Overriding run changes the loop itself, as MultithreadedMapper does to run several threads per task.

public class StatusCountMapper extends Mapper<LongWritable, Text, Text, IntWritable> {
    private static final IntWritable ONE = new IntWritable(1);
    private final Text outKey = new Text();                 // reused for every record
    private Counter malformed;
    private Set<String> ignoredPaths;

    @Override
    protected void setup(Context ctx) throws IOException {
        malformed = ctx.getCounter("app", "MALFORMED_LINES");
        ignoredPaths = loadIgnoreList(ctx.getConfiguration()); // once per task, not per record
    }

    @Override
    protected void map(LongWritable offset, Text line, Context ctx)
            throws IOException, InterruptedException {
        String[] f = line.toString().split(" ");
        if (f.length < 9) { malformed.increment(1); return; } // count, do not throw
        if (ignoredPaths.contains(f[6])) return;
        outKey.set(f[8]);                                    // HTTP status
        ctx.write(outKey, ONE);                              // serialized now, safe to reuse
    }
}

Two reuse rules follow from the framework's design. The key and value objects passed to map are reused by the record reader, so storing a reference to line across calls keeps a buffer that will be overwritten; copy it if you must keep it. Conversely, output objects can be reused safely because context.write serializes them into the buffer immediately. Reusing output objects and avoiding a new String per field are the difference between a mapper that is CPU bound and one that is GC bound; check the GC_TIME_MILLIS counter.

Where map output goes

When the job has reducers, context.write applies the partitioner to each pair and places it in the in-memory collection buffer, 100 MiB by default. When the buffer passes its threshold, a background thread sorts it by partition and key, optionally runs the combiner, and spills it to local disk; at the end of the task the spills are merged into one file per task with an index. Sizing that buffer and reading spill counters are covered in the shuffle and sort article.

When the job sets mapreduce.job.reduces=0, none of that happens. Map output goes straight through the job's OutputFormat to the output file system, with no partitioning, sorting, spilling or shuffle. Map-only jobs are the right shape for filtering, format conversion, record-level enrichment and loading data, and they are often several times faster than a job with a trivial identity reducer. Each map task writes its own output file, so plan split sizes with the file count in mind.

Liveness, retries and speculation

A task that neither reads input, writes output, updates a counter nor reports progress for mapreduce.task.timeout, ten minutes by default, is killed as hung. A mapper doing long work per record, such as calling an external service, must call context.progress() to show it is alive. A failed attempt is retried on another node up to mapreduce.map.maxattempts times, four by default, after which the job fails. Because retries rerun the whole split, map functions must be deterministic and free of external side effects, or those side effects must be idempotent.

Speculative execution launches a duplicate attempt of a map that runs much slower than its peers and keeps whichever finishes first. It helps with a slow disk or a noisy neighbour, not with skew such as the gzip straggler, where the duplicate has the same work to do. Disable it per job with mapreduce.map.speculative=false in that case.

Sizing map tasks

Two numbers define a map task's memory. mapreduce.map.memory.mb is the YARN container size, which YARN enforces by killing the container if it exceeds it. mapreduce.map.java.opts sets the heap. The heap must leave room for stacks, metaspace and native codec memory; about 80 percent of the container is the usual start, and Hadoop 3 derives the heap at that ratio when -Xmx is unset.

# Container and heap for map tasks (heap roughly 80% of the container)
mapreduce.map.memory.mb=3072
mapreduce.map.java.opts=-Xmx2458m

# Fewer, larger splits for a job with a cheap map function
mapreduce.input.fileinputformat.split.minsize=268435456

# Many small files: pack them, up to 256 MiB per split
mapreduce.job.inputformat.class=org.apache.hadoop.mapreduce.lib.input.CombineTextInputFormat
mapreduce.input.fileinputformat.split.maxsize=268435456

# Liveness and retries (these are the defaults)
mapreduce.task.timeout=600000
mapreduce.map.maxattempts=4
mapreduce.map.speculative=true

Aim for map tasks that run from about a minute to a few minutes. Much shorter and container start-up dominates; much longer and a single retry or straggler delays the whole job.

Failure modes

  • Single straggler. An unsplittable or skewed input makes one task run for hours. Check the split size spread in the task list and recompress or repartition.
  • Task explosion. Tens of thousands of tiny splits overload the ApplicationMaster and the scheduler. Combine inputs or raise the minimum split size.
  • Physical memory kills. Containers killed for exceeding physical memory, with a healthy heap, mean the heap leaves too little native headroom.
  • Hung task kills. Tasks killed after the timeout while doing slow per-record work need progress reporting, not a longer timeout.
  • Held references. Collecting input Text objects in a list yields a list of identical values, all pointing at the last record.
  • Silent bad records. Throwing on a malformed line fails the task four times and then the job. Count bad records with a counter and set a threshold instead.

Trade-offs

Larger splits mean fewer containers and less scheduling overhead but less parallelism and costlier retries; smaller splits mean the reverse. Combining files trades locality for task count. Map-only jobs trade grouping for speed. Speculation trades cluster capacity for tail latency. The defaults, one block per map, suit a well-organized HDFS dataset of large splittable files; everything else is a deliberate move away from them.

What to do next

  1. For your largest job, list the input files by size and codec and predict its map count with the split formula; compare with the job's actual map count.
  2. Find any unsplittable files larger than a block and rewrite them in a splittable container format.
  3. Switch small-file inputs to CombineTextInputFormat with a maximum split size near your block size.
  4. Check the data-local, rack-local and other-local counters, and investigate if data-local maps are a minority on HDFS input.
  5. Move per-task initialization into setup, reuse output objects, and compare GC_TIME_MILLIS before and after.
  6. Convert jobs whose reducer only passes records through into map-only jobs.
  7. Set container size and heap together, and add context.progress() to any mapper with slow per-record calls.
Key takeaway: The map phase turns files into splits, splits into tasks placed near their data, and bytes into records through a RecordReader that owns every line exactly once. Split size is max(minSize, min(maxSize, blockSize)) with a 10 percent slop, compression and small files change the count, and Mapper.run fixes where setup, per-record work and cleanup belong. Size splits for tasks of a few minutes, keep inputs splittable, reuse objects, report progress, and drop the reducer when the job does not need one.