A MapReduce job moves data in three expensive steps: map tasks write intermediate records to local disk, reducers fetch those records over the network, and reducers merge them back from disk before calling reduce(). For aggregation jobs, most of those bytes are redundant. A word count over web logs emits the pair (the, 1) millions of times from a single map task, and every copy is sorted, spilled, served, fetched and merged, only for the reducer to add them up.
The combiner is Hadoop's hook for doing part of that addition early, before the data leaves the machine. It is one of the cheapest MapReduce optimisations and one of the easiest to get subtly wrong, because Hadoop treats it as a hint: it may call your combiner zero, one or several times on any subset of a key's values. This article explains the contract a combiner must satisfy, shows where the framework invokes it, works through correct and incorrect examples, and shows how to measure whether it helped. Buffer and spill mechanics are covered in MapReduce shuffle and sort.
What a combiner is, from first principles
A combiner is a class with the same shape as a reducer: it receives a key and an iterable of values and writes key-value pairs. In the org.apache.hadoop.mapreduce API it literally is a Reducer subclass, registered with Job.setCombinerClass(). The difference is where it runs and what it is allowed to assume.
A reducer is guaranteed to see every value for its key, exactly once, after all map tasks finish. A combiner sees an arbitrary slice of those values: whatever was in one map task's sort buffer when it spilled, or whatever a reducer held in memory during a merge. It cannot know whether other values for the key exist elsewhere, or whether its output will be combined again.
So the combiner is a local, partial, optional pre-aggregation. The algebra, type rules and failure modes all follow from those three words.
The contract: why associativity and commutativity
Write the reducer's job as a function R over the full multiset of values for a key. Hadoop may split that multiset into arbitrary groups, apply the combiner C to some groups, possibly apply C again to groups of C's outputs, and then hand everything to R. The job is only correct if the answer is the same no matter how the values were grouped or ordered, and no matter how many times C ran:
- Associative: combining (a with b) and then with c gives the same result as a with (b with c). This covers arbitrary grouping and repeated application.
- Commutative: the order of values does not matter. Values arrive in an order determined by input splits, spill timing and fetch completion, none of which you control.
- Optional: the job must be correct if C never runs at all. In practice this means the reducer must accept raw map output as well as combiner output, which is why both must share a type.
Sum, count, minimum, maximum, bitwise OR and set union all satisfy these properties. So does anything that can be expressed as merging a fixed-size summary, such as a (sum, count) pair or a HyperLogLog sketch. Mean, median, mode and 'count of distinct values' computed naively do not.
Where Hadoop actually calls it
The diagram shows the three points on the shuffle path at which the framework may invoke a configured combiner. None of them is guaranteed.
- On each spill. When the map task's collection buffer reaches its spill threshold (
mapreduce.map.sort.spill.percent, 0.80 ofmapreduce.task.io.sort.mbby default), the buffer is sorted by partition and key and the combiner runs over each run of equal keys before the spill file is written. If a map task produces less output than one buffer, this is the only map-side call. - During the final map-side merge. If the task spilled several times, the spill files are merged into one output file per task. The combiner runs again during that merge only when the number of spills is at least
mapreduce.map.combine.minspills, 3 by default. Below that, re-combining two small files is judged not worth the CPU. - During the reduce-side in-memory merge. Reducers buffer fetched map outputs in memory and periodically merge them to disk. That in-memory merge runs the combiner too, which shrinks what the reducer spills to its own disk. Later merges of on-disk segments do not.
So 'the combiner runs once per map task' is wrong in both directions: a key can be combined on both sides, or never. The reduce side is covered in the MapReduce reduce phase.
The type and key rules
Because the combiner's output goes back into the same sorted, partitioned stream as the map output, three rules follow:
- Input and output types must equal the map output types. A combiner is a
Reducer<K2, V2, K2, V2>. If your reducer emits a different value type, for example aDoubleWritablemean fromIntWritableinputs, it cannot also be the combiner. - The combiner must not change keys. Its output is written into a partition segment that was already chosen and already sorted. A new key would sit in the wrong partition or out of order, and the reducer's grouping would silently break. Emit the key you received.
- No side effects. Writing to HDFS, updating an external counter service or calling an API from a combiner will happen an unpredictable number of times, including zero. Hadoop counters incremented in a combiner are similarly unreliable as business metrics.
Word count reuses its reducer as the combiner only because the sum reducer maps (Text, IntWritable) to (Text, IntWritable). Reuse is a coincidence to check, not a default.
public class WordCount {
// TokenMapper (not shown) emits (word, 1) as (Text, IntWritable).
// Safe as both combiner and reducer: input and output are (Text, IntWritable),
// and integer addition is associative and commutative.
public static class SumReducer extends Reducer<Text, IntWritable, Text, IntWritable> {
private final IntWritable total = new IntWritable();
@Override
protected void reduce(Text key, Iterable<IntWritable> vals, Context ctx)
throws IOException, InterruptedException {
int sum = 0;
for (IntWritable v : vals) sum += v.get();
total.set(sum);
ctx.write(key, total);
}
}
public static void main(String[] args) throws Exception {
Job job = Job.getInstance(new Configuration(), "wordcount");
job.setJarByClass(WordCount.class);
job.setMapperClass(TokenMapper.class);
job.setCombinerClass(SumReducer.class); // optional optimisation
job.setReducerClass(SumReducer.class); // required for correctness
job.setMapOutputKeyClass(Text.class);
job.setMapOutputValueClass(IntWritable.class);
job.setOutputKeyClass(Text.class);
job.setOutputValueClass(IntWritable.class);
FileInputFormat.addInputPath(job, new Path(args[0]));
FileOutputFormat.setOutputPath(job, new Path(args[1]));
System.exit(job.waitForCompletion(true) ? 0 : 1);
}
}
Worked example: a correct mean
Suppose map tasks emit (user, latency in milliseconds) and you want the mean latency per user. A tempting combiner computes the mean of its slice and emits it. Consider user u with values 10, 20 and 90. If one map task sees 10 and 20 and another sees 90, the combiners emit 15 and 90, and a reducer that averages those gets 52.5. The true mean is 40. The mean of means is wrong whenever the groups have different sizes, and group sizes depend on spill timing, so the answer would also change from run to run.
The fix is to carry the information needed to merge correctly. Map emits a SumCount pair of (latency, 1); the combiner adds sums and counts; only the reducer divides. With the same split, the combiners emit (30, 2) and (90, 1), the reducer gets (120, 3), and the mean is 40 however the values were grouped. SumCount is a small custom Writable holding two longs.
// Map emits (user, SumCount(latencyMs, 1)). The combiner merges partial pairs;
// only the reducer divides. Combiner output type == map output type.
public static class SumCountCombiner
extends Reducer<Text, SumCount, Text, SumCount> {
private final SumCount out = new SumCount();
@Override
protected void reduce(Text user, Iterable<SumCount> parts, Context ctx)
throws IOException, InterruptedException {
long sum = 0, count = 0;
for (SumCount p : parts) { sum += p.getSum(); count += p.getCount(); }
out.set(sum, count);
ctx.write(user, out);
}
}
public static class MeanReducer
extends Reducer<Text, SumCount, Text, DoubleWritable> {
@Override
protected void reduce(Text user, Iterable<SumCount> parts, Context ctx)
throws IOException, InterruptedException {
long sum = 0, count = 0;
for (SumCount p : parts) { sum += p.getSum(); count += p.getCount(); }
ctx.write(user, new DoubleWritable((double) sum / count));
}
}The pattern generalises: variance as (count, mean, M2) merged with the parallel variance formula, top-k as a list trimmed to k at each step, distinct counts as sketches. Median and exact percentiles have no fixed-size mergeable summary, so they need all values at the reducer, or an approximate sketch such as t-digest.
Measuring it: counters and a worked run
Every task reports Combine input records and Combine output records (the COMBINE_INPUT_RECORDS and COMBINE_OUTPUT_RECORDS task counters). Their ratio is the combiner's effectiveness. Compare them with Map output records, which counts records before combining, Spilled Records, and Reduce shuffle bytes.
Take an illustrative word-count map task over a 128 MB text split that emits 20 million (word, 1) records. With the default 100 MB buffer, suppose it spills seven times, and each spill holds about 150,000 distinct words. The first combiner pass turns 20 million records into about 1.05 million. Because seven spills is at least three, the final merge runs the combiner again over those 1.05 million and produces about 200,000, one per distinct word in the split.
| Counter | Without combiner | With combiner |
|---|---|---|
| Map output records | 20,000,000 | 20,000,000 |
| Combine input records | 0 | about 21,050,000 (both passes) |
| Combine output records | 0 | about 1,250,000 (both passes) |
| Records in file.out | 20,000,000 | about 200,000 |
| Records shuffled to reducers | 20,000,000 | about 200,000 |
Two readings matter. The output-to-input ratio is about 6 percent, so the combiner earns its CPU; a ratio near 1.0 means keys rarely repeat within a spill and the combiner is pure overhead. And the counters include both passes, so do not read combine input as map output. Reduce shuffle bytes is the number that pays for the network and for the ShuffleHandler described in the external shuffle service.
In-mapper combining: the deterministic alternative
A combiner re-serialises and re-sorts data that the mapper already had in memory. The alternative, described by Lin and Dyer in their MapReduce text-processing book, is to aggregate inside the mapper with a hash map and emit partial results only when the map grows too large or when the task ends.
// In-mapper combining: aggregate in a bounded HashMap, flush when it grows too large
// and once more in cleanup(). Runs exactly once per record, unlike a combiner.
public static class InMapperCount extends Mapper<LongWritable, Text, Text, IntWritable> {
private static final int MAX_KEYS = 100_000;
private final Map<String, Integer> partial = new HashMap<>();
@Override
protected void map(LongWritable off, Text line, Context ctx)
throws IOException, InterruptedException {
for (String t : line.toString().split("\\s+")) {
if (!t.isEmpty()) partial.merge(t.toLowerCase(), 1, Integer::sum);
}
if (partial.size() >= MAX_KEYS) flush(ctx);
}
@Override
protected void cleanup(Context ctx) throws IOException, InterruptedException {
flush(ctx);
}
private void flush(Context ctx) throws IOException, InterruptedException {
Text k = new Text(); IntWritable v = new IntWritable();
for (Map.Entry<String, Integer> e : partial.entrySet()) {
k.set(e.getKey()); v.set(e.getValue()); ctx.write(k, v);
}
partial.clear();
}
}In-mapper combining runs exactly once per record and emits far fewer records into the sort buffer, often avoiding spills entirely. It costs heap, so the map must be bounded and flushed; an unbounded map is a classic cause of map-task out-of-memory failures. It is also deterministic and testable. Many jobs use both: in-mapper combining for the hot path and a combiner as a safety net.
The same idea appears in other engines. Spark's reduceByKey and aggregateByKey combine map-side by default, while groupByKey does not; see Spark execution. Hadoop Streaming jobs can pass any executable with -combiner, provided it reads and writes the mapper's key-value format.
Failure modes
- Non-associative logic. Means, medians, 'first value seen' and ratios in the combiner give answers that change with
mapreduce.task.io.sort.mbor split size. Run once with the combiner disabled and diff. - Mismatched types. A reused reducer fails with a collector type mismatch, or a custom
Writablereads garbage. Declare combiner generics explicitly. - Key rewriting. Normalising or truncating keys breaks sort order within a spill. Normalise in the mapper.
- Filtering in the combiner. Dropping keys below a partial-count threshold loses data; only the reducer knows the full count.
- Assuming it ran. A reducer that expects one value per map task fails on the day the combiner is skipped.
Trade-offs
The combiner trades map-side CPU for less disk I/O, network and reduce-side merging. It pays when keys repeat heavily within a map task, and not for joins, raw-value reducers or nearly unique keys. Map output compression, discussed in the MapReduce overview, is complementary: the combiner removes records, compression shrinks the rest.
What to do next
- List the aggregations in your slowest MapReduce jobs and classify each as mergeable (sum, count, min, max, set union, sketch) or not.
- For every mergeable one, add a combiner with explicit
Reducer<K2, V2, K2, V2>generics, and restructure means and variances into (sum, count) style summaries. - Run the job once with and once without the combiner on the same input and diff the output to prove correctness.
- Compare Combine output records with Combine input records; remove any combiner whose ratio stays near 1.0.
- Check Reduce shuffle bytes and Spilled Records before and after, and enable map output compression alongside.
- For the hottest mappers, try bounded in-mapper combining and watch map heap usage.