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.