PySpark
PySpark is the Python API for Apache Spark, and the most popular way Spark is used today. It lets data engineers and data scientists write distributed data pipelines, SQL, streaming jobs, and ML workflows in Python, while execution happens in Spark's JVM engine across a cluster.
That split is the key to using PySpark well. DataFrame operations written in Python are just instructions to build a plan, so they run at full JVM speed. Python code that touches individual rows (UDFs, rdd.map, collect() loops) moves data between the JVM and Python processes, which is where most PySpark performance problems come from. Knowing which side of that boundary your code runs on is most of the craft.
TL;DR
- PySpark DataFrame and SQL code builds JVM plans, so it's as fast as Scala for built-in operations.
- Python UDFs run row-by-row in Python workers, which is slow. Use built-in functions first, then pandas UDFs (Arrow-vectorized).
- Spark Connect decouples the Python client from the cluster via a gRPC protocol, for thin clients and better isolation.
- Arrow speeds up
toPandas()andcreateDataFrame(pandas_df). Still, nevertoPandas()big data. - The pandas API on Spark (
pyspark.pandas) offers pandas syntax at scale. - Package dependencies consistently (conda/venv archives, container images) and unit-test transformations with small local sessions.
Quick Example
Core Concepts
Architecture: Python Driver, JVM Engine
In classic PySpark, your Python driver program talks to a JVM driver through Py4J. DataFrame calls create JVM plan objects, and nothing is computed in Python. When a plan contains Python UDFs, executors launch Python worker processes and stream rows (or Arrow batches) between the JVM and Python, which adds serialization overhead.
Spark Connect
Spark Connect (Spark 3.4+, and the default direction in Spark 4) introduces a client-server protocol: a thin Python client sends unresolved plans over gRPC to a Spark Connect server. Benefits: lightweight clients (notebooks, IDEs, apps) without a local JVM, better isolation between users, easier upgrades, and remote development against shared clusters. Most DataFrame APIs work the same. Some low-level APIs (RDDs, SparkContext access) aren't available.
UDF Options
Python UDFs are opaque to Catalyst: filters after them can't be pushed down, and optimizations are limited.
Arrow and pandas Interop
With Arrow enabled, df.toPandas() and spark.createDataFrame(pdf) transfer columnar batches efficiently. But toPandas() still collects all rows to the driver, so only use it on small, aggregated results. For larger data, write to Parquet and read it with pandas or Polars, or keep processing in Spark.
pandas API on Spark
import pyspark.pandas as ps gives a pandas-compatible API backed by Spark: psdf.groupby(...).mean(), ps.read_parquet, and so on. It's handy for scaling existing pandas code and for analysts, but some pandas semantics (row order, index operations) are expensive in a distributed engine, so learn the differences.
Dependencies and Packaging
Executors need the same Python packages as the driver. Options:
- Container images (Kubernetes, EMR on EKS, Dataproc) with pinned dependencies. It's the most reproducible.
- Packed environments:
conda-pack,venv-pack, orpexarchives shipped with--archives. --py-filesfor your own modules (zip or wheel).- Managed platforms offer cluster libraries and serverless environment specs.
Pin versions: a mismatch between the driver and executor Python or library versions causes confusing errors. See Python packaging.
Testing PySpark Code
Structure pipelines as pure functions from DataFrames to DataFrames, and test them with a local SparkSession and tiny inputs:
Reduce shuffle partitions in tests for speed, and use a session-scoped fixture.
Best Practices
Stay in the DataFrame API
Express logic with pyspark.sql.functions, SQL expressions, and higher-order array functions (transform, filter, aggregate) before reaching for Python. Most "I need a UDF" cases have a built-in solution.
Use Type Hints and Small Modules
Typed, well-named transformation functions compose cleanly and make pipelines readable. Avoid giant notebooks full of chained code with no tests. See type hints.
Broadcast Models and Lookup Data Wisely
Load ML models once per executor process (module-level caching inside mapInPandas, or broadcast variables for small objects), not once per row or batch.
Mind the Driver
Keep driver work light: avoid collect() loops, large toPandas(), and building huge Python lists of paths or IDs. Size driver memory for broadcasts and results.
Common Mistakes
Looping Over Rows in Python
Creating Many Small Spark Jobs in a Loop
Running df.filter(col == x).count() for each of 1,000 values triggers 1,000 jobs. Use one groupBy(...).count() instead.
Mismatched Python Environments
Different pandas, NumPy, or PyArrow versions on the driver and executors cause serialization errors or subtle bugs. Pin versions, and ship consistent environments.
FAQ
Is PySpark slower than Scala Spark?
Not for DataFrame and SQL operations using built-in functions, which compile to the same JVM plans. PySpark is slower when Python code processes rows (UDFs, RDD lambdas), because of serialization between the JVM and Python. Pandas UDFs with Arrow narrow that gap considerably.
What is Spark Connect?
A client-server architecture for Spark that lets thin clients (Python, Scala, Go, and others) send DataFrame plans to a remote Spark cluster over gRPC. It decouples client environments from the cluster and simplifies remote development and multi-tenant use.
When should I use pandas UDFs?
When you need custom Python logic, like ML model inference, complex string or numeric processing, or library functions, that isn't available as a built-in Spark function. Pandas UDFs process data in vectorized Arrow batches, and they're far faster than row-by-row Python UDFs.
Should I use PySpark or pandas?
Use pandas (or Polars or DuckDB) when data fits comfortably on one machine. It's simpler and often faster for small to medium data. Use PySpark when data or processing exceeds one machine, when you need distributed processing, or when integrating with a Spark-based lakehouse.
Related Topics
- Apache Spark — Pillar overview
- Spark DataFrames — The core API
- Spark Performance Tuning — UDF and memory costs
- Pandas — Single-machine DataFrames
- Pandas Performance — When to move beyond pandas
- Python — The language