Unit Testing DataFrames

Unit Testing DataFrame Transformations

Testable Spark 129 code keeps reading, transforming and writing apart: transformations are plain functions from DataFrames to a DataFrame, with no paths, sessions or actions inside, so a test can feed them a few rows.

jobs/booknest/transforms.py: pure DataFrame-in, DataFrame-out functionsPython
"""BookNest's order transformations: pure DataFrame-in, DataFrame-out functions."""
from pyspark.sql import DataFrame, functions as F
def explode_lines(orders: DataFrame) -> DataFrame:
    """One row per order line; orders without items produce no lines."""
    return (orders.select("order_id", "customer_id", "order_ts", "status", F.inline("items"))
            .withColumn("amount", F.col("qty") * F.col("unit_price")))
def genre_revenue(lines: DataFrame, books: DataFrame) -> DataFrame:
    """Delivered revenue per genre; lines for unknown books count as 'Unknown'."""
    catalog = books.select(F.col("id").alias("book_id"), "genre")
    return (lines.where(F.col("status") == "delivered")
            .join(F.broadcast(catalog), "book_id", "left")
            .withColumn("genre", F.coalesce("genre", F.lit("Unknown")))
            .groupBy("genre").agg(F.sum("amount").alias("revenue")))
def order_size(orders: DataFrame) -> DataFrame:
    """Label each order small, medium or large by its total."""
    return orders.withColumn("size", F.when(F.col("total") >= 50, "large")
                             .when(F.col("total") >= 20, "medium").otherwise("small"))

Catalyst still optimizes chained functions as one plan, so the structure costs nothing at run time. Test the rules a reader cannot see at a glance: which statuses count, unknown books, and the size boundaries.