Simulation/Spark

Educational simulation
Spark

UDFs vs built-ins

Watch rows cross the JVM to Python boundary, then batch them.

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.