Lesson 12 / 25

Transformations, Actions and Lazy Evaluation

Distinguish transformations from actions and see why nothing runs until an action.

Nothing happens until you ask for a result

Transformations (select, filter, withColumn, join, groupBy) describe a computation and return instantly; they only add to a logical plan. Actions (count, show, collect, write, first, take) force Spark to run the plan and return or save a result. This lazy evaluation lets Spark see your whole pipeline and optimise it as one, for example pushing a filter before a join or reading only needed columns. The cost is that errors appear late (at the action), and calling an action twice recomputes everything unless you cache.

Plan first, run later

Transformations build a plan; an action optimises it, cuts it into stages at shuffles and runs the tasks.

Four ideas: lazy, shuffle, plan, partitions.
Figure 4.1 — Lazy, shuffle, plan and partitions.

A million rows, one job, run

I ran this on Apache Spark 4.0.0 (PySpark, local mode, in the official Docker image). The filter and withColumn lines return immediately without doing work; only count() runs a job. Half of the numbers 0 to 999,999 are even.

big = spark.range(0, 1000000)
t = big.filter("id % 2 = 0").withColumn("sq", F.col("id") * 2)   # nothing runs yet
print("no job yet; count =", t.count())          # the action triggers the job

Output:

no job yet; count = 500000

Use take() or limit() while developing

Test logic on a small sample (df.limit(1000)) to get fast feedback, then run on the full data.

Quick check: Which of these is an action?

  • filter()
  • count()
  • select()
  • withColumn()
Answer

count() — count() forces execution; the others only extend the plan.