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

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:

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

References