Spark was designed around HDFS, a real file system where renaming a directory is one cheap, atomic metadata operation. Amazon S3 is an object store. It has no directories, only keys with slashes in them; it has no rename, only copy then delete; listing is a paged HTTP call; and every request costs money and counts against rate limits. The S3A connector hides those differences well enough that a job pointed at an s3a:// path runs unchanged. Hiding them is not the same as removing them, and most slow or flaky Spark-on-S3 jobs come from a handful of places where the object store leaks through.
This article covers the four that matter most: committing output safely and quickly, reading columnar files efficiently, staying inside S3's request limits, and producing a file layout that keeps the next job fast. It assumes Spark 3 or later with the Hadoop S3A connector, the usual setup outside managed platforms that ship their own committers.
Why the classic commit protocol fails on S3
Spark writes output through a commit protocol. Each task writes to a temporary attempt directory; when a task succeeds, its output is promoted; when the whole job succeeds, everything appears in the destination. Hadoop's classic FileOutputCommitter implements promotion with renames. Algorithm version 1 renames each task's directory into a job temporary directory at task commit, then renames again at job commit. Version 2 renames task output straight into the destination at task commit.
On S3, each rename is a copy of every object followed by a delete, done file by file by the client. Copy time grows with data size, so a job that wrote 2 TB in ten minutes can spend much longer committing. Worse, the copy is not atomic: if the driver dies part-way, the destination holds some files and not others. With version 2, partial output from a failed job is visible to readers. The Hadoop documentation is direct: using the classic committer to write to S3 risks loss or corruption of data.
S3 has offered strong read-after-write consistency since December 2020, which removed an older class of bugs where listings missed new files. It did not add rename. The fix is a committer that never renames.
Why the classic commit protocol fails on S3
Spark writes output through a commit protocol. Each task writes to a temporary attempt directory; when a task succeeds, its output is promoted; when the whole job succeeds, everything appears in the destination. Hadoop's classic FileOutputCommitter implements promotion with renames. Algorithm version 1 renames each task's directory into a job temporary directory at task commit, then renames again at job commit. Version 2 renames task output straight into the destination at task commit.
On S3, each rename is a copy of every object followed by a delete, done file by file by the client. Copy time grows with data size, so a job that wrote 2 TB in ten minutes can spend much longer committing. Worse, the copy is not atomic: if the driver dies part-way, the destination holds some files and not others. With version 2, partial output from a failed job is visible to readers. The Hadoop documentation is direct: using the classic committer to write to S3 risks loss or corruption of data.
S3 has offered strong read-after-write consistency since December 2020, which removed an older class of bugs where listings missed new files. It did not add rename. The fix is a committer that never renames.
How the S3A committers avoid rename
All three S3A committers rely on a feature of S3 multipart upload: you can upload every part of an object and then decide later whether to complete the upload, which makes the object visible in one step, or abort it, which discards the parts. The committers upload task output as uncommitted multipart uploads and defer completion to job commit. Completing an upload is a single small request per file, no matter how large the file.
| Committer | Where task output goes | Strengths | Watch out for |
|---|---|---|---|
directory | Local disk of the executor, uploaded at task commit | Mature; tolerant of any destination layout | Needs local disk for a task's full output; conflict-mode applies to the whole destination |
partitioned | Same as directory | Conflict resolution per partition, so replace overwrites only the partitions written | Spark-oriented; still stages to local disk |
magic | Directly to S3, under a special __magic path that the connector redirects | No local staging; fastest for large output | Requires magic support in the S3A filesystem, on by default since Hadoop 3.3.1 |
With the magic committer, a task that opens a file under the job's magic path actually starts a multipart upload to the final key, and writes a small .pending manifest describing it. At task commit those manifests are aggregated; at job commit the driver reads them and completes every upload. Failed and speculative task attempts never have their uploads completed, so their data never appears. The Hadoop documentation notes that with S3 now consistent there are fewer reasons not to use the magic committer, and suggests trying both.
One consequence for operations: a job killed between task and job commit leaves uncommitted multipart uploads behind. They are invisible in listings but billed as storage until aborted.
Worked example: switching a daily job to the magic committer
Enabling the magic support in S3A does not make Spark use it. Spark keeps its default commit protocol until you bind it to the Hadoop path-output committers, which live in the spark-hadoop-cloud module. Without that module on the classpath, the classes below fail to load.
# spark-defaults.conf -- needs the spark-hadoop-cloud module on the classpath
# (artifact spark-hadoop-cloud_2.13, or the jar shipped with your distribution)
spark.hadoop.fs.s3a.committer.name magic
spark.sql.sources.commitProtocolClass org.apache.spark.internal.io.cloud.PathOutputCommitProtocol
spark.sql.parquet.output.committer.class org.apache.spark.internal.io.cloud.BindingParquetOutputCommitter
# Read path for columnar formats: footers and column chunks are random reads
spark.hadoop.fs.s3a.experimental.input.fadvise random
# Connection pool and upload threads: size them together, with headroom over executor cores
spark.hadoop.fs.s3a.connection.maximum 200
spark.hadoop.fs.s3a.threads.max 64
# Throttling: retry 503 Slow Down responses with back-off rather than failing tasks
spark.hadoop.fs.s3a.retry.throttle.limit 20
spark.hadoop.fs.s3a.retry.throttle.interval 1000msThe connection pool and thread numbers above are illustrative, not defaults. Defaults vary between Hadoop releases; check yours, then size the pool larger than the number of concurrent tasks per executor plus upload threads, because an executor that runs out of connections stalls on a pool wait that looks like S3 latency.
A job then writes exactly as before.
from pyspark.sql import SparkSession, functions as F
import json
spark = SparkSession.builder.appName("events-daily").getOrCreate()
events = (spark.read.parquet("s3a://lake/raw/events/")
.where(F.col("event_date") == "2026-10-01") # partition filter: prunes the listing
.drop("event_date")) # the value lives in the path
out = "s3a://lake/curated/events_daily/event_date=2026-10-01/"
(events
.repartition(64) # 64 files of a sensible size, not 4,000 tiny ones
.write.mode("overwrite")
.option("maxRecordsPerFile", 5_000_000)
.parquet(out))
# An S3A committer writes a JSON _SUCCESS manifest; the classic committer writes 0 bytes.
manifest = json.loads(
spark.sparkContext.wholeTextFiles(out + "_SUCCESS").first()[1])
print(manifest["committer"], len(manifest.get("filenames", [])), "files")Check the _SUCCESS file on the first run. S3A committers write a JSON manifest naming the committer and the files committed; a zero-byte _SUCCESS means the classic committer still ran and your bindings did not take effect. Make that check part of the job's tests, because a misconfigured classpath falls back silently.
The dynamic partition overwrite trap
Many pipelines rewrite only the partitions present in the new data by setting spark.sql.sources.partitionOverwriteMode=dynamic and writing in overwrite mode. Spark implements this by staging output and renaming partitions into place, and its cloud-integration documentation states that the S3A committers do not meet the conditions it needs. The write fails with an error saying the path output committer does not support dynamic partition overwrite.
There are three workable answers. Write each partition to its explicit path, as the example does, so a plain overwrite of that path is enough. Use the partitioned committer with conflict mode replace, which replaces only the partitions the job wrote. Or use a table format: Apache Iceberg and Delta Lake commit by writing new data files and atomically swapping table metadata, so they need no rename and support partition-level overwrites, concurrent writers and time travel. For new data lakes the table format is usually the right default.
Reading efficiently: seeks, footers and vectored reads
Parquet and ORC are read out of order. A reader fetches the footer at the end of the file, then jumps to the column chunks it needs. Every jump on S3 is a new ranged GET, and if the stream was opened expecting a sequential read the connector may have requested far more bytes than will be used, then must abort that connection.
S3A's fs.s3a.experimental.input.fadvise sets the policy: sequential for whole-file reads such as CSV or JSON, random for columnar formats, and normal, which starts sequential and switches to random after a backward seek. For Parquet-heavy jobs, setting random avoids the first wasted read. Newer Hadoop releases also implement vectored reads, where a reader hands the connector a list of byte ranges and S3A merges nearby ranges and fetches them in parallel; fs.s3a.vectored.read.min.seek.size and fs.s3a.vectored.read.max.merged.size tune the merging. Whether your Parquet reader uses the vectored API depends on its version, so measure rather than assume.
The biggest read optimisation is reading less. Partition pruning removes whole directories from the listing and predicate pushdown skips row groups using footer statistics; both are covered in partition pruning in Spark. On S3 each skipped file is also a skipped HEAD, LIST and GET.
Request limits and throttling
S3 scales request throughput per prefix: AWS documents at least 3,500 PUT, COPY, POST or DELETE and 5,500 GET or HEAD requests per second per partitioned prefix, and it partitions busy prefixes further as load grows. A sudden burst against one prefix, for example 4,000 tasks starting at once and all reading files under one date path, can exceed what that prefix currently supports and draw HTTP 503 Slow Down responses.
S3A retries throttled requests with back-off, which is why the example raises the throttle retry settings: a retried read is slower but a failed task is far worse. Beyond that, spread hot data across prefixes where the layout allows, avoid jobs that list or HEAD the same keys repeatedly, ramp large jobs up rather than starting every executor at once, and remember that other jobs share the same limits. Fewer, larger files reduce request counts more than any setting.
Measure before and after any change. S3 request metrics in CloudWatch, which you enable per bucket or per prefix filter, show request counts by operation, 4xx and 5xx errors and first-byte latency; server access logs give per-request detail when you need to attribute load to a job. A useful habit is to record GET, LIST and PUT counts for one representative run of each important job. When a change such as a new fadvise policy or a compaction job lands, the request counts should fall along with the run time; if run time falls but requests rise, you have moved cost rather than removed it.
File layout: the next job's performance
Each output file costs at least one PUT to write, one GET and often one HEAD to read, and an entry in every listing. A job writing 200 partitions with 500 tasks each can create 100,000 files of a few hundred kilobytes, and every downstream job then pays per file. Aim for files in the low hundreds of megabytes for Parquet, set by repartitioning on the partition columns before writing, capping with maxRecordsPerFile, or letting adaptive query execution coalesce shuffle partitions. The trade-offs between repartition and coalesce are covered in coalesce versus repartition.
Compaction is the long-term answer for streaming and incremental writers that produce small files by nature. Table formats ship compaction procedures; with plain Parquet you schedule a job that rewrites small files in a partition into larger ones and swaps them in during a quiet window.
Failure modes
- Commit takes longer than the job. The classic committer is renaming. Check for a zero-byte
_SUCCESSand fix the bindings. - ClassNotFoundException for PathOutputCommitProtocol. The spark-hadoop-cloud jar is missing.
- Dynamic partition overwrite error after enabling S3A committers. Expected; see the alternatives above.
- Storage bill grows with no visible data. Uncommitted multipart uploads from killed jobs. Add a lifecycle rule to abort them.
- Bursts of 503 Slow Down. Too many requests on one prefix; raise throttle retries, reduce file counts, ramp up executors.
- Executors idle waiting on connections. Connection pool smaller than concurrent readers plus upload threads.
- Directory staging fills local disk. The directory and partitioned committers buffer whole task output locally; give executors disk or use magic.
{
"Rules": [{
"ID": "abort-stale-multipart-uploads",
"Status": "Enabled",
"Filter": {"Prefix": ""},
"AbortIncompleteMultipartUpload": {"DaysAfterInitiation": 7}
}]
}Apply that rule with aws s3api put-bucket-lifecycle-configuration on every bucket Spark writes to. Seven days comfortably exceeds any job's commit window. Impala on S3 faces related problems from the query side; see Impala on HDFS and S3.
What to do next
- Find which committer your jobs use today: read one
_SUCCESSfile. - Add spark-hadoop-cloud and the three committer settings to one job, run it, and confirm the JSON manifest.
- Replace any use of dynamic partition overwrite with explicit partition paths, the partitioned committer in replace mode, or a table format.
- Set fadvise to random for Parquet and ORC readers and compare the job's S3 request counts before and after.
- Size the S3A connection pool and threads against executor cores.
- Add the abort-incomplete-multipart-upload lifecycle rule to every output bucket.
- Measure average output file size per table and repartition or compact anything averaging under about 64 MB.