# Case Study: A Daily Sales Pipeline — Apache Spark: Big Data Processing with DataFrames

Source: https://www.geekswithgeeks.com/en/spark/wrap-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_id`s), **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.

![Four stages: ingest, transform, publish, observe.](assets/figures/spark/section-8-map.svg) — Figure 8.1 — Ingest, transform, publish and observe.

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

```python
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.

**Quiz:** What makes re-running the same date safe in this design?

- [ ] Using random file names
- [ ] Appending to the table each time
- [x] 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.
