Pandas UDFs map columns to a column. The pandas 16,086 function APIs hand you whole DataFrames instead and return DataFrames of any shape: groupBy(...).applyInPandas() (grouped map) runs a function on all rows of each group, mapInPandas() runs one on an iterator of batches, and groupBy(...).cogroup(...).applyInPandas() pairs the groups of two DataFrames. You declare the output schema, because Spark 129 cannot infer it from Python.
import pandas as pd
def gaps(pdf: pd.DataFrame) -> pd.DataFrame: # all orders of one customer
ts = pdf["order_ts"].sort_values()
days = ts.diff().dt.total_seconds() / 86400
return pd.DataFrame({"customer_id": [pdf["customer_id"].iloc[0]], "orders": [len(ts)],
"median_gap_days": [round(days.median(), 1)]})
per_customer = (orders.select("customer_id", "order_ts").groupBy("customer_id")
.applyInPandas(gaps, "customer_id int, orders long, median_gap_days double"))
per_customer.orderBy(F.desc("orders")).show(3)Output
+-----------+------+---------------+ |customer_id|orders|median_gap_days| +-----------+------+---------------+ | 48907| 203| 1.7| | 45227| 202| 2.0| | 27995| 201| 1.8| +-----------+------+---------------+ only showing top 3 rows
A customer's median gap needs that customer's sorted history: natural in pandas, awkward in SQL. The price is a shuffle, and each group is loaded whole into one Python worker, so one huge group can fail the task. mapInPandas (a generator over batches, no shuffle) suits filters, enrichment and model scoring.