spark/core
intermediate
Connecting…

UDFs: the box Catalyst cannot see into

A Python UDF is a black box Catalyst cannot optimize. Built-ins, higher-order functions, then pandas_udf.

Lesson 18 of 29 · Spark path

Explain it at my level

  1. Built-in
  2. Black box
  3. Per-row hop
  4. Arrow batch
  5. Fix ladder
Watch the canvas:JVM built-inoptimizable filterblack box / row hoppython workercheapest optionLive simulation
JVM executor · Tungsten binary rowsTungsten rowsoff-heap, no copyupper(col) · built-inwhole-stage codegenfilterpushed belowrows outnever left the JVMmy_udf(col)BatchEvalPython — opaquefilter cannot move above the UDFevery row pays, even the ones you dropSerialization boundary · pickle or ArrowPython worker · one row at a timepickle → socket → CPython → pickle backone worker process per executor corebatchPython worker · Arrow batchespandas_udf — one Series, thousands of rowsone crossing per batch, not per rowFix ladder — try left to right1. Built-inupper, split, when2. Higher-ordertransform, filter, aggregate3. pandas_udfArrow batches4. Python UDFlast resortCatalyst knows upper(). That is why it can reorder around it.
A built-in stays in the JVM. A Python UDF blocks the optimizer and pays a process hop per row; pandas_udf pays it per batch.