SpiralTrain
Exercises › Block 2 · Exercise 1

Block 2 · Exercise 1

Spark Architecture

Open the notebook 13-spark-architecture-starter in your workspace and run the setup cell, which starts the Spark session. Most steps end in a question whose answer is on your own screen, in the Spark UI, so open the UI in a second browser tab and keep it there.

Do one part at a time, 1a, 1b and so on, and compare with the expected output where there is one. The solution is in 13-spark-architecture-solution, under every part on the exercise site, and at the back of the exercises PDF.

Step 1Spark UI

1a

Find the Spark UI for your session. In a Fabric notebook it sits under the session indicator at the bottom of the notebook, and in the Monitor hub under your running application. Open it in a second browser tab.

Hint 1

The slide Reading the Spark UI lists the tabs and what each one shows.

Show solutionHide solution

Session indicator at the bottom of the notebook, then the link to the Spark UI. The Monitor hub lists the same application under Monitor.

1b

Open the Executors tab. Write down how many executors you have, how many cores each has, and how much memory is listed against each. Run the cell of Step 1, which prints what the session reports about itself, and compare.

Hint 1

The Executors tab has a row for the driver as well. Count the executors separately from it.

Check your output
executor entries : <a small whole number>
hosts and memory : <a map with one entry per executor and the driver>
cores for tasks  : <a whole number>

The numbers depend on the pool your session runs on.

Show solutionHide solution
python
status = spark.sparkContext._jsc.sc().getExecutorMemoryStatus()
print("executor entries :", status.size())
print("hosts and memory :", status)
print("cores for tasks  :", spark.sparkContext.defaultParallelism)

A local Spark reports one executor entry, the driver. A Fabric session on the course pool reports the driver plus the pool's executors. Keep the three numbers : they explain every task count you see later.

Step 2Single stage

2a

Run the cell that counts the flights. Go to the Jobs tab. One count() started two jobs. How many stages does the older of the two have?

Hint 1

The slide Single Stage Job shows what a job without a shuffle looks like.

Show solutionHide solution

One stage : the tasks that read the files and count each partition.

2b

Click into the older job, then into its stage. How many tasks are in that stage? Does the number match the number of files, the number of cores, or neither? Then print the partition count of flights and compare.

Hint 1

A task processes one partition. flights.rdd.getNumPartitions() gives the count.

Check your output
partitions : <a whole number>
rows       : <the number of flights>

The partition count equals the task count of the stage. Both depend on the data files and the session.

Show solutionHide solution
python
print("partitions :", flights.rdd.getNumPartitions())
print("rows       :", flights.count())

The task count is the partition count of the parquet files. Nothing you wrote said how many partitions to use, so Spark chose when it read them.

2c

Open the newer job. Which of its stages is marked skipped, and how many tasks ran in the other one?

Hint 1

The slide Skipped Stages and Caching says when Spark lists a stage as skipped.

Show solutionHide solution

A count on a DataFrame counts each partition, then sends the partial counts through a shuffle to one task that adds them up. With adaptive query execution on, which is the default, Spark runs the work before a shuffle as a job of its own, measures what it wrote, and then plans the rest. The later job lists the stage that already ran as skipped. The stage that ran in the later job has one task, the one that adds up the partial counts. The job with one stage is the first of the two.

Step 3Shuffle

3a

Run the cell that groups the flights by carrier and shows the counts. Go to the Jobs tab. The cell started two jobs again, the same shape as in Step 2. Open the newer job. How many stages does it have, and which one is skipped?

Hint 1

Compare with what you saw for the newer job in 2c.

Show solutionHide solution

Two stages. The stage before the shuffle already ran in the older job and is listed as skipped ; the stage after it ran.

3b

Which operation caused the shuffle this time, and which one caused it in Step 2? Print the plan of by_carrier and find the shuffle in it.

Hint 1

by_carrier.explain() prints the plan. A shuffle shows up as a node named Exchange.

Check your output
== Physical Plan ==
AdaptiveSparkPlan ...
+- HashAggregate(keys=[carrier...], functions=[count(1)])
   +- Exchange hashpartitioning(carrier...)
      +- HashAggregate(keys=[carrier...], functions=[partial_count(1)])
         +- FileScan parquet ...

The exact text and numbers vary. The shape is an aggregate, an Exchange, and an aggregate again.

Show solutionHide solution
python
by_carrier.explain()

groupBy adds the shuffle, because every row for one carrier has to meet on one executor. The Exchange node in the plan is that shuffle. In Step 2 the count itself needed one.

3c

Open the Stages tab and find the stage with the shuffle. How much data was written in Shuffle Write, and how much was read on the other side? Compare the shuffle write with the size of the data the job read.

Show solutionHide solution

The shuffle write is far smaller than the data the job read. Spark aggregates within each partition before the shuffle, so only the partial counts cross the network.

Step 4Slowest task

4a

Stay on the stage you opened in 3c. Look at the task table at the bottom, or at the summary metrics above it. What was the duration of the Max task, and of the Median task? How far apart are they? Write both numbers down : Step 6 compares against them.

Hint 1

The summary metrics table has a row per statistic, with Min, 25th percentile, Median, 75th percentile and Max as columns.

Show solutionHide solution

A large gap between median and max in one stage means the work was not spread evenly : the data is skewed.

4b

Count the flights per origin airport and show the ten busiest.

Hint 1

groupBy("origin").count(), then sort on the count, descending.

Check your output
+------+------+
|origin| count|
+------+------+
|   ...|   ...|
...
+------+------+
only showing top 10 rows

Ten airports, the count falling down the column. Which airports and what counts depend on the dataset size.

Show solutionHide solution
python
(flights.groupBy("origin").count().orderBy(F.desc("count"))).show(10)

The row counts per airport are what the gap between median and max reflects : a few hub airports carry far more flights than the rest.

Step 5Skipped stage

5a

Run the cell that caches by_carrier and then calls count() on it twice. Look at the jobs this cell started in the Jobs tab. How many jobs did the first count() start, and how many did the second?

Hint 1

The jobs of one cell sit together in the Jobs tab, the newest at the top.

Show solutionHide solution

On Fabric the first count ran as four jobs and the second as two. Your numbers may differ with the session, so compare the first against the second : the second starts fewer.

5b

Open the first job of the second count. Which stage is marked skipped, and how many tasks ran in the stage after it?

Hint 1

Two things cause skipped stages here. One is the same as in 2c. Read the slide Skipped Stages and Caching for the other.

Show solutionHide solution

Skipped stages appear in both counts, for two reasons. Inside one action, every job after the first lists the stages already run, as in Step 2. In the second count, the stage that reads the flights is skipped as well, because the cache already holds what it produced ; the stage after it reads the cached partitions instead.

A skipped stage is Spark reusing what it already has : shuffle output from an earlier job, or a cached copy. Without the cache, a second count() reads the flights again, because count() builds a new query each time.

5c

Open the Storage tab. What is listed there, how many partitions does it hold, and what fraction is cached in memory?

Show solutionHide solution

The Storage tab lists the cached by_carrier result, with its partitions and the fraction held in memory.

5d

Empty the cache again and watch the Storage tab empty.

Hint 1

The opposite of cache() is a method on the same DataFrame.

Check your output
cache released
Show solutionHide solution
python
by_carrier.unpersist()
print("cache released")

5e

Cache by_carrier once more and time two count() calls with time.perf_counter. Print one line per call.

Hint 1

Time each call the way you would time any Python function : take time.perf_counter() before and after.

Check your output
first action       <seconds>s
second action      <seconds>s

The second count is faster than the first. The times themselves vary from run to run.

Show solutionHide solution
python
import time

by_carrier.cache()

for label in ("first action", "second action"):
    start = time.perf_counter()
    by_carrier.count()
    print(f"{label:16} {time.perf_counter() - start:6.2f}s")

Step 6Bigger job

6a

Run the cell with the join of flights and airports, the largest thing you run today. While it runs, watch the Jobs tab. How many stages did the whole thing take?

Check your output
+------------+---------+-------+
|municipality|avg_delay|flights|
+------------+---------+-------+
|         ...|      ...|    ...|
...
+------------+---------+-------+

Up to ten cities, avg_delay falling down the column, every flights above 10000. The cities and numbers depend on the dataset size.

Show solutionHide solution
python
worst = (flights
    .join(airports, flights.origin == airports.iata_code)
    .groupBy("municipality")
    .agg(F.avg("dep_delay").alias("avg_delay"),
         F.count("*").alias("flights"))
    .filter(F.col("flights") > 10000)
    .orderBy(F.desc("avg_delay")))

worst.show(10)

6b

Which stage took the longest, and what was it doing? Compare the Max and Median task time in that stage against the two numbers from 4a.

Hint 1

Open the Stages tab and sort on the Duration column.

Show solutionHide solution

Read the stage's name and its input and shuffle columns to see what it was doing. If its Max task time is far above its median, the work was uneven, as in Step 4.

6c

Open the SQL tab, click the query, and look at the plan diagram. Can you find the join in it? Does the plan say BroadcastHashJoin or SortMergeJoin? Is there a node in the plan you did not write? Print the plan in the notebook as well.

Hint 1

worst.explain() prints the same plan as text. Look for Exchange nodes : they mark the stage boundaries.

Check your output
== Physical Plan ==
AdaptiveSparkPlan ...
+- Sort ...
   +- Exchange rangepartitioning(avg_delay DESC ...)
      +- Filter (flights > 10000)
         +- HashAggregate(keys=[municipality...])
            +- Exchange hashpartitioning(municipality...)
               +- HashAggregate(keys=[municipality...])
                  +- BroadcastHashJoin ...

The shape is a scan and a join at the bottom, two aggregates with an Exchange between them, then the filter and the sort. Names and numbers vary.

Show solutionHide solution
python
worst.explain()

The airports table is small, so Spark may have chosen to broadcast it instead of shuffling both sides. Look for BroadcastHashJoin or SortMergeJoin, and for the Exchange nodes that mark the stage boundaries. On Fabric the SQL tab may also carry entries the platform itself created ; yours is the one whose time matches the job you just ran.

Step 7Stop session

7a

Click the session indicator and stop the session.

Show solutionHide solution

Nothing to run. A session holds its executors until the idle timeout, which is 60 minutes on the course pool, whether you are running anything or not. With the whole class on one capacity, a forgotten session is the reason somebody else cannot start theirs.

If time permits

  • Rerun Step 6 with spark.conf.set("spark.sql.adaptive.enabled", False) and compare the stage count and the plan against the adaptive run. What did AQE change? Set it back to True afterwards.
  • Repartition the flights to 4 partitions with flights.repartition(4) and rerun Step 2. What happens to the task count, and what happens to the time?
  • In the Stages tab, open the event timeline for the longest stage. Which colour dominates the bars : compute, shuffle read, or scheduler delay?

Tried it yourself first?

The solution is a spoiler. Work through the hints first : a wrong attempt teaches more than a solution you only read.