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.

Two paths for the same column logic: a built-in stays inside Catalyst, a UDF is a black boxF.lower(col)or SQL: lower(email)udf(fn)(col)Python or Scala UDFAnalyzername -> registry -> ExpressionOptimizerfold, simplify, push downCodegenfused JVM loopAnalyzerresolves as opaque callOptimizercannot look insideEval operatorrows out to a workerBuilt-in: null handling, types and nullability are known, so filters can reach the file scan.Built-in under ANSI mode: bad input raises a named error, or use a try_ variant to get NULL.UDF: Spark only knows the declared return type; Python UDFs also serialize every row (or Arrow batch).Look for BatchEvalPython or ArrowEvalPython in explain() output to find them.
The same column logic as a built-in and as a UDF. Only the top path is visible to the optimizer and code generator.

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:

FamilyExamplesWhere it is legal
Scalarlower, regexp_extract, to_date, coalesceAnywhere a column is accepted
Aggregatesum, count_if, max_by, approx_count_distinctgroupBy().agg(), or over a window
Window-onlyrow_number, rank, lag, leadOnly with over(Window...)
Generatorexplode, posexplode, inlineOne per select; changes the row count
Table-valuedrange, explode in FROMFROM 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 and count(*) does not. A ratio built from the two mixes populations.
  • concat returns NULL if any argument is NULL, while concat_ws skips NULL arguments.
  • when(cond, v) without otherwise yields NULL for unmatched rows.
  • sum and avg ignore 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 in coalesce(sum(x), 0) on purpose.
  • a = b is NULL when either side is NULL. Use a <=> b in SQL or eqNullSafe in the DataFrame API when NULL should equal NULL, for example in change detection or join keys.
  • greatest and least skip NULLs, unlike arithmetic, so greatest(a, b) and a + b disagree 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:

ExpressionANSI offANSI on (4.0+ default)
2147483647 + 1 on intwraps to -2147483648ARITHMETIC_OVERFLOW error
CAST('a' AS INT)NULLCAST_INVALID_INPUT error
10 / 0NULLdivision-by-zero error
element_at(arr, 99)NULLindex-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 vanish

The 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_list and collect_set return elements in whatever order the shuffle delivered them. That order can change between runs. Sort inside the array with array_sort, or aggregate a struct whose first field is the sort key.
  • first and last in a groupBy are not deterministic for the same reason. To get the value of x at the latest ts, use max_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 smaller rsd for 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) and sum(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 codegen

In 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

SymptomLikely causeFix
Job fails after a 4.0 upgrade with CAST_INVALID_INPUTANSI mode now defaultQuarantine bad rows with try_cast; keep ANSI on
Daily totals differ between regionsSession time zone from the JVMSet spark.sql.session.timeZone explicitly
Rows missing after a != filterNULL comparison is NULLAdd OR col IS NULL, or use <=>
Rerun gives different first() valuesNon-deterministic aggregatemax_by, or a window with a total order
Stage slow, plan shows BatchEvalPythonRow-at-a-time Python UDFRewrite with built-ins or a pandas UDF
AnalysisException: UNRESOLVED_ROUTINEFunction not in this version, or a typoSHOW FUNCTIONS / DESCRIBE FUNCTION
Revenue silently lower than sourceNULLs from try_ ignored by sumCount 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() and SET spark.sql.ansi.enabled on each cluster you use, and write down the answers.
  • Use DESCRIBE FUNCTION EXTENDED for 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.timeZone explicitly in every job that buckets by day or hour.
  • Search your codebase for first(, last( and collect_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.
Key takeaway: A Spark SQL function call builds a Catalyst expression, and built-ins are fast because the optimizer and code generator can see inside them. Learn the NULL rules, keep ANSI mode on and use try_ functions with a count of nulled rows, set the session time zone, make aggregates deterministic, and use explain() to find the UDFs worth rewriting.