TT Lab
Get started
Learn Learning paths Courses

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

Continue in TT Lab

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

  1. Save the output of spark-submit --version (including standard error) to /root/spk/first/version.txt.
  2. Create /root/spk/first/count.py, count the rows of /data/shop/orders.csv (with a header row) under the application name spk-first-count, and write them to /root/spk/first/out/count.json as {"rows": 정수} (the value is an integer).
  3. Create /root/spk/first/lazy.py, read with an explicitly given schema under the application name spk-first-lazy, and stack filter (or where) and withColumn, but call no action. After you run it, that application must have 0 jobs.
  4. Create /root/spk/first/actions.py, call actions three or more times (for example count, take and write) under the application name spk-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.
  5. Create /root/spk/first/agg.py and write the number of orders per status as a CSV with a header row to /root/spk/first/out/by_status under the application name spk-first-agg. Write the number of stages that completed in that application to /root/spk/first/out/stages.txt as an integer.
  6. Create /root/spk/first/par.py so that it takes the application name as its first argument, run spk-first-par1 with --master local[1] and spk-first-par2 with --master local[2], and write master, default_parallelism and partitions (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.
  7. Create /root/spk/first/fail.py, make it read the nonexistent path /data/shop/order.csv under the application name spk-first-fail, catch the exception, and write the error condition name (getCondition()) to the first line of /root/spk/first/out/error.txt.
  8. 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

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.