A MapReduce job does not run each task once. It runs task attempts: a reduce that fails is retried, a slow one may get a speculative twin, and an ApplicationMaster that crashes is restarted and has to pick up whatever the first one finished. Every one of those attempts writes files. Something has to make sure that, when the job reports success, the output directory contains the output of exactly one successful attempt per task, and that nothing half-written is visible if the job fails. That something is the output committer.
This page explains the commit protocol from first principles: what the OutputCommitter contract promises, who calls each method and when, how the AppMaster arbitrates between competing attempts, what FileOutputCommitter versions 1 and 2 actually do to the directory tree, and why object stores need different committers. The shuffle and reduce side that produces the files are covered in the MapReduce reduce phase; this page starts at the moment a task has bytes to publish.
The problem a committer solves
Take a job with 400 reduce tasks writing to /data/out. Reducer 3 is slow, so the AppMaster launches a speculative attempt on another node, as described in speculative execution. Both attempts write part-r-00003. Reducer 117 fails twice on a bad disk and succeeds on its third attempt. If tasks wrote straight into /data/out, it would hold complete files, truncated files from dead attempts, and two writers racing on one name.
The committer gives the job three properties. Isolation: every attempt writes into its own private directory, so attempts cannot see or damage each other. Single winner: exactly one attempt per task has its output promoted. Job atomicity, ideally: downstream readers see either nothing or the complete output, signalled by a _SUCCESS marker. How well each property holds depends on the algorithm and on what the filesystem's rename can promise, which is the theme of the rest of this page.
The OutputCommitter contract
The contract is the abstract class org.apache.hadoop.mapreduce.OutputCommitter, returned by the job's OutputFormat. Its methods split into job-level calls made once by the AppMaster and task-level calls made per attempt.
| Method | Called by | Purpose |
|---|---|---|
setupJob | AM, once | Prepare job-level state, for example create the job temporary directory. |
setupTask | Task attempt | Prepare attempt-level state; FileOutputCommitter creates the attempt directory lazily. |
needsTaskCommit | Task attempt | Return false when the attempt wrote nothing, so commit is skipped. |
commitTask | Task attempt, once permitted | Promote this attempt's output to committed task output. |
abortTask | Task or AM cleanup | Discard an attempt's output. |
commitJob | AM, after all tasks | Make the job's output final and visible. |
abortJob | AM, on failure or kill | Discard the job's output and temporary state. |
isRecoverySupported(JobContext) | Restarted AM | Whether committed task output survives an AM restart. |
recoverTask | Restarted AM | Carry a previously committed task forward instead of rerunning it. |
isCommitJobRepeatable | Restarted AM | Whether an interrupted job commit can be retried. |
Two rules make the contract work. commitTask may only run after the AppMaster has granted permission, and commitJob runs only after every task has a committed attempt. Everything else, including which directories exist and what a rename means, is an implementation detail of the particular committer.
Arbitration: only one attempt may commit
The arbitration lives in the AppMaster, not in the filesystem. When an attempt finishes writing and needsTaskCommit returns true, the task reports commitPending over the task umbilical protocol, then polls canCommit. The AM grants commit to the first attempt that asks and refuses every other attempt of the same task. The refused attempt is killed, and its directory is deleted through abortTask. Only then does the winner call commitTask.
This is why speculative execution is safe with a correct committer and dangerous without one: the AM guarantees one commitTask per task, but it cannot undo a partial commitTask once it has started. If the winning attempt dies in the middle of commitTask, the AM must decide whether the task's output is still clean. With an atomic directory rename, the commit either happened or did not. With a sequence of per-file renames, it may be half done.
Worked example: the directory tree, step by step
FileOutputCommitter, the default for FileOutputFormat subclasses, implements the contract with renames under a _temporary directory inside the destination. The paths below follow the current mapred-default.xml description of both algorithms, for application attempt 1 and the first attempt of reduce task 3.
# Job output /data/out, application attempt 1, reduce task 3, first task attempt
# while the attempt runs (both algorithms)
/data/out/_temporary/1/_temporary/attempt_1727_0001_r_000003_0/part-r-00003
# v1: commitTask renames the attempt dir to a task dir under the job attempt
/data/out/_temporary/1/task_1727_0001_r_000003/part-r-00003
# v1: commitJob merges every task dir into /data/out, deletes _temporary, writes _SUCCESS
/data/out/part-r-00003
/data/out/_SUCCESS
# v2: commitTask renames each file straight into the destination
/data/out/part-r-00003 # visible before the job has finished
# v2: commitJob only deletes _temporary and writes _SUCCESSIn version 1 a task commit is a single directory rename from the attempt directory to a task directory. On HDFS that is one atomic NameNode operation, so a task is either committed or not. Job commit then walks every task directory and moves its files into the destination. The documentation notes that this merge is single-threaded and starts only after every task completes, so a job producing many files can spend minutes in commitJob (MAPREDUCE-4815). If the job fails before commitJob, the destination contains nothing new.
Version 2 moves each file straight into the destination at task commit. commitJob has almost nothing left to do, which removes the long serial tail. The cost is visibility and atomicity: committed tasks appear in /data/out while the job is still running, a failed job leaves partial output behind, and a task commit is a series of file renames rather than one directory rename.
Why v2 is a correctness trade, not a speed setting
The current mapred-default.xml lists 2 as the default. MAPREDUCE-7282 proposed deprecating v2 and making it non-default; it was closed as Won't Fix, so the default stays, but the reasons it was filed still apply. Because v2 task commit moves files one by one, an attempt that fails or is partitioned from the cluster in the middle of commitTask can leave some of its files in the destination while the AM gives the task to another attempt. If file names differ between attempts, the job ends with extra files; if they are the same, a late writer can overwrite a winner's file. Readers that start before _SUCCESS exists also see a partial dataset.
Spark reached the same conclusion from its side: SPARK-33019 made Spark set mapreduce.fileoutputcommitter.algorithm.version to 1 by default. On a plain HDFS cluster, the practical rule is to use v1, accept the slower job commit, and attack the file count instead, since a job writing tens of thousands of small files has a problem beyond commit time; see the HDFS small files problem. Use v2 only for jobs whose consumers wait for _SUCCESS, whose output paths are disposable on failure, and which run without speculation.
AM restart and recovery
The AppMaster is itself an attempt. mapreduce.am.max-attempts defaults to 2, so YARN restarts a failed AM once, as described in the YARN ApplicationMaster. A restarted AM asks the committer whether recovery is supported, and if it is, does not rerun tasks whose output was already committed. With v1 that means renaming each committed task directory from the previous application attempt's tree into the new one, which is why the paths include the attempt number. With v2 the files are already in the destination, so recoverTask has nothing to move.
Job commit is the delicate case. If the AM dies during commitJob, the new AM cannot know how far the merge got by looking at the destination alone. The MapReduce AM records commit start, success and failure markers in the job staging directory so that a restarted AM can tell whether a job commit was in progress, and isCommitJobRepeatable tells it whether retrying the commit is safe. Any code you add to commitJob must therefore tolerate being run twice.
Object stores break the rename assumption
Both FileOutputCommitter algorithms assume rename is cheap and atomic. On Amazon S3 a rename is a copy plus a delete per object, proportional to the data size and not atomic; the Hadoop S3A documentation says the classic committer risks loss or corruption there. S3A ships its own committers that use multipart uploads instead of renames, covered in Spark S3 optimisation.
For Azure ADLS Gen2 and Google Cloud Storage, Hadoop 3.3.5 added the intermediate manifest committer. Each task attempt writes its files under its attempt directory as usual, but task commit writes a manifest, a file listing those files, instead of renaming a directory. Job commit loads every committed manifest, creates the destination directories in parallel and renames the files with a thread pool. Its documentation describes v1 as correct but slow on deep trees, and v2 as unsafe on GCS, which lacks atomic directory rename. The manifest committer is bound per filesystem scheme, and it also works on HDFS.
# Job-level settings (Configuration, or -D on the command line)
mapreduce.fileoutputcommitter.algorithm.version=1 # mapred-default.xml lists 2
mapreduce.fileoutputcommitter.marksuccessfuljobs=true # write _SUCCESS at job commit
mapreduce.am.max-attempts=2 # AM restarts that can recover output
mapreduce.reduce.speculative=true # needs a committer that arbitrates
# Object stores: bind a committer per filesystem scheme (Hadoop 3.3.5 or later)
mapreduce.outputcommitter.factory.scheme.abfs=org.apache.hadoop.fs.azurebfs.commit.AzureManifestCommitterFactory
mapreduce.outputcommitter.factory.scheme.gs=org.apache.hadoop.mapreduce.lib.output.committer.manifest.ManifestCommitterFactory
mapreduce.manifest.committer.io.threads=64 # rename and create threads at job commitThe _SUCCESS file doubles as a fingerprint. The classic FileOutputCommitter writes it as an empty file. The manifest and S3A committers write a JSON document with the committer name, a list of committed files and I/O statistics, which makes it the quickest way to confirm which committer a job really used.
Writing a custom committer
Most custom committers exist to do one extra thing after the data is safe, typically to register a partition in a catalog or to flip a pointer that readers follow. The safe pattern is to extend FileOutputCommitter, let the parent finish its job commit, then publish, and to make the publish step idempotent because job commit can be retried.
public class CatalogPublishingCommitter extends FileOutputCommitter {
private final Path output;
public CatalogPublishingCommitter(Path output, TaskAttemptContext ctx) throws IOException {
super(output, ctx);
this.output = output;
}
@Override
public void commitJob(JobContext ctx) throws IOException {
super.commitJob(ctx); // files promoted and _SUCCESS written first
String partition = ctx.getConfiguration().get("publish.partition");
// Your metadata client. Must be idempotent: a restarted AM can repeat job commit.
CatalogClient.upsertPartition("events", partition, output.toString());
}
@Override
public void abortJob(JobContext ctx, JobStatus.State state) throws IOException {
super.abortJob(ctx, state); // delete _temporary; never touch the catalog
}
}
public class PublishingTextOutputFormat<K, V> extends TextOutputFormat<K, V> {
private OutputCommitter committer;
@Override
public synchronized OutputCommitter getOutputCommitter(TaskAttemptContext ctx)
throws IOException {
if (committer == null) {
committer = new CatalogPublishingCommitter(getOutputPath(ctx), ctx);
}
return committer;
}
}Overriding getOutputCommitter bypasses the per-scheme factory binding, so this committer always renames; use it on HDFS only. On object stores, wrap the committer the factory creates instead. Three mistakes are common. Publishing in commitTask exposes partial jobs. Publishing before super.commitJob can register a partition whose files are still in _temporary. And a non-idempotent publish, such as appending a row, duplicates the partition when a restarted AM repeats the commit.
Diagnostics
# Leftover job temporary data: a failed job, a job still running, or a job killed mid-commit
$ hdfs dfs -ls /data/out/_temporary
# Zero bytes: the classic FileOutputCommitter wrote it. JSON: a manifest or S3A committer did.
$ hdfs dfs -stat "%b" /data/out/_SUCCESS
# Committed output: one part file per task, no attempt directories
$ hdfs dfs -ls /data/outA leftover _temporary directory after a job reports success usually means a cleanup failure and is harmless but wasteful; one beside a missing _SUCCESS means the job failed or is still committing. A v2 job that failed leaves its committed files in the destination, so rerun it into a fresh path rather than on top of them.
Failure modes
- Readers see partial data: the consumer polls the directory instead of waiting for _SUCCESS, or the job runs v2 and has not finished.
- Duplicate or extra part files: v2 with speculation or task retries, where a dying attempt committed some files before losing.
- Job commit takes longer than the job: v1 with tens of thousands of output files, or any rename-based committer on an object store.
- Silent data loss on S3: the classic committer writing through S3A; switch to an S3A committer.
- Downstream registered an empty partition: a custom committer published before or without super.commitJob.
- Output from a dead AM: a recovered job reran tasks whose output was not recoverable, and a non-idempotent commit hook ran twice.
Trade-offs at a glance
| Committer | Task commit | Job commit | Failed job leaves | Use on |
|---|---|---|---|---|
| FileOutputCommitter v1 | Atomic directory rename on HDFS | Serial merge, slow with many files | Nothing in the destination | HDFS, by default |
| FileOutputCommitter v2 | Per-file renames into the destination | Delete temp, write _SUCCESS | Partial output | HDFS jobs that tolerate it |
| Manifest committer | Write a manifest | Parallel renames from manifests | Nothing in the destination | ABFS, GCS, also HDFS |
| S3A committers | Upload parts, defer completion | Complete multipart uploads | Nothing visible | S3 |
What to do next
- Check which algorithm your cluster and your Spark jobs use; mapred-default.xml and Spark disagree.
- Set mapreduce.fileoutputcommitter.algorithm.version=1 for any job that runs with speculation or feeds readers that do not wait for _SUCCESS.
- Make every downstream consumer wait for _SUCCESS, and check its size to confirm the committer in use.
- On ABFS or GCS, bind the manifest committer per scheme; on S3, use an S3A committer and never the classic one.
- Move any catalog or pointer update into commitJob after super.commitJob, and make it idempotent.
- Alert on _temporary directories older than your longest job, and reduce output file counts where job commit dominates runtime.