Scalar vs Grouped UDFs

Scalar Versus Grouped Pandas UDFs

The type hints select the kind of pandas 16,086 UDF. Scalar forms map a Series to a Series of the same length, or an iterator of Series to an iterator, which lets you load something expensive (a model, a lookup table) once per task rather than once per batch. A Series-to-scalar hint makes a grouped aggregate UDF for groupBy().agg() and windows.

A grouped aggregate pandas UDF against the built-in percentileJavaScript
import time
import pandas as pd
from pyspark.sql.functions import pandas_udf
@pandas_udf("double")
def p90_pd(total: pd.Series) -> float:                # Series to scalar: a grouped aggregate
    return float(total.quantile(0.9))
amounts = orders.select("channel", F.col("total").cast("double").alias("total"))
variants = [("built-in", F.percentile("total", 0.9)), ("pandas UDF", p90_pd("total"))]
for name, agg in variants * 2:
    t0 = time.perf_counter()                          # second round: warm
    rows = amounts.groupBy("channel").agg(agg.alias("p90")).orderBy("channel").collect()
    secs = time.perf_counter() - t0
    print(f"{name:<10} {secs:5.2f}s", [(r.channel, round(r.p90, 2)) for r in rows])
Output
built-in    9.21s [('android', 67.15), ('ios', 67.05), ('web', 67.15)]
pandas UDF 10.16s [('android', 67.15), ('ios', 67.05), ('web', 67.15)]
built-in    1.96s [('android', 67.15), ('ios', 67.05), ('web', 67.15)]
pandas UDF  3.84s [('android', 67.15), ('ios', 67.05), ('web', 67.15)]

Warm, the grouped UDF took about twice as long. The built-in percentile merges partial states from every task; the pandas UDF must first shuffle all of a channel's amounts to one place, so three channels keep at most three Python workers busy. Use grouped aggregates only when no built-in exists.