These ten questions cover the chapter, and in each the code runs, or fails, in a way that surprises people who have used Spark 129 for years: a filter that remembers an old value, a maximum that is not the largest number, a date that depends on where the session thinks it is, partition counts that ignore what you asked for, a hint that beats a setting, a write that destroys its own input, a UDF that returns nothing, and an overwrite that deletes more than it replaces. Every question ran on the book's workstation with Spark 4.2.0 in local mode, four cores, Python 3.14 and the session time zone set to UTC, against BookNest's sample order history at ten times the size of JSON, Columnar and Binary Formats's files (Spark at BookNest and Loading Order History). Write down what each line prints and why, then check Appendix I, which gives the real output, the reason, the fix and the section each question comes from.
Running the questions
The script reads data/orders.parquet, data/order_lines.parquet and data/books.parquet, the tables Loading Order History wrote, and writes its own scratch files under out/quiz/. Run it from the folder that holds data/ with python questions.py in the book's virtual environment (Python Virtual Environment); it takes under a minute. Every result line starts with its question's number, and a failure prints the exception's class and, where Spark gives one, its error condition, so you can compare your notes line by line.
import contextlib, io, os
from pyspark.sql import SparkSession, Window, functions as F
from pyspark.sql.functions import udf
spark = (SparkSession.builder.master("local[4]").appName("quiz")
.config("spark.sql.session.timeZone", "UTC").config("spark.ui.enabled", "false")
.getOrCreate())
spark.sparkContext.setLogLevel("OFF")
orders = spark.read.parquet("data/orders.parquet")
lines = spark.read.parquet("data/order_lines.parquet")
books = spark.read.parquet("data/books.parquet")
def q(n, fn): # print a result, or the error's condition
try:
print(f"{n}:", fn())
except Exception as e:
print(f"{n}:", type(e).__name__, getattr(e, "getCondition", lambda: "")())
# 1. The threshold changes after the filter is defined. Which one does count() use?
# (Section 5.5.7)
threshold = 50
big = orders.where(F.col("total") > threshold)
threshold = 100
q(1, big.count)
# 2. Order totals go to CSV and come back without a schema. What is the "largest" total?
# (Sections 5.5.2 and 5.5.3)
orders.select("order_id", "total").write.mode("overwrite").csv("out/quiz/csv", header=True)
q(2, lambda: spark.read.csv("out/quiz/csv", header=True).agg(F.max("total")).first()[0])
# 3. The last order date, in two session time zones. (Section 5.5.1)
def last_day(zone):
spark.conf.set("spark.sql.session.timeZone", zone)
return str(orders.agg(F.max(F.to_date("order_ts"))).first()[0])
q(3, lambda: (last_day("UTC"), last_day("Asia/Kuala_Lumpur")))
spark.conf.set("spark.sql.session.timeZone", "UTC")
# 4. Partitions before, after coalesce(64) and after repartition(64). (Section 5.8.2)
q(4, lambda: [d.rdd.getNumPartitions() for d in (lines, lines.coalesce(64),
lines.repartition(64))])
# 5. Partitions after a groupBy on channel, with AQE on and off. (Sections 5.3.10, 5.10.6)
def by_channel(aqe):
spark.conf.set("spark.sql.adaptive.enabled", aqe)
return orders.groupBy("channel").count().rdd.getNumPartitions()
q(5, lambda: (by_channel(True), by_channel(False)))
spark.conf.set("spark.sql.adaptive.enabled", True)
# 6. Automatic broadcasts are off. Which join does the plan use? (Section 5.6.5)
spark.conf.set("spark.sql.autoBroadcastJoinThreshold", -1)
plan = io.StringIO()
with contextlib.redirect_stdout(plan):
lines.join(F.broadcast(books), lines.book_id == books.id).explain()
q(6, lambda: [w for w in plan.getvalue().split() if w.endswith("Join")][0])
spark.conf.unset("spark.sql.autoBroadcastJoinThreshold")
# 7. Drop cancelled lines from a daily summary by overwriting it in place. (Section 5.5.5)
(lines.groupBy(F.to_date("order_ts").alias("day"), "status").count()
.write.mode("overwrite").parquet("out/quiz/daily"))
daily = spark.read.parquet("out/quiz/daily")
q(7, lambda: daily.where("status <> 'cancelled'").write.mode("overwrite")
.parquet("out/quiz/daily"))
q(7, lambda: os.listdir("out/quiz/daily"))
# 8. A UDF declared DOUBLE returns a Python int: the sum of the title lengths, pickled
# and with Arrow. (Sections 5.7.1 and 5.7.3; localCheckpoint() as in 5.7.3's warning)
titles = books.select("title").localCheckpoint()
q(8, lambda: [titles.select(F.sum(udf(len, "double", useArrow=a)("title"))).first()[0]
for a in (False, True)])
# 9. Number every order by time. How many partitions does the result have? (Section 5.6.8)
q(9, lambda: orders.withColumn("n", F.row_number().over(Window.orderBy("order_ts")))
.rdd.getNumPartitions())
# 10. Monthly copies, then June 2026 rebuilt with a plain overwrite. How many month
# directories remain? (Section 5.8.3)
monthly = (lines.groupBy(F.date_format("order_ts", "yyyy-MM").alias("month"))
.agg(F.sum("qty").alias("copies")))
monthly.write.mode("overwrite").partitionBy("month").parquet("out/quiz/monthly")
(monthly.where("month = '2026-06'")
.write.mode("overwrite").partitionBy("month").parquet("out/quiz/monthly"))
q(10, lambda: len([d for d in os.listdir("out/quiz/monthly") if d.startswith("month=")]))When an answer surprises you, ask Spark rather than guessing: explain() shows the join and the shuffle, and rdd.getNumPartitions() shows the layout.