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.
"""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.