The pandas API on Spark, imported as pyspark.pandas, lets you write pandas-style code that runs as distributed Spark jobs. It began as the Koalas project and has shipped inside PySpark since Spark 3.2. The promise is attractive: take a notebook that works on a sample, change the import, and run it on the full dataset. Often that is exactly what happens.
The cases where it does not are predictable once you see what the library does underneath. pandas assumes one machine, an ordered row index and eager execution; Spark assumes many machines, unordered partitions and lazy plans. This article explains how pyspark.pandas bridges those models, where the bridge costs a shuffle or a single-partition bottleneck, and how to write code that stays fast. Option names and defaults below are from the PySpark 4.2 documentation; defaults have changed between releases, so check yours with ps.get_option.
What a pandas-on-Spark DataFrame really is
Each pandas-on-Spark DataFrame holds an internal frame: a Spark DataFrame whose columns include the data columns and one or more index columns, plus metadata mapping pandas column labels to Spark columns. Almost every method returns a new wrapper around a new Spark plan. Nothing executes until an action needs a result: printing, len, to_pandas, writing a file, or a method whose output depends on data, such as computing the shape.
That makes the performance model Spark's, not pandas'. Filters, projections and aggregations become Catalyst plans, as described in Spark DataFrame and Dataset. A groupby().sum() becomes a partial aggregation, a shuffle and a final aggregation. A sort_values becomes a global sort with a range-partitioning shuffle. The pandas syntax is the same; the cost is what Spark would pay.
import pyspark.pandas as ps
ps.get_option("compute.default_index_type") # check the defaults of YOUR release
ps.get_option("compute.ops_on_diff_frames")
orders = ps.read_parquet("s3://lake/orders/", index_col="order_id") # a real column as index
orders["net"] = orders["amount"] - orders["discount"] # lazy: builds a plan
daily = (orders[orders["status"] == "paid"]
.groupby(["country", "order_date"])["net"]
.sum()
.reset_index())
daily.spark.explain() # inspect the Spark plan before running it
top = daily.sort_values("net", ascending=False).head(20) # global sort, then limit
print(top) # the action: this runs Spark jobs
The index problem
pandas requires every DataFrame to have an index. Spark tables have no row order and no index. When you read data without naming an index column, pyspark.pandas must invent one, and how it does so is set by compute.default_index_type.
| Type | How it is built | Cost and caveats |
|---|---|---|
sequence | A window function over the whole dataset without a partition | Moves all data into one partition on one executor; unusable on large data |
distributed-sequence | A distributed group-by and group-map that assigns consecutive numbers globally (the 4.2 default) | Extra indexing work and a cached intermediate (storage level from compute.default_index_cache); the rows behind each number can differ between computations |
distributed | Spark's monotonically_increasing_id() | Cheap, but values are non-deterministic and not consecutive; two frames almost never share index values |
The fix is usually to avoid the default index entirely. Most datasets already have a key. Pass it as index_col when reading or converting, and the index is simply a column, with no extra indexing step. If you need no meaningful index at all, the distributed type is the cheapest, as long as you never align two frames on it.
Worked example: a daily job reads 2 billion click rows from Parquet, computes features and writes them back. With the default index, the read adds an indexing step that numbers all 2 billion rows through a distributed group-map and keeps a cached intermediate. Reading with index_col="event_id" removes that step and the cache, and nothing downstream changes because the job never used positional access.
Ordering: the pandas assumption Spark does not keep
In pandas, row order is part of the data. In Spark, order exists only immediately after an explicit sort. pyspark.pandas makes this visible in a few places. head(n) returns some n rows, not necessarily the first ones, unless compute.ordered_head is True, which costs an ordering step. Positional operations such as shift, diff, cumsum and rank are implemented with window functions ordered by the index.
That last point is the most expensive trap in the library. A window with no partition key must see all rows in order, so Spark moves the whole dataset to one task. psdf["prev"] = psdf["price"].shift(1) on a large frame runs on a single core. The same operation after a groupby, for example psdf.groupby("symbol")["price"].shift(1), partitions the window by the group key and runs in parallel. Before using any positional method on the whole frame, ask whether it really means "per entity"; it almost always does.
Operations between different frames
In pandas, df1["a"] + df2["b"] aligns on the index. In Spark, two DataFrames are separate plans, so aligning them requires a join on the index columns. Whether this is allowed is set by compute.ops_on_diff_frames (True in the 4.2 documentation). When enabled, each cross-frame operation is a join, often with a shuffle, and it is only correct if both frames have a deterministic, shared index. With the distributed index type it silently pairs unrelated rows, which the documentation explicitly warns about.
Prefer to keep related columns in one frame, and join explicitly with merge on a real key when you need data from two sources. An explicit join states intent and makes the shuffle visible in the plan.
Moving between pandas, pandas-on-Spark and PySpark
from pyspark.sql import functions as F
sdf = spark.read.table("lake.events") # PySpark DataFrame
psdf = sdf.pandas_api(index_col="event_id") # to pandas-on-Spark, no default index job
psdf["hour"] = psdf["ts"].dt.hour
sdf2 = psdf.to_spark(index_col="event_id") # back to PySpark, keeping the index column
sdf2 = sdf2.withColumn("day", F.to_date("ts")) # use native Spark where it is clearer
# SQL over a pandas-on-Spark frame; {events} is substituted from the keyword argument
by_hour = ps.sql("SELECT hour, count(*) AS n FROM {events} GROUP BY hour", events=psdf)
small = by_hour.to_pandas() # collects to the driver: only small resultsTreat the three APIs as one toolbox. pandas_api(index_col=...) and to_spark(index_col=...) convert without copying data, and naming the index column avoids the default-index job and keeps the key. ps.sql runs SQL against pandas-on-Spark frames passed as keyword arguments. ps.from_pandas distributes a local pandas frame, and to_pandas collects a distributed one to the driver. The last is the most frequent cause of driver out-of-memory errors: it is only for results you have already reduced, such as a daily summary or a chart's data.
Some pandas methods are not implemented and raise an error. The 4.2 options list compute.pandas_fallback (default False), which falls back to pandas' implementation. Leave it off in production: a silent fallback on a large frame is a hidden collect.
Running your own pandas code
Built-in methods compile to Spark expressions. When you need arbitrary pandas logic, pyspark.pandas ships your function to Python workers, feeding them pandas DataFrames converted from Spark's columnar data with Apache Arrow.
import pandas as pd
import pyspark.pandas as ps
def sessionize(pdf: pd.DataFrame) -> ps.DataFrame["user_id": int, "session": int, "events": int]:
# Runs in a Python worker on ONE user's rows, as a plain pandas DataFrame.
pdf = pdf.sort_values("ts")
gap = pdf["ts"].diff() > pd.Timedelta(minutes=30)
pdf["session"] = gap.cumsum()
out = pdf.groupby("session").size().rename("events").reset_index()
out.insert(0, "user_id", pdf["user_id"].iloc[0])
return out
events = ps.read_parquet("s3://lake/clicks/", index_col="event_id")
sessions = events.groupby("user_id").apply(sessionize) # type hint: no schema-inference run
# Batch-wise pandas over arbitrary partitions (no grouping), output length may differ:
cleaned = events.pandas_on_spark.apply_batch(lambda pdf: pdf.dropna(subset=["ts"]))groupby(...).apply(f) gives your function all rows for one group as a pandas DataFrame. It is the tool for per-entity logic that has no Spark equivalent, such as sessionisation. Two cautions. First, every group must fit in one worker's memory, so a skewed key with millions of rows fails on one task while the others finish. Second, the output schema must be known. Without a return type hint, pyspark.pandas infers it by running your function on a sample limited by compute.shortcut_limit (1,000 rows by default), which runs your function an extra time and can infer the wrong types if the sample is unrepresentative. The return type hint shown above avoids both.
pandas_on_spark.apply_batch and pandas_on_spark.transform_batch apply a function to arbitrary batches of rows. Use them for row-local work such as cleaning or scoring with a model. transform_batch must return the same length as its input; apply_batch may not. Never compute a statistic that should be global inside a batch function: a mean over one batch is not the mean of the column.
Reading plans, caching and tuning
Because the cost model is Spark's, debug it with Spark's tools. psdf.spark.explain() prints the physical plan; count the Exchange nodes, since each is a shuffle, and look for a window without a partition key. The Spark UI shows which stage is slow and whether one task dominates. Adaptive query execution coalesces small shuffle partitions and splits skewed joins at runtime, as described in Spark AQE, and the mechanics of each exchange are in the Spark shuffle.
Lazy plans are recomputed for every action. If one expensive frame feeds several outputs, cache it explicitly and release it when done.
features = build_features(events) # expensive lineage, reused three times below
with features.spark.cache() as cached: # cached until the block exits
train = cached[cached["split"] == "train"]
test = cached[cached["split"] == "test"]
stats = cached.describe()
train.to_parquet("s3://lake/features/train/")
test.to_parquet("s3://lake/features/test/")Partition layout matters as much as in PySpark. Writing with to_parquet(path, partition_cols=[...]) on a low-cardinality column lets later reads prune files; Spark partitioning covers choosing partition counts and keys.
Semantics that differ from pandas
- Missing values. pandas uses NaN for missing floats and has several missing markers; Spark has SQL null. pyspark.pandas maps between them, but results of comparisons and aggregations over missing data can differ in edge cases. Test them.
- ANSI mode. Spark 4 enables SQL ANSI mode by default, under which some invalid operations raise errors rather than returning null. The 4.2 options include
compute.ansi_mode_supportandcompute.fail_on_ansi_mode(both True by default) to control how pandas-on-Spark behaves under it. Test division by zero, overflow and invalid casts in your pipeline instead of assuming pandas results. - Result size limits.
compute.max_rowsanddisplay.max_rows(1,000 each) bound some shortcut computations and printed output. They do not limit data processed. - Plotting. The default backend is plotly; top-n plots are limited by
plotting.max_rows, and sample-based plots useplotting.sample_ratio. A chart of a billion rows is always a summary.
Failure modes
| Symptom | Cause | Fix |
|---|---|---|
| Job stuck on one task | Positional method or sequence index over the whole frame | Partition the window by an entity key; use index_col |
| Extra indexing work and cached data on read | Default distributed-sequence index | Pass index_col, or use the distributed type if no alignment is needed |
| Driver out of memory | to_pandas or pandas fallback on large data | Aggregate first; keep compute.pandas_fallback off |
| Wrong values after combining frames | Cross-frame alignment on a non-deterministic index | Use merge on a real key |
| One groupby.apply task fails | Skewed group larger than worker memory | Salt or pre-aggregate the hot key; filter outliers |
| Function runs twice, wrong dtypes | Schema inference sample | Add a return type hint |
When to use it, and when not
Use pyspark.pandas to scale existing pandas analysis and feature code, for teams fluent in pandas, and for exploratory work on large tables. Drop to PySpark or SQL for production pipelines where you want the plan explicit, for joins and window logic that pandas expresses awkwardly, and anywhere positional semantics would force a single partition. Mixing them in one job is normal and costs nothing when you name index columns at each boundary.
What to do next
- Print
ps.get_optionfor the index type, cross-frame and ANSI options on your cluster and record them in the job's documentation. - Add
index_colto every read and conversion, and compare the Spark UI stages and storage tab before and after. - Search your code for
shift,diff,cumsumandrankon whole frames, and partition each by an entity key. - Add return type hints to every
groupby().applyand batch function, and check group sizes for skew. - Run
spark.explain()on your heaviest frame, count shuffles, and cache frames that feed more than one output. - Replace every
to_pandason unaggregated data with an aggregation or a file write.