Interview/Spark

Spark interview questions: UDFs: the box Catalyst cannot see into

Interactive Spark interview questions on UDFs: the box Catalyst cannot see into. Practice with the matching lesson. A Python UDF is a black box Catalyst cannot optimize. Built-ins, higher-order functions, then pandas_udf.

Lesson · Simulation

You wrote a Python UDF to uppercase a column. Why does a reviewer push back?

Answer it out loud, then reveal. Play steps through like the simulators.

Production scenario

UDFs: the box Catalyst cannot see into

Nightly job went from 12 minutes to 5 hours after a one-line phone-number cleanup UDF

Symptoms

  • The only change in the PR was a Python UDF on one column
  • Input Size / Records on the scan jumped by two orders of magnitude
  • Executor CPU is mostly in Python worker processes

All questions on this page

Indexed as FAQ. Open any item if you prefer a list to Play.

beginner

You wrote a Python UDF to uppercase a column. Why does a reviewer push back?

UDFs: the box Catalyst cannot see into · tap to open the answer

Short: Every row leaves the JVM for a Python process and comes back; upper() never leaves the JVM at all.

Detailed: A Python UDF serializes rows to a Python worker, evaluates them one at a time, and serializes results back — it shows up as BatchEvalPython in the plan. Built-ins like upper, regexp_replace, and to_timestamp compile into whole-stage codegen and operate on Tungsten binary rows directly.

Common mistake: Assuming the UDF is fine 'because it's just string formatting'.

Follow-up: Name a built-in you would use instead of a UDF you have written before.

Lesson · Simulation

intermediate

You moved a filter into a UDF and the scan went from 4 GB to 900 GB. Explain.

UDFs: the box Catalyst cannot see into · tap to open the answer

Short: Catalyst cannot see inside a UDF, so the predicate is not pushed into the scan and no partitions get pruned.

Detailed: A UDF is an opaque expression: it can never become a PushedFilter or a PartitionFilter, so Spark reads every file and filters afterwards above BatchEvalPython. explain('formatted') shows PushedFilters empty and the FileScan's Input Size / Records covering the whole table.

Common mistake: Blaming 'Python is slow' when the real cost is the extra 896 GB of I/O.

Follow-up: How do you keep the custom logic but restore pruning?

Lesson · Simulation

senior

Walk me up the fix ladder for a row-at-a-time Python UDF, cheapest first.

UDFs: the box Catalyst cannot see into · tap to open the answer

Short: Built-in expression, then higher-order functions, then pandas_udf, and only then a plain Python UDF.

Detailed: Built-ins keep codegen and pushdown. Higher-order functions (transform, filter, exists, aggregate) handle array and map logic without leaving SQL. pandas_udf ships columnar Arrow batches, so you pay serialization per batch instead of per row and get vectorized pandas or NumPy code. A row-wise Python UDF is the last resort.

Senior: It is still opaque to Catalyst, so nothing pushes through it, and on Databricks it still knocks the query off Photon. Faster, not free.

Common mistake: Jumping straight to pandas_udf when composing two built-ins would have done it.

Follow-up: What does a pandas UDF still not fix?

Lesson · Simulation

architect

A Scala UDF runs inside the JVM, so is it safe to use anywhere?

UDFs: the box Catalyst cannot see into · tap to open the answer

Short: It skips the Python round trip, but it is still a black box to the optimizer.

Detailed: A Scala or Java UDF executes in the executor JVM with no cross-process serialization, so throughput is far better than Python. Catalyst still cannot rewrite or push it, it can break whole-stage codegen for the operators around it, and null handling is entirely your responsibility. A native Expression or a composition of built-ins is still preferable when the logic can be expressed.

Senior: Rerun the query with that column replaced by a literal. If runtime collapses, it is the UDF; if not, you were about to optimize the wrong operator.

Common mistake: Treating 'we made it a Scala UDF' as equivalent to a native expression in a design review.

Follow-up: How do you prove the UDF is or is not your bottleneck?

Lesson · Simulation