A MapReduce job hands you one hook into the read side of the computation: the mapper. The framework reads each input split, turns it into records, and calls your map() once per record. Whatever you pass to context.write() becomes the job's intermediate data, and every byte of it is partitioned, sorted, spilled, copied and merged before a reducer sees it. So the mapper decides most of what a job costs.
This page is about writing mappers: the contract of the Mapper class, the Context object, overriding run(), map-only jobs, map-side joins, Streaming and testing. How the framework splits files, schedules map tasks and sizes them is covered in the map phase in depth; read that for splits, locality and the lifecycle of a map task.
Where the Mapper sits
For one map task the framework builds an InputSplit (usually a byte range of one HDFS file), asks the job's InputFormat for a RecordReader over it, creates one instance of your mapper class by reflection, and calls its run() method with a Context that wraps the reader and the output collector. With TextInputFormat a record is a byte offset and a line.
What happens to your output depends on the reducer count. With one or more reducers, each pair goes to the partitioner and a sort buffer; the shuffle and sort page follows those bytes. With zero reducers the output format writes pairs straight to part-m-NNNNN files.
The contract: four type parameters and four methods
The class is Mapper<KEYIN, VALUEIN, KEYOUT, VALUEOUT>. The input types are fixed by the input format, and the output types must match what you declare with setMapOutputKeyClass and setMapOutputValueClass (or the job's output classes, if you leave those unset). A mismatch is not caught at compile time; the task fails at the first write with a type mismatch error. Output keys must be WritableComparable when there are reducers, because they are sorted.
This is the whole base class, abridged from the current Hadoop source:
// org.apache.hadoop.mapreduce.Mapper, current trunk (abridged)
public class Mapper<KEYIN, VALUEIN, KEYOUT, VALUEOUT> {
protected void setup(Context context) throws IOException, InterruptedException {
// NOTHING
}
@SuppressWarnings("unchecked")
protected void map(KEYIN key, VALUEIN value, Context context)
throws IOException, InterruptedException {
context.write((KEYOUT) key, (VALUEOUT) value); // identity by default
}
protected void cleanup(Context context) throws IOException, InterruptedException {
// NOTHING
}
public void run(Context context) throws IOException, InterruptedException {
setup(context);
try {
while (context.nextKeyValue()) {
map(context.getCurrentKey(), context.getCurrentValue(), context);
}
} finally {
cleanup(context);
}
}
}Three consequences follow. First, cleanup() runs in a finally block, so it runs even when map() throws; code there must cope with half-initialized state. Second, the default map() is the identity, which is why a job with no mapper class configured passes its input through unchanged. Third, one instance handles every record in the split on one thread, so instance fields survive across calls. That enables per-task state, and it is why the framework may reuse key and value objects between calls: copy anything you keep.
Worked example: parsing access logs
The job: for every HTTP status code, the total bytes served. Input is raw web server logs, some lines truncated or corrupt. The mapper parses each line, skips bad ones visibly, drops paths an operator wants ignored, and emits status and byte count.
public class StatusBytesMapper
extends Mapper<LongWritable, Text, Text, LongWritable> {
enum Logs { MALFORMED, PARSED, FILTERED }
private static final Pattern LINE = Pattern.compile(
"^(\\S+) \\S+ \\S+ \\[[^\\]]+\\] \"(\\S+) (\\S+) [^\"]*\" (\\d{3}) (\\d+|-)$");
private final Text outKey = new Text(); // reused: the framework
private final LongWritable outVal = new LongWritable(); // serializes on write()
private Set<String> ignoredPaths;
@Override
protected void setup(Context ctx) {
Configuration conf = ctx.getConfiguration();
ignoredPaths = new HashSet<>(conf.getTrimmedStringCollection("logs.ignore.paths"));
}
@Override
protected void map(LongWritable offset, Text line, Context ctx)
throws IOException, InterruptedException {
Matcher m = LINE.matcher(line.toString());
if (!m.matches()) {
ctx.getCounter(Logs.MALFORMED).increment(1);
return; // skip, never throw
}
if (ignoredPaths.contains(m.group(3))) {
ctx.getCounter(Logs.FILTERED).increment(1);
return;
}
outKey.set(m.group(4)); // HTTP status
outVal.set("-".equals(m.group(5)) ? 0L : Long.parseLong(m.group(5)));
ctx.write(outKey, outVal);
ctx.getCounter(Logs.PARSED).increment(1);
}
}The pieces worth copying are the habits, not the regex. Configuration is read once in setup(), not per record. Bad input increments a counter and returns, because a thrown exception fails the task attempt, the attempt is retried on the same split, it fails the same way, and after the configured number of attempts the whole job fails on one bad line. Output objects are allocated once and reused, which is safe because write() serializes the pair into the buffer immediately.
Add a summing reducer, and the same class as combiner, so each task ships one pair per status code.
Context: counters, configuration and progress
The Context passed to every method is your connection to the framework. The parts you will use:
write(key, value)emits a pair.getCounter(enum)orgetCounter(group, name)returns a counter that is aggregated across all tasks and shown in the job history. Use them to report data quality, with a bounded set of names; never create a counter per key.getConfiguration()reads job parameters set by the driver withconf.set(...)or-Don the command line. This is how you pass small settings into tasks.progress()andsetStatus(String)tell the framework the task is alive. A task that neither reads input, writes output nor reports progress formapreduce.task.timeout(600,000 ms by default) is killed. A mapper that calls a slow external service or builds a large structure insetup()should callprogress()as it goes.
Overriding run()
Most mappers override map() and leave run() alone. Override run() when the per-record call is the wrong unit of work. The common cases are batching (collect 500 records, send one request to a service, emit the results), stopping early (a sampling job that needs the first N matching records of each split and should not read the rest), and wrapping the whole loop in a resource such as a connection that setup() and cleanup() would otherwise have to share through fields.
If you override it, keep the try/finally shape, keep calling context.nextKeyValue() until you are done, and flush any partial batch in the finally path. Stopping early is safe.
Map-only jobs and multiple outputs
Many jobs do not need a reduce at all: format conversion, filtering, parsing, enrichment against a small table, or any per-record transformation. Setting the reducer count to zero removes the partition, sort, spill, shuffle and merge entirely, which is usually the largest cost in a MapReduce job.
Job job = Job.getInstance(conf, "status-bytes");
job.setJarByClass(StatusBytesMapper.class);
job.setMapperClass(StatusBytesMapper.class);
job.setMapOutputKeyClass(Text.class); // must match KEYOUT
job.setMapOutputValueClass(LongWritable.class);
// Map-only variant: no partition, no sort, no shuffle.
job.setNumReduceTasks(0);
job.setOutputKeyClass(Text.class);
job.setOutputValueClass(LongWritable.class);
MultipleOutputs.addNamedOutput(job, "errors", TextOutputFormat.class,
Text.class, LongWritable.class);
LazyOutputFormat.setOutputFormatClass(job, TextOutputFormat.class); // no empty part filesTwo things change in a map-only job. Output is not sorted, and there is one output file per map task, so a job over 4,000 splits writes 4,000 files; if they are small, follow it with a compaction step or choose a larger split size. And the output types are now the job's output types, because the mapper writes to the output format directly.
MultipleOutputs lets one mapper write to several named outputs, such as good records to the main output and 5xx lines to an errors directory: create it in setup(), call mos.write("errors", key, value, "errors/part") in map(), and close it in cleanup(), or the files are truncated. Its files go through the output committer like the main output.
Map-side joins
When one side of a join is small enough to hold in every task's memory, do the join in the mapper and skip the shuffle. The driver ships the small file with the job through the distributed cache; each task loads it in setup() and probes it per record:
// Driver: ship the small side with every task.
job.addCacheFile(new URI("/ref/countries.tsv#countries")); // symlinked as ./countries
// Mapper: load it once per task, join every record against it.
private final Map<String, String> countries = new HashMap<>();
@Override
protected void setup(Context ctx) throws IOException {
try (BufferedReader r = Files.newBufferedReader(Paths.get("countries"))) {
String row;
while ((row = r.readLine()) != null) {
String[] f = row.split("\t", 2);
countries.put(f[0], f[1]);
}
}
}
@Override
protected void map(LongWritable k, Text v, Context ctx)
throws IOException, InterruptedException {
String[] f = v.toString().split("\t");
String name = countries.get(f[2]);
if (name == null) { ctx.getCounter("join", "unmatched").increment(1); return; }
ctx.write(new Text(f[0]), new Text(name));
}Size the table honestly. A file that is 200 MB as tab-separated text can easily take several times that as Java strings in a HashMap, and it must fit inside the task heap set by mapreduce.map.java.opts along with the sort buffer, which itself sits inside the container size mapreduce.map.memory.mb. Count unmatched keys, so a stale reference file shows up as a number.
Several mappers in one task: MultithreadedMapper and ChainMapper
MultithreadedMapper runs your mapper class on a thread pool inside one task: you set it as the job's mapper and name the real class with MultithreadedMapper.setMapperClass and the pool with setNumberOfThreads, which defaults to 10. It helps only when each record waits on something slow, such as a network call. Reading input and writing output are serialized, so it does not help CPU-bound work, and any state shared between threads must be thread-safe.
ChainMapper runs several mapper classes in sequence within one task, passing each one's output to the next in memory, so a parse step, a filter step and an enrichment step can stay as separate tested classes without writing intermediate data to HDFS.
Streaming: a mapper in any language
Hadoop Streaming runs any executable as the mapper. The framework writes each record to the program's standard input as a line, reads its standard output, and splits each output line at the first tab into key and value. A line on standard error of the form reporter:counter:group,name,amount increments a counter, and reporter:status:message sets the task status. Here is the access-log mapper in Python:
#!/usr/bin/env python3
"""Hadoop Streaming mapper: emit (status, bytes) per access-log line."""
import re
import sys
LINE = re.compile(r'^(\S+) \S+ \S+ \[[^\]]+\] "(\S+) (\S+) [^"]*" (\d{3}) (\d+|-)$')
def main(stdin=sys.stdin, stdout=sys.stdout, stderr=sys.stderr):
for raw in stdin:
m = LINE.match(raw.rstrip("\n"))
if not m:
# Streaming turns this stderr line into a job counter.
stderr.write("reporter:counter:logs,malformed,1\n")
continue
status, size = m.group(4), m.group(5)
stdout.write(f"{status}\t{0 if size == '-' else int(size)}\n")
if __name__ == "__main__":
main()Because the contract is just lines in and lines out, you can test it with a pipe before touching a cluster. Running it on five sample lines prints:
$ python3 mapper.py < sample.log 2> err.log # 5 lines; the 3rd is truncated garbage
200 5120
404 312
500 0
200 20480
$ cat err.log
reporter:counter:logs,malformed,1Submit it with mapred streaming -files mapper.py -mapper "python3 mapper.py" -numReduceTasks 0 plus input and output paths. Streaming costs a process per task and a text round trip per record, so it suits logic that already exists in another language more than hot paths. A script that exits non-zero fails the attempt.
Testing a mapper
MRUnit was retired to the Apache Attic in 2016, so do not build new tests on it. Test the parsing and business logic as plain functions with no Hadoop types at all, which covers most of the risk. Then test the wiring with a mocked Context:
// Lives in the same package as StatusBytesMapper (setup/map are protected).
@Test
void malformedLinesAreCountedNotEmitted() throws Exception {
Mapper<LongWritable, Text, Text, LongWritable>.Context ctx = mock(Mapper.Context.class);
Counter malformed = mock(Counter.class);
when(ctx.getCounter(StatusBytesMapper.Logs.MALFORMED)).thenReturn(malformed);
when(ctx.getConfiguration()).thenReturn(new Configuration());
StatusBytesMapper m = new StatusBytesMapper();
m.setup(ctx);
m.map(new LongWritable(0), new Text("garbage line"), ctx);
verify(malformed).increment(1);
verify(ctx, never()).write(any(), any());
}Add one end-to-end run in local mode (mapreduce.framework.name=local) over a small file checked into the test resources. It catches mapper-driver type mismatches that unit tests miss.
Failure modes
- Exception on bad input. One malformed record fails the attempt every time, so the job dies after the retry limit. Catch, count and skip.
- Holding references to reused objects. Storing the incoming
Textin a list keeps a pointer to one object that is overwritten each call, so the list ends up full of the last value. Copy withnew Text(value)ortoString(). - Unbounded state in fields. An in-memory map keyed by user ID grows with the split and ends in an out-of-memory error or a container kill for exceeding
mapreduce.map.memory.mb. Bound it and flush when full. - Side effects outside the committer. Writing to a database or a fixed HDFS path from the mapper is repeated by retries and by speculative attempts. Make such writes idempotent, or write through the output format.
Trade-offs
| Choice | Gain | Cost |
|---|---|---|
| Map-only job | No sort or shuffle at all | Unsorted output, one file per task |
| Map-side join | Join without moving the big side | Small side must fit in every task's heap |
| Overriding run() | Batching, early exit, scoped resources | You own the loop and the cleanup path |
| MultithreadedMapper | Overlaps slow per-record I/O | Thread-safety burden; no help for CPU-bound work |
| Streaming | Any language, simple to test | Process and text-serialization overhead |
What to do next
- Open one production mapper and check that malformed input is counted and skipped rather than thrown.
- Move any per-record configuration lookup or object allocation into
setup()or a reused field. - Compare the map output bytes and map input bytes counters for your heaviest job; if output is larger than input, look for a combiner or a smaller output record.
- List jobs with a reducer that only concatenates or passes records through, and make them map-only.
- For any join against a table under a few hundred megabytes, try a map-side join and measure task heap use.
- Pipe sample input through each Streaming mapper locally before every deploy.
- Add a mocked-Context unit test and one local-mode run for each mapper you change.