Lesson 24 / 25
Case Study: A Daily Sales Pipeline
Design a Spark batch job that turns raw order files into a partitioned revenue table.
The design
A job runs once a day for a given run_date. It reads the raw CSVs for that date with an explicit schema and fail-fast mode, cleans them (trim strings, cast types, fill or flag nulls, drop duplicate order_ids), joins with a small customer dimension using a broadcast hint, aggregates revenue per city per day, writes Parquet partitioned by order_date in overwrite mode for that partition only (idempotent: re-running the same date replaces that date), and runs data-quality checks (row count above a minimum, no null keys) before publishing. It is scheduled by an orchestrator such as Airflow, retries on transient failures, logs row counts and durations, and its Spark UI is checked when runtimes drift. Tests cover each transformation function with small DataFrames.
A daily pipeline end to end
Read, clean, join, aggregate, write partitioned output and monitor, as one reliable job.
The job in code (illustrative)
Dynamic partition overwrite replaces only the partitions present in the output. Paths, column names and the run_date argument are examples.
spark.conf.set("spark.sql.sources.partitionOverwriteMode", "dynamic")
raw = (spark.read.schema(SCHEMA).option("header", True).option("mode", "FAILFAST")
.csv(f"s3a://lake/raw/orders/{run_date}/"))
clean = (raw.dropDuplicates(["order_id"])
.withColumn("customer", F.trim("customer"))
.withColumn("amount", F.col("amount").cast("double")))
revenue = (clean.join(F.broadcast(dim_customers), "customer", "left")
.groupBy("order_date", "city")
.agg(F.sum("amount").alias("revenue"), F.count("*").alias("orders")))
assert revenue.count() > 0, "no rows produced"
(revenue.write.mode("overwrite").partitionBy("order_date")
.parquet("s3a://lake/gold/daily_revenue/"))Make reruns the normal case
Design every job so rerunning a date is safe. Then a failure at 3 a.m. is fixed by a retry or a backfill, not by manual surgery on tables.
Quick check: What makes re-running the same date safe in this design?
- Using random file names
- Appending to the table each time
- Overwriting only that date's partition produces the same result each time
- Skipping data-quality checks
Answer
Overwriting only that date's partition produces the same result each time — Idempotent, partition-scoped writes make retries and backfills harmless.