Transformations vs actions
Why nothing runs until you call an action, and what lazy evaluation buys Spark.
On this page
You will learn
- The difference between a transformation and an action
- What lazy evaluation is and why Spark uses it
- Why the same DataFrame can be computed twice, and how to avoid it
- How laziness affects where errors appear
select, filter, join, ...) only describe a new DataFrame. Actions (count, collect, write, show) make Spark actually run. Waiting until an action lets Spark optimise the whole chain at once.Two kinds of operations
| Transformations (lazy) | Actions (run the job) |
|---|---|
select, withColumn, filter | count(), collect(), take(n) |
groupBy().agg(), join, orderBy | show(), toPandas() |
union, distinct, repartition | write...save(), saveAsTable |
cache() (marks only) | foreach, foreachPartition |
A transformation returns a new DataFrame immediately, in microseconds, whatever the data size, because it only adds a step to a logical plan. An action returns a value or writes output, so Spark has to compute.
Walk through a pipeline
Line 1 · read
The command
orders = spark.read.parquet("s3://shop/orders")
Line 2 · filter
The command
paid = orders.filter(F.col("status") == "PAID")
Line 3 · aggregate
The command
daily = paid.groupBy("order_date").agg(F.sum("amount").alias("revenue"))
Line 4 · action
The command
daily.write.mode("overwrite").parquet("s3://shop/daily_revenue")
Why lazy?
Seeing the whole pipeline before running it is what makes Catalyst's optimisations possible: filters move down next to the scan, unused columns are never read, consecutive projections merge into one, and joins can pick a strategy knowing what comes after. Running each line eagerly, as pandas does, would read everything first.
The recomputation trap
A DataFrame is a recipe, not a result. Every action re-runs the recipe from the source:
clean = raw.filter(...).withColumn(...) # expensive clean.count() # job 1: reads and cleans clean.write.parquet("out") # job 2: reads and cleans again
If a DataFrame is used by several actions, either cache it (see the Caching lesson) or restructure so a single action does the work. A common unnecessary action is a count() just to log a number.
Where errors show up
- Analysis errors (a missing column, a type mismatch) appear immediately when you write the transformation, because Spark resolves names eagerly.
- Runtime errors (a bad cast in ANSI mode, a division by zero, a corrupt file, an out-of-memory) appear only at the action, often many lines later. The stack trace points at the action, not at the transformation that caused it.
limit(10).show() after a suspicious step to force evaluation there. Remove it afterwards.Common mistakes
Timing transformations
Calling several actions on an uncached DataFrame
Debug count() calls left in production
Key takeaways
- Transformations build a plan; actions run it.
- Laziness lets Catalyst optimise the whole pipeline.
- Each action recomputes from the source unless you cache.
- Runtime errors surface at the action, not where they were caused.
Check yourself
3 questions1. Which is an action?
Show the answer
count(). count returns a value to the driver, so Spark must compute.
2. A DataFrame is written and then counted, with no cache. How many times is the source read?
Show the answer
Twice. Each action re-runs the plan from the source.
3. When does a missing column error appear?
Show the answer
Immediately, when the transformation is defined. Spark analyses names eagerly, so unresolved columns fail right away.
Go deeper
Jobs, stages and tasksSpark internals
Caching and persistenceSpark internals
Catalyst and physical plans
Primary sources: RDD programming guide: lazy evaluation · Spark SQL guide