SpiralTrain
Exercises › Block 1 · Exercise 11

Block 1 · Exercise 11

PySpark Intro

Starter notebook
11-pyspark-intro-starter
Fabric path
/lakehouse/default/Files/data/solutions/11.pyspark-intro/

Open the notebook 11-pyspark-intro-starter in your workspace and run its first cell. It takes the session Fabric started and puts the SparkContext in sc, and it can take a few minutes. Each step has a cell with a comment where your code goes.

Do one part at a time, 1a, 1b and so on, and compare with the expected output before you go on. The solution is in 11-pyspark-intro-solution, under every part on the exercise site, and at the back of the exercises PDF.

Step 1Session

1a

The variable spark is the session. Print it, spark.version, spark.sparkContext.appName and spark.sparkContext.master, each on its own line, with a label in front of the last three.

Hint 1

The slide Developing in PySpark says why spark exists before your first cell runs.

Check your output
<pyspark.sql.session.SparkSession object at 0x...>
Spark version : 3.5.x
application   : <name of the application>
master        : <master>

The address, the application name and the master vary with the environment. Read the master : does it start with local?

Show solutionHide solution
python
print(spark)
print("Spark version :", spark.version)
print("application   :", spark.sparkContext.appName)
print("master        :", spark.sparkContext.master)

1b

Make a DataFrame from the list data with the column names Name and Age, and show it.

Hint 1

The session has a method that turns a list of tuples and a list of column names into a DataFrame.

Check your output
+-------+---+
|   Name|Age|
+-------+---+
|  Alice| 25|
|    Bob| 30|
|Charlie| 35|
+-------+---+
Show solutionHide solution
python
df = spark.createDataFrame(data, ["Name", "Age"])
df.show()

Step 2Context

2a

Take the SparkContext out of the session with spark.sparkContext and call it sc. Print sc, then its version and master on one line.

Check your output
<SparkContext master=<master> appName=<name of the application>>
3.5.x <master>
Show solutionHide solution
python
sc = spark.sparkContext
print(sc)
print(sc.version, sc.master)

2b

Make an RDD from the numbers 1 to 5 with sc.parallelize. Print its count and its sum, with the labels count : and sum :.

Hint 1

An RDD has the methods count() and sum().

Check your output
count : 5
sum   : 15
Show solutionHide solution
python
numbers = sc.parallelize([1, 2, 3, 4, 5])
print("count :", numbers.count())
print("sum   :", numbers.sum())

2c

Ask the builder for a second session : SparkSession.builder.appName("A Second Session").getOrCreate(). Print whether it is the same object as spark, with is.

Hint 1

getOrCreate gets the session that exists, or creates one when there is none.

Check your output
True
Show solutionHide solution
python
other = SparkSession.builder.appName("A Second Session").getOrCreate()
print(other is spark)

Step 3Configuration

3a

Read the setting spark.sql.shuffle.partitions with spark.conf.get and print it.

Check your output
200

The value depends on the environment. It is a number, given as text.

Show solutionHide solution
python
print(spark.conf.get("spark.sql.shuffle.partitions"))

3b

Change spark.sql.shuffle.partitions to 4 with spark.conf.set. Read it again and print it.

Check your output
4
Show solutionHide solution
python
spark.conf.set("spark.sql.shuffle.partitions", 4)
print(spark.conf.get("spark.sql.shuffle.partitions"))

3c

Set spark.executor.memory to "2g" inside the try. Run the cell and read what the except prints.

Hint 1

The memory of an executor is decided when the executor starts. Ask whether a running session can still change that.

Check your output
AnalysisException : [CANNOT_MODIFY_CONFIG] Cannot modify the value of the Spark config: "spark.executor.memory".

Some lines with a link may follow.

Show solutionHide solution
python
try:
    spark.conf.set("spark.executor.memory", "2g")
except Exception as error:
    print(type(error).__name__, ":", error)

Step 4Spark UI

4a

Add up the numbers 0 to 999,999 with spark.range and selectExpr("sum(id)"), and show the result.

Hint 1

spark.range(n) makes a DataFrame with a column id. selectExpr takes SQL text.

Check your output
+------------+
|     sum(id)|
+------------+
|499999500000|
+------------+
Show solutionHide solution
python
spark.range(1_000_000).selectExpr("sum(id)").show()

4b

Open the Spark jobs panel under the cell of 4a. Find the job that computed the sum, and follow it into the Spark UI.

Hint 1

The slide Spark Web UI says where the panel is.

4c

In the Spark UI, open the Environment tab and search for spark.sql.shuffle.partitions. Is the value the one you set in Step 3? Print the value in force now, and write a comment that says what the tab shows.

Hint 1

The slide Spark Web UI lists configuration as one of the things the UI shows. Ask when that configuration was read.

Check your output
4
Show solutionHide solution
python
print(spark.conf.get("spark.sql.shuffle.partitions"))
# No : the Environment tab lists the configuration the application started with. A key set at startup
# shows its startup value, a key that was not set does not appear, and spark.conf.set in a cell
# changes neither. spark.conf.get reads the value in force now.

4d

Open the Executors tab. How many executors does your session have, and how many cores does each have? Write the numbers down for the performance exercise. Then print the same numbers from code.

Hint 1

The Spark UI counts the driver in its own row. The number of executors is the number of rows without it.

Hint 2

spark.sparkContext.defaultParallelism gives the cores available to tasks. The executors are behind spark.sparkContext._jsc.sc().getExecutorMemoryStatus(), and the driver counts as one entry there.

Check your output
executor entries (the driver counts as one) : <number>
hosts and memory                            : Map(<host:port> -> (<maximum memory>,<free memory>), ...)
cores available to tasks                    : <number>

The numbers depend on the session. Compare them with the Executors tab.

Show solutionHide solution
python
status = spark.sparkContext._jsc.sc().getExecutorMemoryStatus()
print("executor entries (the driver counts as one) :", status.size())
print("hosts and memory                            :", status)
print("cores available to tasks                    :", spark.sparkContext.defaultParallelism)

If time permits

  • Run spark.sparkContext.getConf().getAll() and read the list. How many settings are there, and how many of them mention Fabric or Synapse?
  • Run the job of 4a again and watch the Spark jobs panel. Is a second job created, or is the first result reused?
  • Open Monitor in the Fabric sidebar and find your own session in the list of running applications.
  • Stop your session from the session menu, then run any cell. Time how long the new session takes to come up.
  • On a laptop a program builds its own session with SparkSession.builder.appName("My Application").master("local[*]").getOrCreate(). Run that on Fabric, and check whether the result is spark and whether the master changed.
  • Values such as the memory per executor are set before the session starts. On Fabric a %%configure cell as the first cell of the notebook does it, and running it restarts the session. Do not run it in this notebook. Read the shape :
python
%%configure
{
    "conf": {
        "spark.sql.shuffle.partitions": "8"
    }
}

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.