Lesson 22 / 25
Testing Spark Jobs
Structure transformations as functions and test them with a local SparkSession.
Small functions, small test data
Write transformations as plain functions that take a DataFrame and return a DataFrame (def add_revenue(df): ...), keeping reading and writing at the edges. Then unit-test each function with a few rows in a local SparkSession (one shared fixture per test run, because starting Spark is slow) and compare results using collect() or helper assertions. Add data-quality checks inside the pipeline: row counts, null rates, uniqueness of keys, value ranges and referential checks, failing the job (or quarantining bad rows) when expectations break. Test against realistic edge cases: nulls, duplicates, empty input, timezone boundaries.
A unit test with a local session (illustrative)
The function under test is pure: DataFrame in, DataFrame out.
import pytest
from pyspark.sql import SparkSession, functions as F
@pytest.fixture(scope="session")
def spark():
s = SparkSession.builder.master("local[2]").config("spark.ui.enabled", "false").getOrCreate()
yield s
s.stop()
def add_revenue(df):
return df.withColumn("revenue", F.col("qty") * F.col("price"))
def test_add_revenue(spark):
df = spark.createDataFrame([(2, 10.0), (3, 5.0)], ["qty", "price"])
out = add_revenue(df).orderBy("qty").collect()
assert [r["revenue"] for r in out] == [20.0, 15.0]Quick check: Why write transformations as functions of DataFrames?
- It removes the need for tests
- Spark forbids anything else
- Functions run faster by definition
- They are easy to unit-test with small local data
Answer
They are easy to unit-test with small local data — Pure functions separate logic from I/O, so tests need no cluster.