Back to selected work

04 · Data engineering · ETL

Apple Purchase Data with PySpark

A tiny tutorial dataset rebuilt as a proper mixed-source pipeline with explicit schemas, auditable corrections, reusable transformations and automated checks.

PySpark pipeline connecting CSV, Delta and Parquet inputs to tested analytics outputs
3source formats behind one reader interface
5analytical workflows and output tables
0unmatched dimensions after standardisation
3Databricks notebooks executed end to end

Small data, proper engineering

The original learning exercise asks several customer-purchase questions using a small Apple dataset. Rather than leave the work as notebook cells, I rebuilt it as a package with a command-line workflow, mixed storage formats, explicit schemas, modular transformations and tests.

The dataset is deliberately tiny and synthetic. Its numerical findings are proof that the pipeline behaves correctly, not evidence about real Apple customers. The value of the project lies in the structure and checks.

Using a tiny dataset makes expected answers easy to calculate independently. That turns the data into a useful test fixture rather than pretending it demonstrates scale.

One interface across three source formats

Transactions remain in CSV, customer data is prepared as a Delta table and product data is written as Parquet. A Factory Pattern reader selects the correct implementation while the transformation layer receives consistent Spark DataFrames.

01
Transaction CSVEvent records with customer, product and purchase time.
RAW
02
Customer Delta + Product ParquetDimension-style sources deliberately stored in different formats.
PREPARE
03
Reader FactoryFormat-specific readers behind one extraction contract.
INGEST
04
Validation + transformationsKeys, nulls and relationships checked before business logic.
PROCESS
05
Parquet + Delta outputsFive repeatable analytical results written to storage.
LOAD

Explicit schemas replace type inference. Spark windows provide deterministic purchase ordering, while broadcast joins are used for the tiny customer and product dimensions.

A source mismatch left visible

The public source uses different names for the same products. The product master contains labels such as “iPhone SE”, “AirPods Pro” and “MacBook Air”, while transactions contain “iPhone”, “AirPods” and “MacBook”. An uncorrected join would match no product rows.

The raw files remain unchanged. During source preparation, the three names are standardised to the transaction vocabulary and the original value is retained in source_product_name.

CHECK 01

Null validation

Required identifiers and analytical fields are checked before processing.

CHECK 02

Duplicate keys

Dimension keys are tested before they are trusted in a join.

CHECK 03

Unmatched dimensions

Customer and product references must resolve after documented standardisation.

CHECK 04

Expected answers

Reference results are independently checkable without Spark.

Quietly replacing the raw product file would have made the example look neater while hiding the exact data-quality issue the pipeline needed to handle.

Five workflows, each with a clear rule

  • Immediate next purchase: use lead() over ordered customer transactions to find an iPhone followed directly by AirPods.
  • Exact product set: use collect_set() to find customers whose entire history contains only iPhone and AirPods, regardless of order.
  • Subsequent purchases: rank each customer’s events and retain an ordered list after the first purchase.
  • Average delay: calculate the number of days between the qualifying iPhone and immediately following AirPods purchase.
  • Top products: join product prices and rank revenue across the supplied records.

The distinction between “bought only these products” and “bought one immediately after the other” is tested directly. Customer 107 qualifies for the exact-set result but not the sequence result because that customer purchased AirPods before the iPhone.

Verified locally, in CI and in Databricks

The repository includes unit tests and a mixed-source integration test, plus a dependency-free Python reference script for checking expected business answers before Spark is installed.

All three notebooks were executed on 16 August 2026 using Databricks serverless compute and a Unity Catalog volume. The run completed source preparation, validation, all five workflows, Parquet and Delta writes and Spark optimisation-plan inspection.

The third notebook is intentionally explanatory. It displays formatted plans for repartitioning, coalescing, predicate pushdown and broadcast joins so optimisation choices can be inspected rather than merely named.

PySparkSpark SQLDelta LakeParquetDatabricksUnity CatalogpytestGitHub Actions
Next case studyOrbit QA Framework