skip to content

What does a Hadoop MapReduce job do differently when mapreduce.job.reduces is set to 0?

level: middleimportance: nice to knowfreq 30%

answer

  1. the whole right-hand side disappears
  2. no partitioner, no sort, no network
  3. the file prefix changes letter
  4. perfect for filters and conversions
  5. the default value is one, not zero

basics

~10 s

It becomes a map-only job: no partitioning, no sort, no shuffle. Each map task writes its records straight through the OutputFormat to HDFS as part-m-NNNNN files, and any configured combiner never runs.

solid answer

~50 s

Setting `mapreduce.job.reduces=0` (equivalently `job.setNumReduceTasks(0)`) removes the entire reduce side. Map output no longer goes into the sort buffer to be partitioned and sorted on local disk; the `RecordWriter` from the job's `OutputFormat` is attached directly to each map task, which writes to HDFS as it reads, producing `part-m-00000` upward — one file per map task. Because the shuffle is where the network transfer, local disk spills and merge-sort cost live, a map-only job is dramatically cheaper and finishes as fast as the cluster can read and write the data. It suits any per-record transformation: filtering, format conversion, projection, enrichment from a side file in the distributed cache, and file copies. The counterpart gotcha is the default: `mapreduce.job.reduces` is **1**, so a job that needs reducers but never sets a count funnels the whole dataset through a single reduce task.

code

bash · 9 lines
bash
# Map-only format conversion: no shuffle, output lands as part-m-00000...
hadoop jar convert.jar com.example.ToParquet \
  -D mapreduce.job.reduces=0 \
  /raw/events /curated/events

hdfs dfs -ls /curated/events
#  _SUCCESS
#  part-m-00000
#  part-m-00001

go deeper

for a junior

Know that setting the reduce count to zero produces a map-only job with no shuffle, and that such a job is right for per-record work like filtering or converting a file format.

for a middle

Explain exactly which machinery disappears — partitioner, sort buffer and spills, combiner, copy and merge — and that map tasks then write directly through the OutputFormat as part-m files.

for a senior

Show the judgment side: recognizing when a shuffling job could be restructured as map-only, using a distributed-cache map-side join for a small dimension, and catching the default-of-one reducer that quietly serializes a large job.

for a principal

Own the design rule that shuffle is the scarce resource in a batch platform, and that pipeline design should push work into shuffle-free stages wherever the computation is genuinely per-record. The same principle governs broadcast joins in every engine you might migrate to.

## The setting and what it removes `mapreduce.job.reduces` controls how many reduce tasks a job runs. Its default is 1. Setting it to 0 — in code with `job.setNumReduceTasks(0)`, or on the command line with `-D mapreduce.job.reduces=0` for a `Tool`-based driver — turns the job into a **map-only** job, and that removes more machinery than people expect: - No `Partitioner` is consulted; there are no partitions to assign records to. - No sort: the in-memory sort buffer and its spills to local disk are bypassed entirely. - No combiner: a class registered with `setCombinerClass` is silently never invoked, because combining exists only to shrink what the shuffle carries. - No copy phase, no reduce-side merge, no `reduce()` call. Instead, the map task's `Context.write` goes straight to a `RecordWriter` created by the job's `OutputFormat`, writing to a task-attempt temporary directory in the output path. The `OutputCommitter` promotes each successful attempt to the final location, so a retried or speculative attempt cannot leave a half-written file behind. Output files are named `part-m-00000`, `part-m-00001`, … — the `m` is the tell that no reducer ran, and a useful thing to notice when reading someone else's output directory. ## Why it is so much faster Everything expensive about MapReduce lives in the shuffle: serializing and sorting every intermediate record, spilling to local disk, merging spills, transferring bytes across the network to every reducer, and merge-sorting them again. A map-only job pays none of it. Runtime becomes roughly the time to read the input and write the output, and the job scales linearly with map parallelism. Where a shuffling job's cost grows with intermediate data volume and skew, a map-only job's does not. ## What it is good for Map-only jobs are the right shape whenever the computation is per-record and needs no cross-record grouping: - filtering rows against a predicate - format conversion, for example text to Avro, ORC or Parquet - projection and column pruning - masking, tokenizing or otherwise transforming fields - lookup-style enrichment where a small dimension file is shipped via the distributed cache and held in a map in each mapper — a map-side join - bulk file copying, which is exactly how `distcp` works The distributed-cache map-side join is worth naming explicitly: if one side of a join is small enough to fit in a mapper's memory, joining it in the map phase eliminates the shuffle entirely, which is the same idea a broadcast join expresses in later engines. ## The other half of the knob Because the default is 1, a job that genuinely needs a reduce phase but never sets a count sends every key through one reduce task on one node. That job will appear to progress quickly to "map 100%, reduce 33%" and then crawl. Sizing reducers is a matter of aiming for a reasonable amount of data per reducer — enough that container startup is amortized, small enough that a reducer's merge and output stay comfortable — and staying within the parallelism the queue will actually grant you. Too many reducers is also a real cost: each one produces a file, so a large reduce count over a small dataset manufactures the small-files problem in the output directory. ## What an interviewer is checking This is a small question that separates people who have used MapReduce from people who have only read about it. Knowing that the shuffle is optional, that `part-m` versus `part-r` tells you which shape ran, and that a combiner is dead code in a map-only job, all signal real exposure.

  • How can you tell from an output directory alone whether a reduce phase ran?
    Look at the file prefix. `part-r-NNNNN` files come from reduce tasks; `part-m-NNNNN` files are written directly by map tasks in a map-only job. Both sit beside a zero-byte `_SUCCESS` marker written by the OutputCommitter once the job commits, so the marker tells you the job finished cleanly but not which shape it had.
  • What is a map-side join and why does it fit a map-only job?
    When one input is small enough to fit in a mapper's memory, you ship it through the distributed cache, load it in `setup()`, and join each record against it inside `map()`. No key needs to travel to a common reducer, so the shuffle disappears and the job becomes map-only. The constraint is the small side's memory footprint per map task.
  • What goes wrong if a job that needs reducers leaves mapreduce.job.reduces at its default?
    The default is 1, so the entire intermediate dataset is fetched, merged and reduced by a single task on a single node. The job shows all maps complete and then crawls through one long-running reducer, and it may fail on local disk or memory. The fix is to set a reduce count sized to the intermediate data volume and the queue's available parallelism.

saying these in an interview costs you the question

  • Says the combiner still runs in a map-only job
  • Thinks map output is still sorted when there are no reducers
  • Assumes the default reduce count is zero
  • Expects output files to be named part-r in a map-only job
  • Believes a map-only job cannot write to HDFS without a reducer

context