Apache Spark — The answer to a slow job is in the plan and the event log
Writes cannot be undone, so decide the mode and file layout first
In one line
A Spark write is determined by three things: the save mode (what to do if it already exists), the partition directories (in what shape to lay it out) and the commit protocol (how to announce that everything has been written). If you choose wrongly on all three, data disappears or multiplies without an error.
Why learn writing separately
With reads and transformations, the source stays intact however many times you get it wrong. Writes are different. One overwrite erases three months of partitions, and a retried append puts the same day in twice. And both end in success. The job is green and only the report is wrong.
The load and save documentation warns about this first when it explains save modes. Save modes do not use locks and are not atomic. And an overwrite deletes the existing data before writing the new data. This means that if the job dies during an overwrite, there is a moment with neither the old data nor the new data.
How it works: the four save modes
The same documentation lists four modes.
- errorifexists (or error): the default. If data already exists at the path, it throws an exception. In the lab it shows up as the
PATH_ALREADY_EXISTScondition. - append: adds new files next to the existing data. If you run it twice with the same data, you get exactly double.
- overwrite: deletes the existing data and writes anew.
- ignore: does nothing if it already exists. It resembles SQL's
CREATE TABLE IF NOT EXISTS.
That the default is an error is good design. It keeps a write that did not say what to do from touching someone else's data. The problem starts when people habitually tack on overwrite to get rid of that error.
Partition directories and the number of files
If you write with partitionBy("day"), a directory like day=2026-01-03/ is created for each value, and that column goes in the path rather than inside the file. The reader restores the value from the path, and if a condition is placed on that column, it skips the directory whole. The same documentation notes that because this approach creates a directory structure, it is not suited to a column with a very large number of values. If you partitionBy user ID, as many directories are created as there are users.
The number of files is a quieter trap. One write task writes one file per partition value of the rows it holds. If all 200 post-shuffle tasks hold a little of 30 days, there can be up to 200 × 30 = 6,000 files: 200 files of a few KB each per day.
(df.repartition("day") # 같은 날은 한 태스크로 — 날마다 파일 하나
.write.partitionBy("day")
.option("maxRecordsPerFile", 50000) # 너무 큰 날은 5만 줄씩 끊는다
.mode("overwrite")
.option("partitionOverwriteMode", "dynamic")
.parquet("/data/lake/orders"))
If you repartition by the partition column before writing, the rows of the same day gather in one task and there is one file per day. In exchange, a large day becomes one huge file. What puts a cap on that is spark.sql.files.maxRecordsPerFile in the configuration documentation. It is the maximum number of records to write to one file, and the default 0 means no limit. You can also give it as a write option.
The alternative when you want to divide by a column with many values is bucketing. According to the same documentation, bucketBy hashes the data into a fixed number of buckets regardless of the number of distinct values, so it can be used even on a column whose distinct values grow without bound. In exchange, buckets and sorting apply only to persistent tables (saveAsTable). You cannot leave buckets with save(), which writes only files to a path. A way to remember the division: a column with few distinct values that always goes into query conditions, like a date, gets partitionBy, and a column with many distinct values that is used as a join key, like a user ID, gets bucketBy.
Static overwrite and dynamic overwrite
When you overwrite a partitioned path, what gets deleted? spark.sql.sources.partitionOverwriteMode in the configuration documentation decides this, and the default is STATIC. Static mode deletes the partitions matching the target in advance before writing. If you overwrite the whole path with a DataFrame, the target is the whole path. This is the accident where you overwrite with a DataFrame containing only one day's rows to fix that day and all the other dates disappear.
Dynamic mode does not delete in advance and replaces only the partitions where data was actually written. If you overwrite with a one-day DataFrame, only that day changes and the other days stay. The documentation says the write option partitionOverwriteMode takes precedence over the session setting. So it is safer to write it as an option at the place of writing, as in the example above. A session setting can be changed by someone else, but an option written in the code travels with that write.
The commit protocol and _SUCCESS
Several tasks write to one directory at the same time, so when can the reader trust that everything has been written? A task writes not to the final location but to its own place under _temporary/. When a task succeeds, its output is committed, and when all the tasks have finished, the job commit moves the outputs to the final location. Hadoop's mapred-default.xml explains that the job commit of algorithm version 1 merges and moves the task outputs, deletes _temporary, and then writes _SUCCESS. So _SUCCESS is a marker saying "this directory was written to the end". It is not in the directory of a failed job.
Spark chooses this commit algorithm by configuration. In the configuration documentation, the default of spark.hadoop.mapreduce.fileoutputcommitter.algorithm.version is 1, and it says 2 can cause correctness problems such as MAPREDUCE-7282. Version 1 relies on moving (renaming). The cloud integration documentation warns that renaming in an object store is very slow and that if it fails the state becomes unknown. This is why something cheap on a local disk or HDFS becomes expensive on S3.
What it looks like in the field
First, trying to fix one day, you delete everything. This is the classic case of a static overwrite. In code that overwrites a partitioned lake, put the dynamic mode in as an option.
Second, a retry makes it double. If a job that writes with append fails midway and runs again, the same data goes in again on top of the part committed earlier. For a job that can be rerun by day, a dynamic overwrite of that day's partition is safer than append, because the result is the same however many times you run it.
Third, the next job reads a half-written directory. If a following job starts reading as soon as it sees files, it sees a state in the middle of writing. One rule, to read after checking _SUCCESS, prevents this. Also remember that it is created at the top path of the write, not in each partition directory. The reader looks for this marker not in the date directories but at the top of the lake.
What really matters in practice
- The default mode is an error. Do not tack on overwrite out of habit.
- Save modes are not atomic. overwrite deletes before writing.
- For overwriting a partitioned lake, write the dynamic mode as an option at the place of writing.
- The number of files is the number of tasks × partition values. Repartition by the partition column and cap with maxRecordsPerFile.
- _SUCCESS is the marker that it was written to the end. The reader checks it.
What you will do in the next lab
You write a lake by partitioning the orders by date with partitionBy and check the error condition raised when you write again to the same path in the default mode. After an overwrite, you write once more with append and count that the rows double while the distinct order numbers stay the same, see that the other dates disappear when you do a static overwrite with one day's data, and then fix only that day with a dynamic overwrite. Finally, you compare the number of files of a lake written from clicks scattered into 40 pieces with a lake rewritten divided by the partition column, and put a cap on the number of records that go into one file with maxRecordsPerFile.