Almost every line of Spark you write calls a function from org.apache.spark.sql.functions or pyspark.sql.functions, or uses the same function by name in SQL. There are several hundred of them, and most bugs in a Spark pipeline involve one of them. It returned NULL when you expected a value, raised an error after an upgrade, parsed a date in the wrong time zone, or turned out slower than a hand-written UDF.
This page explains what a built-in function is to Spark, how to find the one you need, and the rules that decide its result: NULL propagation, ANSI mode, the try_ family, time zones and the determinism of aggregates. It then shows how to prove from the plan that a built-in beats a UDF, using a worked cleaning job. Arrays, maps and higher-order functions such as transform are covered in Spark DataFrame Complex Types, so they are only mentioned here.
What a function is to Spark
A call such as F.lower(F.col("email")) computes nothing. It builds a node in an expression tree and wraps it in a Column. The same holds for F.expr("lower(email)") and for the SQL text SELECT lower(email). All three produce an unresolved function call that the analyzer looks up by name in the session's function registry. Built-ins, temporary functions you register and catalog functions are all found this way. The result is a Catalyst Expression that knows its input types, its output type, whether it can return NULL, and whether it is deterministic.
That knowledge is what makes built-ins fast. The optimizer folds constant subexpressions, removes redundant casts and simplifies CASE branches. It can also push a filter like lower(country) = 'de' towards the scan when the expression allows it. Whole-stage code generation then compiles a chain of built-ins into one tight JVM loop with no per-row virtual calls, as described in Whole-Stage Codegen. A user-defined function stops all of this, because Spark knows only its declared return type.
Two practical consequences follow. First, a function that exists in SQL but has no Python or Scala wrapper in your version can still be called with F.expr, or with F.call_function("name", col), which was added in Spark 3.5. Second, a name in a string is resolved at analysis time, so a typo fails when the plan is built, before any task runs. That makes analysis errors cheap to catch in unit tests.
Finding the function you need
Do not guess at function names or argument orders. Ask the session, which reports exactly what your version supports:
SHOW FUNCTIONS LIKE '*date*';
DESCRIBE FUNCTION EXTENDED date_trunc; -- usage, arguments, examples, "Since" version
SELECT current_date(), version();The output of DESCRIBE FUNCTION EXTENDED includes a usage line, examples and the version that introduced the function. Check it before you rely on something you read in a blog, including this one. The functions fall into five families, and the family decides where a function may appear:
| Family | Examples | Where it is legal |
|---|---|---|
| Scalar | lower, regexp_extract, to_date, coalesce | Anywhere a column is accepted |
| Aggregate | sum, count_if, max_by, approx_count_distinct | groupBy().agg(), or over a window |
| Window-only | row_number, rank, lag, lead | Only with over(Window...) |
| Generator | explode, posexplode, inline | One per select; changes the row count |
| Table-valued | range, explode in FROM | FROM clause in SQL |
NULL semantics
Spark SQL follows three-valued logic. Almost every scalar function returns NULL if any input is NULL, and a comparison with NULL is NULL rather than false. A filter keeps a row only when its predicate is true, so WHERE status != 'test' silently drops rows whose status is NULL. Most surprising results come from a handful of rules:
count(col)skips NULLs andcount(*)does not. A ratio built from the two mixes populations.concatreturns NULL if any argument is NULL, whileconcat_wsskips NULL arguments.when(cond, v)withoutotherwiseyields NULL for unmatched rows.sumandavgignore NULLs, and both return NULL when every input is NULL. A total of zero and a total of nothing look the same unless you wrap it incoalesce(sum(x), 0)on purpose.a = bis NULL when either side is NULL. Usea <=> bin SQL oreqNullSafein the DataFrame API when NULL should equal NULL, for example in change detection or join keys.greatestandleastskip NULLs, unlike arithmetic, sogreatest(a, b)anda + bdisagree about a NULL input.
from pyspark.sql import functions as F
df = spark.createDataFrame([(1, None), (2, 5), (None, None)], "a int, b int")
df.select(
(F.col("a") + F.col("b")).alias("plus"), # NULL, 7, NULL
F.greatest("a", "b").alias("greatest"), # 1, 5, NULL
F.col("a").eqNullSafe(F.col("b")).alias("eq"), # false, false, true
F.concat_ws("-", "a", "b").alias("ws"), # "1", "2-5", ""
).show()
ANSI mode and the try_ functions
From Spark 4.0, spark.sql.ansi.enabled defaults to true. The setting has existed since 3.0, but it used to default to false, which gave Hive-compatible behaviour. The difference matters for functions, because ANSI mode turns silent corruption into errors:
| Expression | ANSI off | ANSI on (4.0+ default) |
|---|---|---|
2147483647 + 1 on int | wraps to -2147483648 | ARITHMETIC_OVERFLOW error |
CAST('a' AS INT) | NULL | CAST_INVALID_INPUT error |
10 / 0 | NULL | division-by-zero error |
element_at(arr, 99) | NULL | index-out-of-bounds error |
Do not respond to these errors by turning ANSI mode off for the whole job. The errors are telling you that some rows are bad. Decide row by row which failures are expected and use the matching try_ function there. Each one behaves like its base function but returns NULL instead of raising. The ANSI documentation lists, among others, try_cast, try_add, try_divide, try_sum, try_avg, try_element_at, try_to_timestamp and try_to_date. Several were added in different releases, so confirm with DESCRIBE FUNCTION on your cluster.
raw = spark.read.json("s3://bucket/events/") # amount_raw and ts_raw are strings
parsed = raw.select(
"event_id", "amount_raw",
F.expr("try_cast(amount_raw AS DECIMAL(12,2))").alias("amount"),
F.expr("try_to_timestamp(ts_raw, 'yyyy-MM-dd HH:mm:ss')").alias("ts"),
)
# A NULL that was not NULL in the input is a parse failure, not missing data.
bad = parsed.filter(F.col("amount").isNull() & F.col("amount_raw").isNotNull())
print("unparseable amounts:", bad.count()) # quarantine these; do not let them vanishThe pattern matters more than the syntax. A try_ call turns a crash into a NULL, and an unexamined NULL is the same silent corruption ANSI mode was meant to stop. Always pair it with a count of the rows it nulled out.
Dates, timestamps and time zones
Date and time functions cause more production incidents than any other family, for three reasons.
The session time zone. A TIMESTAMP is an instant. Functions that render or truncate it, such as date_format, hour, date_trunc and to_date, use spark.sql.session.timeZone, which defaults to the JVM's zone. Two clusters in different regions can therefore write different daily partitions from the same data. Set the zone explicitly, usually to UTC, in the job configuration rather than relying on the machine.
Instant or wall-clock time. Spark 3.4 and later also have TIMESTAMP_NTZ, a local date-time with no zone. Use it for values that are wall-clock by nature, such as a store's opening hour. Use TIMESTAMP for events that happened at an instant. Mixing them through implicit casts shifts values by the zone offset.
Pattern letters. Patterns are case-sensitive: MM is month and mm is minute; yyyy is the calendar year. Spark 3.0 replaced the old SimpleDateFormat parser with a stricter one, and inputs the old parser accepted can now fail or, without ANSI, become NULL. Test every pattern against real sample strings, including single-digit days and the last day of the year.
spark.conf.set("spark.sql.session.timeZone", "UTC")
daily = events.groupBy(F.date_trunc("day", "ts").alias("day")).agg(F.count("*").alias("n"))
# from_utc_timestamp / to_utc_timestamp shift instants for display in a named zone.
local = events.withColumn("ts_berlin", F.from_utc_timestamp("ts", "Europe/Berlin"))
Aggregates: determinism and accuracy
Aggregates have their own traps, mostly about determinism and accuracy.
collect_listandcollect_setreturn elements in whatever order the shuffle delivered them. That order can change between runs. Sort inside the array witharray_sort, or aggregate a struct whose first field is the sort key.firstandlastin agroupByare not deterministic for the same reason. To get the value ofxat the latestts, usemax_by(x, ts), which says what you mean and is a single aggregate.approx_count_distinct(col, rsd)uses HyperLogLog++. Its default relative standard deviation is 0.05, so roughly two in three results fall within 5 percent of the true count. Pass a smallerrsdfor tighter answers at more memory per group.percentile_approx(col, p, accuracy)defaults to an accuracy of 10000, and its relative error is about 1/accuracy. Raise it for tail percentiles on large groups, and expect more memory per group.count_if(cond)andsum(when(cond, 1).otherwise(0))give the same answer. The first is easier to read.
Built-ins versus UDFs, proven from the plan
The general rule, prefer built-ins to UDFs, is right, but you should be able to prove it for your own query. Here is the same task written both ways: extract a lower-cased domain from an email column and keep only one domain.
from pyspark.sql.types import StringType
@F.udf(StringType())
def domain_udf(email):
if email is None or "@" not in email:
return None
return email.split("@", 1)[1].lower()
slow = users.filter(domain_udf("email") == "example.com")
fast = users.filter(F.lower(F.substring_index("email", "@", -1)) == "example.com")
slow.explain() # ... BatchEvalPython [domain_udf(email)] ... Filter after the Python step
fast.explain() # ... Filter (lower(substring_index(email, @, -1)) = example.com) inside codegenIn the first plan, every row is serialized, sent to a Python worker, evaluated, and sent back before the filter runs. The code generator also has to stop at the UDF boundary. In the second plan the filter is fused with the scan in one generated loop, and nothing leaves the JVM. One edge case behaves differently: substring_index returns the whole string when there is no @, where the UDF returned NULL. Equivalence tests catch differences like this, so write one before you swap implementations.
A UDF is the right tool when no combination of built-ins expresses the logic, for example calling a library or parsing a proprietary format. Prefer a vectorized pandas UDF, which moves Arrow batches instead of single rows. Its batch-boundary rules are in Spark Pandas UDFs. Mark a UDF asNondeterministic() if it is not deterministic, so the optimizer does not duplicate or reorder it. To read the plans themselves, see Spark Explain Plans.
Worked example: cleaning an orders feed
Here is a realistic cleaning step for a raw orders feed, using only built-ins. It shows the habits from the earlier sections in one place.
from pyspark.sql import functions as F, Window
spark.conf.set("spark.sql.session.timeZone", "UTC")
raw = spark.read.json("s3://bucket/orders/raw/")
clean = (raw
.withColumn("order_ts", F.expr("try_to_timestamp(order_ts_raw, 'yyyy-MM-dd HH:mm:ss')"))
.withColumn("amount", F.expr("try_cast(amount_raw AS DECIMAL(12,2))"))
.withColumn("country", F.upper(F.trim("country")))
.withColumn("email_domain", F.lower(F.substring_index("email", "@", -1)))
.withColumn("bad_row", F.col("order_ts").isNull() | F.col("amount").isNull()))
quality = clean.agg(F.count("*").alias("rows"), F.count_if("bad_row").alias("bad"))
w = Window.partitionBy("order_id").orderBy(F.col("order_ts").desc())
latest = (clean.filter(~F.col("bad_row"))
.withColumn("rn", F.row_number().over(w))
.filter("rn = 1").drop("rn"))
daily = latest.groupBy(F.to_date("order_ts").alias("day"), "country").agg(
F.sum("amount").alias("revenue"),
F.approx_count_distinct("customer_id", 0.02).alias("customers"),
F.max_by("order_id", "order_ts").alias("last_order"))Three choices matter here. Bad rows are counted in quality, and the job should fail if the bad fraction passes a threshold, rather than quietly reporting lower revenue. Deduplication uses row_number with a total order, not first, so reruns pick the same row. The distinct-customer count states its accuracy, and the dashboard can say it is approximate. Running daily.explain() should show a single generated stage before the shuffle with no Python eval operator. If one appears, someone added a UDF.
Failure modes
| Symptom | Likely cause | Fix |
|---|---|---|
Job fails after a 4.0 upgrade with CAST_INVALID_INPUT | ANSI mode now default | Quarantine bad rows with try_cast; keep ANSI on |
| Daily totals differ between regions | Session time zone from the JVM | Set spark.sql.session.timeZone explicitly |
Rows missing after a != filter | NULL comparison is NULL | Add OR col IS NULL, or use <=> |
Rerun gives different first() values | Non-deterministic aggregate | max_by, or a window with a total order |
Stage slow, plan shows BatchEvalPython | Row-at-a-time Python UDF | Rewrite with built-ins or a pandas UDF |
AnalysisException: UNRESOLVED_ROUTINE | Function not in this version, or a typo | SHOW FUNCTIONS / DESCRIBE FUNCTION |
| Revenue silently lower than source | NULLs from try_ ignored by sum | Count nulled rows; fail above a threshold |
Most of these are found by the same habit: a small test DataFrame with NULLs, malformed strings, boundary dates and duplicates, run through the transformation in CI against expected output.
What to do next
- Run
SELECT version()andSET spark.sql.ansi.enabledon each cluster you use, and write down the answers. - Use
DESCRIBE FUNCTION EXTENDEDfor every function you rely on whose behaviour you are unsure of. - Replace any job-wide ANSI opt-out with targeted
try_calls, plus a count of nulled rows that can fail the job. - Set
spark.sql.session.timeZoneexplicitly in every job that buckets by day or hour. - Search your codebase for
first(,last(andcollect_list(, and make each use deterministic. - Run
explain()on your heaviest jobs, list every Python eval operator, and try to rewrite each one with built-ins behind an equivalence test. - Build a NULL-and-garbage fixture DataFrame and run it through your transformations in CI.