AWS Training
Modules Listen Certification
0:00 0:00

← PySpark

Starts this lesson and continues through 8 more to the end of the course.

Lazy evaluation, transformations, and actions

Why your job failed on the wrong line

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 vs actions

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.

The narrow/wide distinction

Transformations divide into two kinds, and the difference decides your job's cost:

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.

The re-computation trap

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():

⚠️ 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 free

if 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.

Reading the plan

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:

  1. Exchange — a shuffle. Count them. Each one is a stage boundary.
  2. Join strategy — a broadcast versus a shuffle-based join is often the entire performance difference on a small-to-large join.
  3. PushedFilters / partition pruning — did your filter reach the scan, or is it happening after reading everything?
  4. Scan sizes — is it reading columns and partitions you don't need?

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.

Check yourself

  1. Your traceback blames the write line but the bug is a bad input path. Why?
  2. Which of filter, groupBy, select, join, withColumn cause a shuffle?
  3. A DataFrame is used by count(), then collect(), then write. What have you paid?
  4. Why is if df.count() > 0: a bad emptiness check?
  5. What's the first thing you look for in explain() output?
Answers
  1. Lazy evaluation. Everything before the action was plan construction and couldn't fail; the action triggered execution, so that's where the error surfaced. Debug backwards from the action.
  2. groupBy and join — both wide, both shuffle. filter, select and withColumn are narrow and move no data.
  3. Three full executions of the whole plan, from the source. A DataFrame is a plan, not a result. Cache it if it's consumed more than once, and unpersist afterwards.
  4. 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.
  5. Exchange nodes — the shuffles. Then the join strategy, then whether filters were pushed down to the scan.

Teaching this section

← PreviousThe Python/JVM boundaryNext →UDFs, Arrow, and pandas UDFs