A production-grade PySpark ETL/ELT pipeline that ingests, validates, cleanses, normalises, and enriches large-scale e-commerce transactional data — deployable on local Spark, Databricks, AWS EMR, Azure Synapse, and HDInsight.
This project implements a full end-to-end batch ETL pipeline for e-commerce sales data, showcasing the exact skills required for modern data engineering roles:
| Skill | Implementation |
|---|---|
| PySpark ETL/ELT pipelines | src/pipeline.py + src/transformation/ |
| Distributed computing (Databricks, EMR, Synapse) | src/utils/spark_session.py — env-aware session factory |
| Data quality: validation, cleansing, normalisation | src/quality/data_quality.py |
| Clean, reusable Python + PySpark code | Modular package structure with typed schemas |
| Large dataset processing | Explicit schemas, adaptive query execution, Parquet I/O |
Raw CSV Sources
│
▼
┌─────────────┐
│ INGEST │ DataLoader — typed schemas, cloud-path aware
└──────┬──────┘
│
▼
┌─────────────┐
│ VALIDATE │ DataQualityChecker — null %, dupes, RI, business rules
└──────┬──────┘
│
▼
┌─────────────┐
│ CLEANSE │ Transformer.cleanse_*() — dedup, null handling, type casting
└──────┬──────┘
│
▼
┌─────────────┐
│ NORMALIZE │ Transformer.normalize_*() — date parsing, rounding, calendar cols
└──────┬──────┘
│
▼
┌─────────────┐
│ ENRICH │ build_fact_table() — joins + KPI derivation
└──────┬──────┘
│
▼
┌─────────────┐
│ LOAD │ Parquet / Delta / CSV → cloud storage or local FS
└─────────────┘
ecommerce-data-pipeline/
├── config/
│ └── pipeline_config.yaml # All env/path/threshold settings
├── data/
│ ├── generate_sample_data.py # Synthetic data generator (with intentional dirt)
│ ├── raw/ # Input CSV files (auto-generated)
│ └── processed/ # Parquet outputs + DQ reports
├── src/
│ ├── pipeline.py # Main ETL orchestrator
│ ├── ingestion/
│ │ └── data_loader.py # SparkDataFrame loaders with explicit schemas
│ ├── transformation/
│ │ └── transformations.py # Cleanse → Normalize → Enrich → Fact/Dim build
│ ├── quality/
│ │ └── data_quality.py # DQ framework: checks, thresholds, JSON report
│ └── utils/
│ ├── spark_session.py # Cloud-aware SparkSession factory
│ ├── config_loader.py # YAML config loader + validator
│ └── logger.py # Rotating file + console logger
├── notebooks/
│ └── pipeline_demo.py # Databricks-ready notebook (EDA + pipeline)
├── tests/
│ └── test_transformations.py # pytest unit tests (local Spark)
├── requirements.txt
└── README.md
- Python 3.9+
- Java 8 or 11 (required by PySpark)
pip install -r requirements.txtpython data/generate_sample_data.pypython -m src.pipelinepytest tests/ -v --tb=short- Upload the repo to DBFS or mount an ADLS/S3 volume.
- Open
notebooks/pipeline_demo.pyas a Databricks notebook. - Update
pathsinpipeline_config.yamltodbfs:/orabfss://paths. SparkSessionis auto-detected — no code changes needed.
spark-submit \
--master yarn \
--deploy-mode cluster \
--py-files src.zip \
src/pipeline.py config/pipeline_config.yamlSet the AZURE_CLIENT_ID environment variable; the session factory will configure Synapse-compatible settings automatically.
The pipeline enforces configurable thresholds before writing any outputs:
| Check | Threshold | Action |
|---|---|---|
| Null % on critical columns | ≤ 5% | FAIL |
| Duplicate rows | ≤ 2% | FAIL |
| Positive amounts & quantities | < 1% bad | FAIL |
| Valid date formats | ≤ 5% malformed | FAIL |
| Email format | < 2% malformed | WARN |
| Referential integrity | < 1% orphans | FAIL |
A JSON quality report is written to data/processed/quality_report/ after every run.
| Table | Description |
|---|---|
sales_facts |
Wide fact table: orders enriched with customer + product dims + KPIs |
dim_customers |
SCD Type-1 customer dimension |
dim_products |
Product dimension |
Derived KPIs:
gross_revenue,estimated_profit,profit_margin_pctis_high_value(order ≥ ₹5,000)customer_segment(Premium / Mid-Tier / Standard)net_unit_price,order_age_days, calendar columns
Tests use PySpark in local mode — no cluster required for CI/CD.
tests/
└── test_transformations.py
├── TestCleanseOrders (4 tests)
├── TestCleanseCustomers (2 tests)
├── TestNormalizeOrders (2 tests)
├── TestBuildFactTable (2 tests)
└── TestDataQuality (3 tests)
- PySpark 3.5 — distributed data processing
- Python 3.9+ — pipeline orchestration, testing, config
- Apache Parquet — columnar output format
- YAML — pipeline configuration
- pytest — unit testing
- Databricks / EMR / Synapse / HDInsight — cloud runtimes
MIT — free to use, modify, and distribute.