Python UDF Costs

Python UDFs and Their Serialization Cost

Executors are JVM processes. For a Python UDF each executor starts Python workers (pyspark.daemon) and streams the input to them over a local socket; the classic UDF pickles rows in small groups and calls your function once per row.

A classic Python UDF and the plan node it createsPython
from pyspark.sql.functions import udf
@udf("double", useArrow=False)            # the classic, pickled row-at-a-time UDF
def net_of_tax(total):
    return None if total is None else round(total / 1.06, 2)
amounts = orders.select(F.col("total").cast("double").alias("total"))
amounts.select(net_of_tax("total").alias("net")).explain()
Output
== Physical Plan ==
*(2) Project [pythonUDF0#37 AS net#36]
+- BatchEvalPython [net_of_tax(cast(total#9 as double))#35], [pythonUDF0#37]
   +- *(1) ColumnarToRow
      +- FileScan parquet [total#9] Batched: true, ..., ReadSchema: struct<total:decimal(10,2)>

BatchEvalPython is the boundary: code generation stops on either side because Catalyst cannot see into the function. Every value becomes a Python object, is pickled, unpickled and passed to a function call, and the result travels back the same way.