How to Set Up Unit Testing for PySpark Applications with pytest
Unit testing for PySpark applications with pytest relies on a session-scoped fixture that instantiates a local SparkSession once and shares it across all tests, enabling fast, deterministic DataFrame assertions without requiring a cluster.
Testing PySpark pipelines traditionally requires heavy infrastructure, but the DataExpert-io/data-engineer-handbook repository demonstrates a lightweight approach using pytest fixtures. By running Spark in local mode and injecting a shared SparkSession into test functions, you can validate transformation logic using small in-memory DataFrames. This guide explains how to implement unit testing for PySpark applications with pytest using the exact patterns found in the source code.
Create a Session-Scoped SparkSession Fixture
The foundation of PySpark unit testing is a reusable SparkSession that lives in conftest.py. According to the source code in intermediate-bootcamp/materials/3-spark-fundamentals/src/tests/conftest.py, the fixture uses session scope to ensure a single local Spark cluster starts once and persists for the entire test run.
# conftest.py
import pytest
from pyspark.sql import SparkSession
@pytest.fixture(scope='session')
def spark():
return SparkSession.builder \
.master("local") \
.appName("chispa") \
.getOrCreate()
Setting scope='session' is critical for performance. Without it, pytest would instantiate a new SparkSession for every test function, causing significant overhead. The master("local") configuration runs Spark in single-node mode, eliminating the need for external cluster resources while maintaining full API compatibility.
Structure a PySpark Unit Test
Each test follows the Arrange-Act-Assert pattern. The test function declares the spark parameter, which pytest automatically injects from the fixture. You then create input DataFrames using spark.createDataFrame(), invoke your production logic, and assert on the results.
Consider a simple transformation in src/jobs/example_job.py:
# src/jobs/example_job.py
def add_one(df):
"""Add a column with value +1."""
return df.withColumn("value_plus_one", df["value"] + 1)
The corresponding test in src/tests/test_example_job.py demonstrates the complete pattern:
# src/tests/test_example_job.py
def test_add_one(spark):
# Arrange – create a tiny DataFrame
data = [(1,), (2,), (3,)]
df = spark.createDataFrame(data, ["value"])
# Act – run the job logic
result = add_one(df)
# Assert – collect and compare
expected = [(1, 2), (2, 3), (3, 4)]
assert result.select("value", "value_plus_one").collect() == expected
Because the DataFrames are small and reside in memory, tests execute in milliseconds while still validating the exact transformation logic used in production.
Testing Production Job Functions
The repository separates concerns cleanly between setup, logic, and verification. Production code lives in the jobs/ package, while tests reside in tests/. For example, intermediate-bootcamp/materials/3-spark-fundamentals/src/jobs/monthly_user_site_hits_job.py contains the aggregation logic, tested by intermediate-bootcamp/materials/3-spark-fundamentals/src/tests/test_monthly_user_site_hits.py.
Similarly, graph-processing jobs using Spark GraphFrames follow the same pattern. The team_vertex_job.py implementation is validated by test_team_vertex_job.py, both located in their respective jobs/ and tests/ directories. This architecture ensures that job functions remain pure and testable, accepting DataFrames as input and returning DataFrames as output, with no dependency on external Spark session management.
Execute the Test Suite
Running the suite requires no additional configuration beyond pytest discovery. From the package root, execute:
python -m pytest
Pytest automatically discovers the tests/ directory, injects the spark fixture into any test function that requests it, and executes assertions against the collected DataFrame rows. For targeted execution of specific modules:
python -m pytest src/tests/test_monthly_user_site_hits.py
Summary
- Session-scoped fixtures in
conftest.pyprovide a shared SparkSession initialized once per test run, dramatically improving performance over function-scoped alternatives. - Local mode (
master("local")) eliminates cluster dependencies while maintaining full Spark API fidelity for unit testing. - Pure job functions in the
jobs/package accept and return DataFrames, enabling isolated testing with small in-memory datasets created viaspark.createDataFrame(). - pytest handles fixture injection and test discovery automatically, requiring only standard
python -m pytestinvocation to validate the entire suite.
Frequently Asked Questions
Why use scope='session' for the SparkSession fixture?
Session scope ensures the SparkSession initializes exactly once for the entire test run, reducing startup overhead from seconds to milliseconds. According to the DataExpert-io/data-engineer-handbook implementation, this pattern prevents the costly operation of creating multiple local Spark clusters when running dozens or hundreds of test cases.
How do I test DataFrame transformations without a Spark cluster?
Configure the fixture with .master("local") to run Spark in local mode using a single JVM process. This approach, implemented in conftest.py, provides full Spark functionality for small datasets without requiring YARN, Kubernetes, or standalone cluster managers.
Where should I place my PySpark test fixtures?
Define the spark fixture in a conftest.py file located at the root of your test directory (e.g., src/tests/conftest.py). pytest automatically discovers and shares this fixture across all test modules in that directory tree, as demonstrated in intermediate-bootcamp/materials/3-spark-fundamentals/src/tests/.
Can I compare DataFrame contents directly in assertions?
Yes, by calling .collect() on the DataFrame to return a list of Row objects, or by converting to pandas using toPandas() for richer comparison utilities. The repository examples use .collect() for simple equality checks against expected tuples, ensuring tests remain fast and dependency-light.
Have a question about this repo?
These articles cover the highlights, but your codebase questions are specific. Give your agent direct access to the source. Share this with your agent to get started:
curl -s "https://instagit.com/install.md" Maintain an open-source project? Get it listed too →