Apache Spark — The answer to a slow job is in the plan and the event log
Write a partitioned lake, overwrite it, and reduce small files
Goal
You write the orders as a Parquet lake divided by date, and confirm what three of the four save modes (errorifexists, append, overwrite, ignore) actually do. You see the difference between a static overwrite and a dynamic partition overwrite in the number of directories, and use two knobs for reducing small files (gathering by the partition column and capping the rows per file).
Why it matters
A Spark write is not something that finishes in one go. Each task writes to a temporary path, and when the job succeeds, the commit protocol moves the result into place and leaves a _SUCCESS marker. That is why a half-written result cannot be read. In exchange, if you choose the save mode wrongly, this honest machine honestly causes an accident.
append puts the same data in twice on one retry. overwrite defaults to static, so what you wrote to fix only one day erases every date under the target path. A dynamic partition overwrite (spark.sql.sources.partitionOverwriteMode=dynamic) replaces only the partitions contained in the data being written. If you use overwrite once on a production lake without knowing this difference, several months vanish at once.
The number of files is the number of writing tasks × the number of partition values each task holds. If post-shuffle tasks each hold a little of every date, each date directory gets as many small files as there are tasks. You can control the size by gathering by the partition column first (repartition("칼럼"), where the argument is a column name) or by putting a cap on rows per file (maxRecordsPerFile).
Steps
- Put a function that reads the orders and attaches
order_datein /root/spk/write/common.py, and with /root/spk/write/lake.py (applicationspk-write-lake), write Parquet divided byorder_dateto /root/spk/write/lake/orders. - With /root/spk/write/exists.py (application
spk-write-exists), make it write again to the same path without giving a mode, and write the error condition name to /root/spk/write/out/exists.txt. - With /root/spk/write/append.py (application
spk-write-append), write the January orders to /root/spk/write/lake/append once with overwrite and once more with append, and write the number of rows read back and the number of distinctorder_idvalues to /root/spk/write/out/append.json. - With /root/spk/write/static.py (application
spk-write-static), write everything to /root/spk/write/lake/static, then write only the one day 2026-02-10 again with overwrite, and write the number of remaining partition directories to /root/spk/write/out/static.json. - With /root/spk/write/dynamic.py (application
spk-write-dynamic,partitionOverwriteMode=dynamic), write everything to /root/spk/write/lake/dynamic, then write with overwrite only the 2026-02-10 orders minus the canceled ones, and write the number of remaining partition directories to /root/spk/write/out/dynamic.json. - With /root/spk/write/small.py (application
spk-write-small), write the clicks withrepartition(40)to /root/spk/write/lake/clicks_many and withrepartition("page")to /root/spk/write/lake/clicks_few, each divided bypage, and write the number of data files of the two lakes to /root/spk/write/out/files.json. - With /root/spk/write/maxrec.py (application
spk-write-maxrec), applyrepartition("channel")to the orders and then write them divided bychannelto /root/spk/write/lake/by_channel withmaxRecordsPerFile=20000. - In /root/spk/write/report.md, write three sections:
## 저장 모드,## 정적과 동적 덮어쓰기and## 파일 수(use exactly these Korean headings in this order; they mean "Save modes", "Static and dynamic overwrite" and "Number of files"). Put the partition counts from steps 4 and 5 in the second section and the two file counts from step 6 in the third.
Notes
- Put the scripts in
/root/spk/writeand run them from there (from common import …). - The data files of a lake start with
part-._SUCCESSis the commit marker and.crcis a checksum (do not count them). - You turn on dynamic overwrite with the session setting (
spark.sql.sources.partitionOverwriteMode) or the write option (.option("partitionOverwriteMode", "dynamic")). - Common mistakes: rewriting one path in several steps and erasing the result of an earlier step (each step has a different path); believing a static overwrite will change only one day; and thinking maxRecordsPerFile sets the file size (in bytes).
- Official documentation: Generic Load/Save Functions — Save Modes · Bucketing, Sorting and Partitioning · Parquet Files · Configuration — spark.sql.sources.partitionOverwriteMode
Write a lake divided by date
In /root/spk/write/common.py, put a function that reads the orders with a schema and attaches order_date = to_date(order_ts), create /root/spk/write/lake.py with the application name spk-write-lake, and write Parquet with partitionBy("order_date") to /root/spk/write/lake/orders (mode("overwrite")).
With 90 days of dates there are 90 directories. _SUCCESS appears only if the job succeeds. A directory without this marker may be half written, so the reader checks it first. The grader looks at the number of directories, the number of rows and the marker.
The default mode stops: errorifexists
Create /root/spk/write/exists.py with the application name spk-write-exists, make it write again to the same path as step 1 without giving a mode, and write the exception's getCondition() to the first line of /root/spk/write/out/exists.txt. The lake from step 1 must remain as it is.
The default save mode is "error if it already exists". It is a safety device that keeps you from writing to the same path by mistake, and it is also why you need the habit of stating the mode explicitly in production jobs.
append: a retry makes it double
Create /root/spk/write/append.py with the application name spk-write-append, write the January 2026 orders to /root/spk/write/lake/append once with overwrite and once more with append, and write the number of rows read back and the number of distinct order_id values to /root/spk/write/out/append.json as {"rows": 정수, "distinct_orders": 정수} (both values are integers).
append does not look at the files that already exist and just adds new files. If you retry the job with the same input, the rows double, and since the file names differ, it looks fine on the surface. To write idempotently, use an overwrite (dynamic if possible) or a merge by key.
Static overwrite: trying to fix one day, you erase everything
Create /root/spk/write/static.py with the application name spk-write-static, write all the orders divided by order_date to /root/spk/write/lake/static, and then write only the one day 2026-02-10 again to the same path with overwrite (default settings). Write the number of remaining order_date= directories to /root/spk/write/out/static.json as {"partitions_after": 정수} (the value is an integer).
The default of partitionOverwriteMode is static. overwrite deletes everything under the target path before writing; it does not look at which dates are in the data being written. The result of this step is an accident. The purpose is to see it with your own eyes.
Dynamic partition overwrite: change only that day
Create /root/spk/write/dynamic.py with the application name spk-write-dynamic and the setting spark.sql.sources.partitionOverwriteMode=dynamic, write all the orders divided by order_date to /root/spk/write/lake/dynamic, and then write with overwrite to the same path only the 2026-02-10 orders minus the canceled ones (cancelled). Write the number of remaining order_date= directories to /root/spk/write/out/dynamic.json as {"partitions_after": 정수} (the value is an integer).
In dynamic mode, only the partitions contained in the data being written (here, one day) are swapped. The other 89 days stay as they are. The grader checks the number of directories, whether that day has no canceled orders, and whether the other days are as in the original.
Small files: number of tasks × partition values
Create /root/spk/write/small.py with the application name spk-write-small, write /data/clicks/clicks.jsonl with repartition(40) to /root/spk/write/lake/clicks_many and with repartition("page") to /root/spk/write/lake/clicks_few, each with partitionBy("page"), and write the number of data files (page=*/part-*) of the two lakes to /root/spk/write/out/files.json as {"many": 정수, "few": 정수} (both values are integers).
One writing task opens one file for each partition value it holds. If 40 tasks all hold all seven pages, that is up to 280. If you gather by the partition column first, one page goes into one task and there is one file per page (in exchange, if one page is large, that task becomes heavy).
Put a cap on rows per file
Create /root/spk/write/maxrec.py with the application name spk-write-maxrec, apply repartition("channel") to the orders, and then write them with partitionBy("channel") and option("maxRecordsPerFile", 20000) to /root/spk/write/lake/by_channel.
When you gather by channel, one channel goes into one task and becomes one file, but the largest channel has over 170,000 rows. With a cap, the task opens a new file every 20,000 rows. The grader uses the Parquet metadata to check that the number of files per channel is the ceiling of (rows ÷ 20000) and that no file has more than 20,000 rows.
Turn the write rules into team rules
In /root/spk/write/report.md, write three sections: ## 저장 모드, ## 정적과 동적 덮어쓰기 and ## 파일 수 (use exactly these Korean headings in this order; they mean "Save modes", "Static and dynamic overwrite" and "Number of files"). Put the remaining partition counts from steps 4 and 5 in the second section and the two file counts from step 6 in the third section.
For a job that writes to a production lake, try writing as rules which mode is the default and what is forbidden. Copy the numbers from your own result files.