This PySpark ETL pipeline processes large-scale Amazon sales data to generate comprehensive business insights and KPIs. The pipeline handles data cleaning, transformation, and aggregation with scalable distributed processing.
- β Distributed data processing with PySpark
- β Comprehensive data quality checks and cleaning
- β 10+ business KPIs calculated automatically
- β Parquet output format for efficient storage
- β HDFS and local filesystem support
- β Optional Kafka integration for real-time processing
- Monthly Revenue - Revenue trends over time
- Profit Margin - Overall profitability analysis
- Region-wise Sales - Geographic performance breakdown
- Average Order Value (AOV) - Customer spending patterns
- Category Performance - Product category analysis
- Fulfillment Performance - Amazon vs Merchant comparison
- B2B vs B2C Analysis - Customer segment insights
- Promotion Impact - Promotional effectiveness
- Top Products - Best-selling items by revenue
- Order Status Distribution - Order completion rates
# Required Software
- Python 3.8+
- Apache Spark 3.x
- Java 8 or 11
- Hadoop (optional, for HDFS)# Using pip
pip install pyspark
# Or with conda
conda install -c conda-forge pyspark# Check Spark installation
pyspark --version
# Test in Python
python3 -c "from pyspark.sql import SparkSession; print('PySpark installed successfully!')"Download Cleaned_Amazon_Sale_Report.csv from the provided SharePoint link and place it in your project directory.
# Navigate to project directory
cd /path/to/project
# Run the pipeline
python pyspark_etl_pipeline.py# Submit to Spark cluster
spark-submit \
--master local[*] \
--driver-memory 4g \
--executor-memory 4g \
pyspark_etl_pipeline.py# First, upload data to HDFS
hdfs dfs -put Cleaned_Amazon_Sale_Report.csv /data/sales/
# Modify the script configuration:
# INPUT_PATH = "hdfs://localhost:9000/data/sales/Cleaned_Amazon_Sale_Report.csv"
# OUTPUT_PATH = "hdfs://localhost:9000/data/output/"
# Run with HDFS paths
spark-submit \
--master yarn \
--deploy-mode cluster \
pyspark_etl_pipeline.pyDatabricks:
- Upload script to Databricks workspace
- Create new notebook
- Import the Python script
- Run all cells
AWS EMR:
aws emr add-steps \
--cluster-id j-XXXXXXXXXXXXX \
--steps Type=Spark,Name="Amazon Sales ETL",\
ActionOnFailure=CONTINUE,\
Args=[--deploy-mode,cluster,--master,yarn,s3://your-bucket/pyspark_etl_pipeline.py]After successful execution, you'll find:
output/
βββ enriched_sales_data/ # Full cleaned dataset
β βββ part-00000.parquet
β βββ part-00001.parquet
β βββ _SUCCESS
βββ kpis/
β βββ monthly_revenue/ # Monthly revenue metrics
β βββ profit_margin/ # Profit analysis
β βββ region_sales/ # Geographic breakdown
β βββ aov/ # Average order value
β βββ category_performance/ # Category metrics
β βββ fulfillment_performance/ # Fulfillment analysis
β βββ b2b_analysis/ # B2B vs B2C
β βββ promotion_impact/ # Promotion effectiveness
β βββ top_products/ # Best sellers
β βββ status_distribution/ # Order statuses
βββ summary_report/ # High-level summary
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("ReadResults").getOrCreate()
# Read enriched data
df = spark.read.parquet("output/enriched_sales_data")
df.show()
# Read specific KPI
monthly_revenue = spark.read.parquet("output/kpis/monthly_revenue")
monthly_revenue.show()import pandas as pd
# Read parquet files
df = pd.read_parquet("output/enriched_sales_data")
print(df.head())
# Read KPI
monthly_rev = pd.read_parquet("output/kpis/monthly_revenue")
print(monthly_rev)-- In Spark SQL or Hive
CREATE EXTERNAL TABLE enriched_sales
STORED AS PARQUET
LOCATION 'hdfs://output/enriched_sales_data';
SELECT * FROM enriched_sales LIMIT 10;Edit these variables in the main script:
# Input/Output paths
INPUT_PATH = "Cleaned_Amazon_Sale_Report.csv"
OUTPUT_PATH = "./output"
# Spark configuration
.config("spark.sql.shuffle.partitions", "200") # Adjust based on data size
.config("spark.driver.memory", "4g")
.config("spark.executor.memory", "4g")# For large datasets (>10GB)
spark = SparkSession.builder \
.config("spark.sql.shuffle.partitions", "400") \
.config("spark.default.parallelism", "400") \
.config("spark.executor.instances", "10") \
.config("spark.executor.cores", "4") \
.config("spark.executor.memory", "8g") \
.getOrCreate()
# For small datasets (<1GB)
spark = SparkSession.builder \
.config("spark.sql.shuffle.partitions", "50") \
.getOrCreate()# Add to imports
from pyspark.sql.streaming import StreamingQuery
# Kafka configuration
def read_from_kafka(spark, kafka_brokers, topic):
"""Read streaming data from Kafka"""
df = spark.readStream \
.format("kafka") \
.option("kafka.bootstrap.servers", kafka_brokers) \
.option("subscribe", topic) \
.option("startingOffsets", "latest") \
.load()
# Parse JSON from Kafka
from pyspark.sql.functions import from_json, col
from pyspark.sql.types import StructType, StructField, StringType, DoubleType
schema = StructType([
StructField("Order_ID", StringType()),
StructField("Amount", DoubleType()),
# Add other fields...
])
parsed_df = df.select(
from_json(col("value").cast("string"), schema).alias("data")
).select("data.*")
return parsed_df
# Use in pipeline
kafka_df = read_from_kafka(spark, "localhost:9092", "orders")spark-submit \
--packages org.apache.spark:spark-sql-kafka-0-10_2.12:3.5.0 \
pyspark_etl_pipeline.py============================================================
π PYSPARK ETL PIPELINE - AMAZON SALES ANALYSIS
============================================================
β Spark Session Created: Amazon_Sales_ETL
β Spark Version: 3.5.0
============================================================
STEP 1: DATA INGESTION
============================================================
β Data loaded from: Cleaned_Amazon_Sale_Report.csv
β Total records: 128,976
β Total columns: 27
============================================================
STEP 2: DATA CLEANING & TRANSFORMATION
============================================================
π Initial Record Count: 128,976
π Missing Values Analysis:
- promotion-ids: 45,231 (35.07%)
π§ Cleaning Operations:
β Removed 234 duplicate orders
β Filled missing values
β Data type corrections applied
β Final cleaned records: 128,742
============================================================
STEP 3: FEATURE ENGINEERING
============================================================
β Added derived columns:
- Cost, Profit, Profit_Margin_Pct
- Price_Per_Item
- Is_B2B, Is_Amazon_Fulfilled, Has_Promotion
============================================================
STEP 4: KPI CALCULATIONS
============================================================
π KPI 1: Monthly Revenue
+----+-----+---------+-------------+-----------+----------------+
|Year|Month|MonthName|Total_Revenue|Order_Count|Total_Units_Sold|
+----+-----+---------+-------------+-----------+----------------+
|2022| 4| Apr| 15234567.89| 98,432| 145,678|
+----+-----+---------+-------------+-----------+----------------+
... [Additional KPI outputs]
============================================================
β
ETL PIPELINE COMPLETED SUCCESSFULLY!
============================================================
1. Memory Errors
# Increase driver/executor memory
spark-submit --driver-memory 8g --executor-memory 8g script.py2. File Not Found
# Check file path
import os
print(os.path.abspath("Cleaned_Amazon_Sale_Report.csv"))3. Java Version Issues
# Check Java version (needs 8 or 11)
java -version
# Set JAVA_HOME
export JAVA_HOME=/usr/lib/jvm/java-11-openjdk4. Slow Performance
# Reduce shuffle partitions for small datasets
.config("spark.sql.shuffle.partitions", "50")- β
PySpark ETL script (
pyspark_etl_pipeline.py) - β Pipeline architecture diagram
- β Complete documentation (this file)
- β Sample outputs in Parquet format
- β Console screenshots showing execution
- β KPI calculation results
- β Data quality report





