PySpark courseLesson 3 of 6
PySpark course · Lesson 3 of 6
PySpark Transformations vs Actions
Understand lazy evaluation in PySpark: transformations build a plan, actions run it, and narrow versus wide transformations decide where Spark shuffles data.
On this page
PySpark code looks like it runs line by line. It does not. Transformations only describe a new DataFrame; actions make Spark plan and execute the whole chain.
Laziness in practice
from pyspark.sql import SparkSession, functions as F
spark = SparkSession.builder.master("local[2]").appName("example").getOrCreate()
events = spark.range(0, 1_000_000).withColumn("bucket", F.col("id") % 10) # transformation
big = events.filter(F.col("id") > 500_000) # transformation
counts = big.groupBy("bucket").count() # transformation
print(type(counts).__name__) # still just a plan
print(counts.orderBy("bucket").first()) # action: now Spark runs a job
DataFrame
Row(bucket=0, count=49999)
Until first(), nothing was computed. Spark could see the whole chain and optimise it: the filter is applied before the grouping, and only the needed columns are read.
Common transformations and actions
| Transformations (lazy) | Actions (run a job) |
|---|---|
select, withColumn, filter |
show, count, collect, first, take |
join, groupBy().agg(), distinct |
write....save(), saveAsTable, foreach |
orderBy, repartition, union |
toPandas, toLocalIterator |
Narrow versus wide
- Narrow: each output partition depends on one input partition (
filter,select,withColumn). No data moves. - Wide: output partitions need data from many input partitions (
groupBy,join,distinct,orderBy). Spark must shuffle, and that creates a stage boundary.
You can see the shuffle in the plan as an Exchange:
counts.explain()
The plan shows a partial aggregation, then Exchange hashpartitioning(bucket, ...), then the final aggregation. Spark pre-aggregates on each partition before shuffling, which is why groupBy().count() moves far less data than the raw rows.
Why repeated actions can be slow
Each action re-executes the plan from the source unless results are cached:
counts.count()
counts.show(3)
Both lines recompute events → filter → groupBy. If an expensive DataFrame is reused by several actions, cache() it (and unpersist() when done). Do not cache everything: memory is finite and caching cheap steps adds overhead.
Common mistakes
- Timing a transformation line and concluding it is fast (nothing ran).
- Calling
count()“just to check” between every step in production code, triggering a full job each time. - Caching without reuse.
collect()on large data.
Interview relevance
This is one of the most common Spark interview questions. Combine it with narrow versus wide and where shuffles occur. See the interview answer.
Key takeaway
Transformations build a plan; actions run it. Know which of your transformations are wide, because those are where Spark shuffles.
Progress is saved in this browser only. No account needed.