Because nothing runs until an action, Catalyst sees the whole pipeline before reading a byte and can rewrite it. The price is that runtime errors surface far from the line that caused them.
orders = spark.read.parquet("data/orders.parquet")
wide = orders.withColumn("year", F.year("order_ts")) # all columns, all rows...
web = wide.where("channel = 'web'").select("order_id", "total") # ...until you narrow it
web.explain()
bad = orders.withColumn("n", F.col("channel").cast("int")) # 'web' is no int: no error yet
print("count:", bad.count()) # n is never computed
try:
bad.agg(F.sum("n")).collect() # now it is
except Exception as e:
print(type(e).__name__, e.getCondition())== Physical Plan ==
*(1) Project [order_id#0L, total#9]
+- *(1) Filter (isnotnull(channel#3) AND (channel#3 = web))
+- *(1) ColumnarToRow
+- FileScan parquet [order_id#0L,channel#3,total#9] Batched: true, DataFilters: [...],
Format: Parquet, ..., PushedFilters: [IsNotNull(channel), EqualTo(channel,web)],
ReadSchema: struct<order_id:bigint,channel:string,total:decimal(10,2)>
count: 1000000
NumberFormatException CAST_INVALID_INPUTThe code computed year for every row; the plan dropped it, reads three of ten columns (ReadSchema) and pushes the channel filter into the Parquet 129 reader, which skips row groups by their statistics (Statistics and Pushdown). The cast shows the other side: under ANSI mode (ANSI SQL Mode by Default) 'web' cannot become an INT, yet count() succeeded because the optimizer pruned the unused column; the error waited for an action that needed the value. Classic PySpark 129 still analyzes each step eagerly, so a misspelled column fails on its own line, while a Spark Connect 129 client (Spark Connect) defers even that.