A Reducer is the part of a MapReduce job you actually write after the framework has done the hard distributed work. By the time your code runs, every value for a key has been copied off the mappers, merged and sorted, and handed to you as one key and a stream of values. Getting it right is much harder: reducers that average averages, joins that hold a whole side in memory, top-N jobs that silently return the top-N per reducer, and side effects that are written twice when a task is retried are all common, and none of them fail loudly.

This article is about writing reduce functions, not about the machinery underneath. The copy, merge and commit internals are covered in MapReduce Reduce Phase, in depth. Here we start from the contract, work through the algebra that decides what is safe, build a real reducer, then cover joins, top-N, skew, side outputs and retries, and finish with a checklist you can apply to your next job.

Advertisement

The contract: what the framework promises and what it does not

In the org.apache.hadoop.mapreduce API a reducer extends Reducer<KEYIN, VALUEIN, KEYOUT, VALUEOUT>. Its run() method calls setup() once, then reduce(key, values, context) once per key group, then cleanup() once. Inside a single reduce task the key groups arrive in sort order, as defined by the sort comparator. Which keys reach which task is decided earlier by the partitioner, so the global ordering across output files is only sorted if your partitioner is order-preserving.

What the framework does not promise matters as much. Values within a group arrive in no particular order unless you engineer one with a secondary sort. The Iterable<VALUEIN> can be walked exactly once, and the framework reuses the same key and value objects as it walks, so storing a reference to a value stores the last value, not the current one. Your reduce function may be executed more than once for the same input, because task attempts fail and are retried, and because speculative execution can run two attempts at the same time. Only one attempt's file output is committed; anything else you do, such as writing to a database, happens once per attempt.

What one reduce task sees: sorted key groups, one reduce() call per groupMap output, partition 3from map 1Map output, partition 3from map 2Map output, partition 3from map NShuffle + mergesorted by keyreduce(c17, [12.0, 3.5, 40.0])reduce(c18, [7.25])reduce(c42, [1.0, 1.0, ... 9 million])setup(context)once: load config, open side outputsreduce() x key groupsstream values, emit, countcleanup(context)once: flush top-N, close outputsRecordWriter -> part-r-00003committed only if the attempt winsa hot key is one call, on one task: skew lives here
One reduce task: merged, sorted input becomes one reduce() call per key group, bracketed by setup() and cleanup(). A hot key is a single call on a single task.

Aggregation algebra: what is safe to split

A combiner is a reducer that the framework may run on map output before the shuffle, zero, one or many times. That only works if your reduction gives the same answer no matter how the values are grouped and in what order they are combined, which is to say the operation is associative and commutative, and the combiner's output type matches the map's output type. The same property decides whether you can split a hot key across reducers later, so it is worth classifying every aggregation before you write it.

AggregationSafe as combiner?How to make it splittable
sum, count, min, maxYesAlready associative and commutative
averageNo: an average of averages is wrongCarry (sum, count) pairs; divide only in the final reducer
varianceNoCarry (count, sum, sum of squares) or a mergeable Welford triple
exact distinct countPartiallyDeduplicate values per key in the combiner; the final count still needs all distinct values
median, exact percentileNoUse a mergeable sketch (t-digest, KLL) or a second sorted pass
top-NYes, if each partial keeps NEach partial keeps its local top N; the merge keeps the top N of the union

The pattern in the right-hand column is general: replace the final answer with a small, mergeable summary, merge summaries everywhere, and compute the answer only at the end. A combiner that breaks this rule does not throw an error. It produces plausible numbers that are wrong, and they are wrong differently from run to run, because whether the combiner runs depends on spill behaviour. See MapReduce Combiner, in depth for when it runs.

Advertisement

Worked example: revenue per customer plus each customer&#x27;s top three products

Input is order lines: customer_id, product_id, amount. The mapper emits (customer_id, product_id:amount). We want one output line per customer with total revenue and the three products that customer spent the most on. A first version that collects every value into a list works in a test and fails in production on the customer with nine million rows, so the reducer below keeps only bounded state per key.

public class RevenueReducer extends Reducer<Text, Text, Text, Text> {
  enum Quality { MALFORMED_VALUE, HOT_KEY }
  private static final int TOP = 3;
  private final Text out = new Text();

  @Override
  protected void reduce(Text customer, Iterable<Text> values, Context ctx)
      throws IOException, InterruptedException {
    double total = 0;
    long seen = 0;
    Map<String, Double> perProduct = new HashMap<>();   // bounded by distinct products per customer
    for (Text v : values) {                             // single pass; v is reused, so copy what you keep
      String s = v.toString();
      int sep = s.lastIndexOf(':');
      if (sep < 0) { ctx.getCounter(Quality.MALFORMED_VALUE).increment(1); continue; }
      double amt = Double.parseDouble(s.substring(sep + 1));
      total += amt;
      perProduct.merge(s.substring(0, sep), amt, Double::sum);
      if (++seen % 100_000 == 0) ctx.progress();       // keep a long group from hitting the task timeout
    }
    if (seen > 1_000_000) ctx.getCounter(Quality.HOT_KEY).increment(1);
    PriorityQueue<Map.Entry<String, Double>> heap =
        new PriorityQueue<>(Map.Entry.comparingByValue());   // min-heap of size TOP
    for (Map.Entry<String, Double> e : perProduct.entrySet()) {
      heap.offer(e);
      if (heap.size() > TOP) heap.poll();
    }
    List<String> top = new ArrayList<>();
    while (!heap.isEmpty()) top.add(0, heap.poll().getKey());
    out.set(String.format("%.2f\t%s", total, String.join(",", top)));
    ctx.write(customer, out);
  }
}

Three details carry the weight. The loop converts each Text to a fresh String before keeping anything, which sidesteps object reuse. The per-product map grows with the number of distinct products for one customer, not with the number of rows, and the heap never holds more than three entries. And the counters turn silent data problems into numbers on the job page: a job that finishes with 40,000 MALFORMED_VALUE records is not a success.

Reduce-side joins

When both inputs are large, the general join in MapReduce happens in the reducer. Each mapper tags its records with the source, for example C for customers and O for orders, and emits them under the join key. The reducer then sees every customer record and every order for that key together. The naive implementation buffers both sides in lists and emits the cross product, which is correct and runs out of memory on the first popular key.

The fix is a secondary sort so the small side arrives first. Make the map output key a composite of (join_key, tag), partition and group on join_key only, and sort so that C precedes O. The reducer reads the one customer record, holds it, and streams the orders past it without buffering them.

protected void reduce(TaggedKey key, Iterable<Text> values, Context ctx)
    throws IOException, InterruptedException {
  String customer = null;
  for (Text v : values) {
    if (key.getTag() == 'C') {             // sorted first by the composite key's comparator
      customer = v.toString();             // copy: v is reused
    } else if (customer != null) {
      ctx.write(new Text(key.getJoinKey()), new Text(customer + "\t" + v));  // inner join
    } else {
      ctx.getCounter("join", "ORPHAN_ORDER").increment(1);
    }
  }
}

Note that key.getTag() changes as you iterate: with grouping on the join key only, the framework updates the key object to match the current value. That is exactly the behaviour secondary sort relies on. If one side fits comfortably in memory, skip all of this and do a map-side join from the distributed cache; it avoids the shuffle entirely.

Global top-N and distinct counts

A reducer only sees its own partition, so a top-N reducer with twenty reduce tasks returns twenty top-N lists. There are two honest ways to get a global answer. The first is to keep a bounded heap in each reducer across all its keys, emit it in cleanup(), and run a second, tiny job with one reducer that merges the twenty partial lists. The second, for small inputs, is to set the reducer count to one, which is simple and becomes a bottleneck as data grows.

Exact distinct counting has the same shape. Counting distinct users per country in one reducer means holding every user ID for the largest country in a set. A cleaner pattern is two jobs: the first uses (country, user) as the key, so the shuffle deduplicates for you and each reducer emits (country, 1) once per group; the second sums.

Choosing the number of reducers

The reducer count is set with job.setNumReduceTasks(n) or mapreduce.job.reduces, which defaults to 1. It determines the number of output files and the parallelism of the reduce stage. Too few reducers and each one processes a huge partition slowly; too many and you pay task start-up costs and create thousands of small files that hurt every downstream reader, as described in the HDFS small files problem.

A practical starting point is to size by bytes, not by cluster slots: estimate the shuffled map output and aim for roughly one to a few gigabytes per reducer, then adjust after looking at the spread of reduce task durations. Setting the count to zero makes the job map-only, which skips the shuffle and sort completely; use it whenever there is no grouping to do.

Skew: handling a hot key in code

The partitioner sends every record for a key to one reducer, and the framework calls reduce() once for that key. No setting splits a single key across tasks. When one key holds a large share of the data, one reduce task runs for hours while the rest finish in minutes.

If the aggregation is splittable by the table above, salt the key. The mapper appends a small random or hash-derived suffix, key#0 to key#15, for keys known or detected to be hot, the first job reduces the salted keys into partial summaries, and a second job strips the suffix and merges the partials. For joins, salt the large side and replicate the matching small-side record to every suffix. If the aggregation is not splittable, such as an exact median, the only real fixes are a sketch or a different algorithm; adding reducers does nothing.

Side outputs, counters and determinism under retries

Write all reducer output through the framework. MultipleOutputs lets one reducer write several named outputs, for example good records and rejected records, and those files go through the same task-attempt directory and commit protocol as the main output, so a failed or losing attempt leaves nothing behind. Create it in setup(), call mos.write(name, key, value) in reduce(), and close it in cleanup(); forgetting the close produces truncated or empty files. Pair it with LazyOutputFormat if you do not want empty default part files.

Anything written outside the framework is not protected. A reducer that inserts rows into a database or calls an external API will repeat those writes when an attempt is retried, and will do them twice in parallel when a speculative attempt runs, as explained in speculative execution in Hadoop. Either make such writes idempotent, keyed by a deterministic ID such as the record key, or disable reduce speculation with mapreduce.reduce.speculative=false for that job and still design for retries. Determinism also means the reducer must not depend on value order, wall-clock time or random numbers without a fixed seed; otherwise a retried attempt produces different output from the one that failed.

Counters are the reducer's built-in metrics. Use them for malformed inputs, orphan join records and hot keys, and fail the job in the driver when a counter crosses a threshold.

Testing reducers

MRUnit, once the standard tool, was retired to the Apache Attic, so test reducers as ordinary Java. Keep the logic in a pure function that takes a key and an iterator of plain values and returns results, unit test that function with edge cases (empty group, one value, a malformed value, a single huge group), and keep the Reducer subclass a thin adapter. For the adapter, a mocked Context that records writes is enough. Then run the whole job in local mode on a small fixture, with the combiner both on and off, and assert the outputs match: if they differ, your aggregation is not as associative as you thought.

Failure modes

SymptomLikely causeFix
All stored values equal the last oneObject reuse in the values iteratorCopy values before keeping them
Numbers change with cluster loadNon-associative combiner, such as averagingCarry mergeable summaries
One reducer runs for hoursHot keySalt and two-stage aggregate; or a sketch
Reducer OOM on one keyBuffering a whole groupBounded state, secondary sort, pre-aggregation
Task killed after 600 s with no outputLong group, no progress reportedCall context.progress() in the loop
Duplicate rows in an external storeRetries and speculative attemptsIdempotent writes or framework outputs
Top-N returns too many rowsPer-reducer top-NMerge partial lists in a second job

What to do next

  1. Classify every aggregation in your job as splittable or not, and change averages and variances to carry mergeable summaries.
  2. Audit each reducer for kept references to values; copy before storing.
  3. Replace any list that buffers a whole group with bounded state or a secondary sort.
  4. Add counters for malformed records, orphans and hot keys, and fail the driver on thresholds.
  5. Size the reducer count from shuffled bytes and check the task-duration spread after the next run.
  6. Move external writes into MultipleOutputs, or make them idempotent and disable reduce speculation.
  7. Run the job locally with the combiner on and off and diff the outputs.
Key takeaway: A reducer is called once per key group on sorted, merged input, may run more than once, and sees only its own partition. Write it as a single streaming pass with bounded state, copy values you keep, carry mergeable summaries so combiners and salting stay correct, use secondary sort for joins, merge partial top-N lists in a second stage, size reducer counts by bytes, and route every output through the framework so retries and speculative attempts cannot duplicate it.