Apache Spark — The answer to a slow job is in the plan and the event log
Launch a first job and check it in the event log — transformations are lazy, actions do the work
Goal
Using Spark in local mode, read and count 300,000 orders, and compare an application that only stacks transformations with one that calls actions, using the event log. You confirm with numbers when jobs, stages and tasks are created and what determines the number of partitions.
Why it matters
When you first read Spark code, it looks as if every line runs on the spot. It does not. A transformation such as filter, withColumn or select only adds one line to the plan, and when you call an action such as count, take or write, the whole plan turns into a job and runs. If you do not know this, you miss that the spot that seems to say "reading is slow" is in fact where all the transformations stacked up earlier run at once.
You can see that difference with your own eyes. When an application finishes, Spark leaves a record of what it did in the event log. How many jobs started, how many stages there were, how many tasks there were, and which physical plan was used are all in there, one line of JSON each. This is also what you look at when tracing a finished job in production (the History Server reads this file).
In local mode, a single driver JVM is itself the executor. local[2] means using two cores, and this number determines both the default parallelism and how many pieces a file is split into.
Steps
- Save the output of
spark-submit --version(including standard error) to /root/spk/first/version.txt. - Create /root/spk/first/count.py, count the rows of
/data/shop/orders.csv(with a header row) under the application namespk-first-count, and write them to /root/spk/first/out/count.json as{"rows": 정수}(the value is an integer). - Create /root/spk/first/lazy.py, read with an explicitly given schema under the application name
spk-first-lazy, and stackfilter(orwhere) andwithColumn, but call no action. After you run it, that application must have 0 jobs. - Create /root/spk/first/actions.py, call actions three or more times (for example
count,takeandwrite) under the application namespk-first-actions, and write 100 rows as JSON to /root/spk/first/out/sample. Count the job start events in that application's event log and write them as a single integer to /root/spk/first/out/jobs.txt. - Create /root/spk/first/agg.py and write the number of orders per
statusas a CSV with a header row to /root/spk/first/out/by_status under the application namespk-first-agg. Write the number of stages that completed in that application to /root/spk/first/out/stages.txt as an integer. - Create /root/spk/first/par.py so that it takes the application name as its first argument, run
spk-first-par1with--master local[1]andspk-first-par2with--master local[2], and writemaster,default_parallelismandpartitions(the number of partitions of the DataFrame that read the orders file) to /root/spk/first/out/spk-first-par1.json and /root/spk/first/out/spk-first-par2.json respectively. - Create /root/spk/first/fail.py, make it read the nonexistent path
/data/shop/order.csvunder the application namespk-first-fail, catch the exception, and write the error condition name (getCondition()) to the first line of /root/spk/first/out/error.txt. - In /root/spk/first/report.md, write your numbers in three sections:
## 잡과 스테이지,## 게으른 실행and## 파티션(as section titles, use exactly these Korean headings in the same order, which correspond to "Jobs and stages", "Lazy execution" and "Partitions"). Put the number of jobs from step 4 in the first section and the two partition counts from step 6 in the third section.
Notes
- Run it like
spark-submit /root/spk/first/count.py. The default settings are in/opt/spark/conf/spark-defaults.conf(local[2], driver memory 1g, event log on). - When an application finishes, its event log is left as
/root/spark-events/local-<숫자>(the placeholder is a number). While it is running,.inprogressis appended to the name. To find out which file belongs to which application, usegrep -l '"App Name":"spk-first-actions"' /root/spark-events/*. - Count the jobs with
grep -c '"Event":"SparkListenerJobStart"' <파일>(the placeholder is the log file), and the stages by counting"Event":"SparkListenerStageCompleted". You can also view the contents withjq -c 'select(.Event=="SparkListenerJobStart")'. - Common mistakes: leaving out
spark.stop()so the log stays as.inprogress; leaving outheader=Trueon a CSV with a header row so the header row is counted as a data row; and reading without a schema so that a job that checks the header row appears (step 3 must have 0 jobs). - Official documentation: Cluster Mode Overview · RDD Programming Guide — RDD Operations · Monitoring — Viewing After the Fact · Submitting Applications — Master URLs
Check which Spark you have
Save the output of spark-submit --version, with standard error merged in, to /root/spk/first/version.txt.
Spark prints the version information to standard error. If you do not merge it with 2>&1, the file ends up empty. Along with the version number, you can also see which Scala and Java it runs on. Spark 4 requires Java 17 or later.
The first job: counting 300,000 rows
Create /root/spk/first/count.py under the application name spk-first-count, count the rows of /data/shop/orders.csv (with a header row), and write them to /root/spk/first/out/count.json as {"rows": 정수} (the value is an integer). Run it with spark-submit.
SparkSession.builder.appName(...) sets the application name. The grader checks whether the event log of that name was closed properly and whether the number in the file equals the number of rows in the original. If you count the header row as a row, you get one too many.
Stacking only transformations launches no job
Create /root/spk/first/lazy.py under the application name spk-first-lazy, read the orders with spark.read.csv and an explicitly given schema, and stack filter (or where) and withColumn, but call no action at all. After you run it, the event log of that application must show 0 jobs.
If you do not give a schema, Spark launches a small job to check the header row. If you pass a DDL string such as schema="order_id STRING, ...", that job disappears too. Printing the schema is fine: the driver already knows the schema, so no data is read.
Every action launches a job: count them from the event log
Create /root/spk/first/actions.py under the application name spk-first-actions, call actions three or more times, and use one of them to write 100 rows as JSON to /root/spk/first/out/sample. After you run it, count the SparkListenerJobStart events in that application's event log and write them as a single integer to /root/spk/first/out/jobs.txt.
Do not assume that one action is one job. Some actions create two or more jobs, such as an action that writes files or schema inference. That is why what you count is the log, not the code. If you ran it several times under the same name, you must count the most recent log.
A wide transformation splits the stages
Create /root/spk/first/agg.py under the application name spk-first-agg and write the number of orders per status as a CSV with a header row (columns: status, count) to /root/spk/first/out/by_status. Count the SparkListenerStageCompleted events in that application's event log and write them as an integer to /root/spk/first/out/stages.txt.
groupBy has to gather the same key into one task, so a shuffle occurs, and the parts before and after the shuffle become different stages. The grader checks whether your CSV matches the values counted from the original, whether that application actually had a stage that used a shuffle, and whether the number of stages you wrote matches the log.
The number of cores determines the number of partitions
Make /root/spk/first/par.py take the application name as its first argument, run it twice with spark-submit --master 'local[1]' par.py spk-first-par1 and spark-submit --master 'local[2]' par.py spk-first-par2, and write {"master": 문자열, "default_parallelism": 정수, "partitions": 정수} (a string and two integers) to /root/spk/first/out/spk-first-par1.json and /root/spk/first/out/spk-first-par2.json. partitions is the rdd.getNumPartitions() of the orders DataFrame read with a schema.
How many pieces a file is split into is not determined by spark.sql.files.maxPartitionBytes (default 128MB) alone. When there are several cores, if the total size divided by the number of cores is smaller, the file is cut to that size. This is a rule to keep cores from sitting idle. See how a 16MB file differs between one core and two.
A failure is also an application: read the error condition name
Create /root/spk/first/fail.py under the application name spk-first-fail, make it read the nonexistent path /data/shop/order.csv, catch the exception, and write the error condition name returned by getCondition() to the first line of /root/spk/first/out/error.txt. The application must exit normally with spark.stop().
A nonexistent path shows up not when you call an action but the moment you read, because the file listing is checked when the plan is built. A Spark 4 error has, besides a sentence for people to read, a name for machines to read (the error condition). It is safer to base alerts and retry rules on this name rather than on the sentence.
Record what you saw as numbers
In /root/spk/first/report.md, write three sections: ## 잡과 스테이지, ## 게으른 실행 and ## 파티션 (use exactly these Korean headings in this order; they mean "Jobs and stages", "Lazy execution" and "Partitions"). Put the number of jobs counted in step 4 in the first section and the two partition counts from step 6 in the third section, as numbers.
Do not just copy the numbers; add a sentence or two on why each number is what it is. How many jobs the three actions produced, why the application that only stacked transformations had 0, and how the number of cores changed the number of partitions are everything this lab showed.