Pair your devices with a code and playback position follows you: pause on this device, hit resume on the other. Position is saved to the site every minute and on pause.
Open this panel on your other device and enter the same code.
Starts this lesson and continues through 8 more to the end of the course.
df = spark.read.parquet("s3://bucket/data/") # line 3
df = df.filter(col("region") == "EU") # line 4
df = df.withColumn("x", col("a") / col("b")) # line 5
df.write.parquet("s3://bucket/out/") # line 6 ← the stack trace points here
The path in line 3 doesn't exist. The traceback blames line 6.
That's lazy evaluation, and it's the first thing to internalise about debugging PySpark: the line in the traceback is where execution was triggered, not where the mistake is. Everything before it was plan construction, which cannot fail on data it has never read.
| Transformations | Actions | |
|---|---|---|
| Examples | select, filter, withColumn, join, groupBy().agg() |
write, count, collect, show, take, toPandas |
| What they do | Append to the plan | Execute the plan |
| Return | A new DataFrame | A value, or a side effect |
| Cost when the line runs | ~nothing | everything |
Transformations are free and actions are not. That's the whole model, and it explains the debugging experience: a script that "runs in two seconds" and then spends forty minutes on the final line was never running at all.
Transformations divide into two kinds, and the difference decides your job's cost:
filter, select,
withColumn. No data movement.groupBy, join, distinct,
orderBy, repartition. Requires a shuffle: data crosses the network, is written to disk, and
re-read.Every wide transformation is a stage boundary. Counting the wide transformations in your script gives you the number of stages, and therefore the number of times your data gets written and re-read.
⚠️ Shuffle is the dominant cost in most jobs. Engine-level mechanics — how the shuffle is written,
partitioned, and read, and how the optimiser picks join strategies — are covered in module E7 of
the sibling db-training repo. What matters here is the PySpark habit: scan your script for wide
transformations, and question each one.
This is the highest-cost lazy-evaluation mistake, and it's invisible on the page:
expensive = spark.read.parquet(big).filter(...).join(other, ...) # plan only
count = expensive.count() # executes the whole thing
sample = expensive.limit(10).collect() # executes the whole thing AGAIN
expensive.write.parquet(out) # and AGAIN
A DataFrame is a plan, not a result. Every action re-executes it from the source. Three actions, three full executions of an expensive join.
expensive = spark.read.parquet(big).filter(...).join(other, ...).cache()
count = expensive.count() # executes and populates the cache
sample = expensive.limit(10).collect() # reads the cache
expensive.write.parquet(out) # reads the cache
expensive.unpersist() # release it when done
Rules for cache() / persist():
cache() is itself lazy. It marks the DataFrame; the data lands on the first action.unpersist() when you're done. Cached data occupies executor memory that your shuffle wants.⚠️ The common anti-pattern is caching everything "to be safe". That evicts what actually needed to be cached and starves execution memory.
count() is not freeif df.count() > 0: # full execution
df.write.parquet(out) # full execution AGAIN
count() executes the plan. Using it as an emptiness check doubles your work. If you only need to
know whether anything is there:
if df.take(1): # short-circuits — reads one partition, usually
df.write.parquet(out)
⚠️ Even take(1) can be expensive if a filter is highly selective — Spark may scan many partitions
before finding a row. The genuinely cheap answer is usually to write unconditionally and let the
output be empty, or to check a metric you already computed.
Two methods, and the second is the one to use:
df.explain() # physical plan
df.explain(True) # parsed → analysed → optimised → physical
df.explain("formatted") # more readable in recent versions
What to look for, in rough priority:
Exchange — a shuffle. Count them. Each one is a stage boundary.PushedFilters / partition pruning — did your filter reach the scan, or is it happening after
reading everything?Point 3 is the one that most often produces a large, cheap win: a filter that could have pruned partitions but didn't, because it was expressed in a way the optimiser couldn't push down (a UDF in the predicate is the classic cause — see lesson 3).
⚠️ I'm deliberately not enumerating physical-operator names here; they vary by Spark version and they're covered properly in db-training E7. Learn the four questions above and read your plan.
write line but the bug is a bad input path. Why?filter, groupBy, select, join, withColumn cause a shuffle?count(), then collect(), then write. What have you paid?if df.count() > 0: a bad emptiness check?explain() output?groupBy and join — both wide, both shuffle. filter, select and withColumn are narrow and
move no data.count() executes the entire plan, and then write executes it again — you've doubled the work
to answer a yes/no question. take(1) short-circuits, though it can still be costly with a highly
selective filter.Exchange nodes — the shuffles. Then the join strategy, then whether filters were pushed down to
the scan..cache() and run it again. The number is more persuasive than the explanation.count() for something other than a count?" Emptiness
checks and logging are the usual culprits, and both are usually removable.