This project implements a complete end-to-end data engineering solution for an e-commerce retail platform. It demonstrates modern big data technologies including Apache Spark, Kafka, MinIO (S3-compatible storage), and Airflow in a production-ready architecture following the Medallion pattern (Bronze → Silver → Gold).
- Domain: Multi-channel e-commerce platform (website, mobile app, marketplace partnerships)
- Challenge: Optimize marketing spend across channels based on customer segments
- ML Use Case: Product recommendation engine for personalized customer experiences
- Data Quality: Hybrid validation with custom checks and Great Expectations framework
| Component | Technology | Purpose |
|---|---|---|
| Storage | MinIO | S3-compatible object storage |
| Processing | Apache Spark | Batch and streaming data processing |
| Streaming | Apache Kafka | Real-time data ingestion and delay monitoring |
| Orchestration | Apache Airflow | Workflow scheduling and monitoring |
| Containerization | Docker Compose | Service orchestration |
| Data Quality | Great Expectations | Data validation and quality assurance |
Data Sources → Kafka → Bronze Layer → Silver Layer → Gold Layer → Analytics/ML
↓ ↓ ↓ ↓ ↓
Streaming Raw Data Cleaned Enriched Business
Events Landing & Valid & Historic Ready
ecommerce-data-platform/
├── README.md # This file
├── KAFKA_DELAY_IMPLEMENTATION.md # Kafka delay monitoring details
├── GREAT_EXPECTATIONS_INTEGRATION.md # Data quality framework details
├── docker-compose.yml # Main orchestration
├── start_platform_demo.sh # Main startup script
├── run_complete_demo.py # Complete demo orchestrator
├── create_sample_data.py # Sample data generation
├── visualize_tables.py # Data visualization
├── orchestration/ # Airflow components
│ ├── docker-compose.yml
│ ├── Dockerfile
│ ├── dags/
│ │ ├── bronze_to_silver_dag.py
│ │ ├── silver_to_gold_dag.py
│ │ ├── ml_feature_generation_dag.py
│ │ ├── data_quality_dag.py
│ │ └── great_expectations_dag.py
├── streaming/ # Kafka and producers
│ ├── docker-compose.yml
│ ├── kafka_delay_monitor.py # Kafka delay monitoring
│ ├── producers/
│ ├── consumers/
│ └── kafka/
├── processing/ # Spark applications
│ ├── docker-compose.yml
│ ├── Dockerfile
│ ├── spark-apps/
│ ├── great_expectations/ # Great Expectations configuration
│ └── utils/
├── sample_data/ # Sample datasets
│ ├── campaigns.json
│ ├── customers.json
│ ├── marketplace_sales.json
│ ├── products.json
│ └── user_events.json
├── storage/ # Data storage
│ ├── data/
│ └── minio/
└── visualizations/ # Generated visualizations
├── campaign_effectiveness_analysis.png
├── customer_segmentation_analysis.png
├── product_performance_analysis.png
├── sales_performance_dashboard.png
├── user_activity_analysis.png
└── interactive_dashboard.html
- Docker: Version 20.10 or higher
- Docker Compose: Version 2.0 or higher
- Memory: Minimum 16GB
- Disk Space: Minimum 16GB free space
- Operating System: Linux, macOS
- CPU: M2 or higher
-
Install Docker Desktop
- Download from docker.com
- Ensure Docker Compose is included
-
Verify Installation
docker --version docker-compose --version
# Clone the repository
git clone <repository-url>
cd ecommerce-data-platform
# Run setup script
./setup.sh
# Start the platform
./start_platform_demo.shThe start_platform_demo.sh script performs the following steps:
- Starts all Docker services
- Waits for services to be ready
- Generates sample data
- Sets up the Kafka delay monitoring
- Runs the platform demo
# Run the complete demo to see all platform capabilities
python run_complete_demo.pyThis demonstrates:
- Data ingestion into Bronze layer
- Transformation to Silver layer
- Aggregation to Gold layer
- Data quality validation
- ML feature generation
- Visualization examples
Once all services are running, access them via:
| Service | URL | Credentials |
|---|---|---|
| Airflow UI | http://localhost:8080 | admin/admin |
| Spark Master UI | http://localhost:8081 | - |
| MinIO Console | http://localhost:9001 | minio/minio123 |
Data Sources → Kafka → Bronze Layer → Silver Layer → Gold Layer → Analytics/ML
↓ ↓ ↓ ↓ ↓
Streaming Raw Data Cleaned Enriched Business
Events Landing & Valid & Historic Ready
- Purpose: Store raw data exactly as received
- Implementation:
- Minimal transformations (parsing, formatting)
- Preserves original data with added metadata
- Combines streaming and batch data sources
- Tables:
raw_user_events- Streaming user activity from Kafkaraw_marketplace_sales- Batch sales data (late arrivals up to 48h)raw_product_catalog- Product informationraw_customer_data- Customer profilesraw_marketing_campaigns- Campaign definitions
- Purpose: Clean, validate, and standardize data
- Implementation:
- Implements data quality checks
- Standardizes schema and data types
- Applies basic transformations and enrichments
- Tables:
standardized_user_events- Cleaned user activitystandardized_sales- Validated sales transactionsdim_customer- Customer dimensiondim_product- Product dimensiondim_marketing_campaign- Marketing campaignsdim_date- date format and info
- Purpose: Create business-ready aggregations and features
- Implementation:
- Builds domain-specific views and metrics
- Creates aggregations for business intelligence
- Generates features for machine learning
- Tables:
fact_sales- Sales fact tablefact_user_activity- User activity fact tablesales_performance_metrics- Channel performance metricscustomer_segmentation- Customer segmentscampaign_effectiveness- Marketing ROI analysiscustomer_product_interactions- Machine learning feature tables
- Implementation: Kafka + Spark Streaming
- Features:
- Real-time data ingestion from various sources
- Event-time processing with watermarking
- Kafka delay monitoring for SLA tracking
- Late data handling (up to 48 hours)
- Continuous processing to Bronze and Silver layers
- Implementation: Airflow + Spark
- Features:
- Scheduled ETL workflows
- Incremental and full processing modes
- Complex transformations for Silver and Gold layers
- Dependency management between processing stages
- Data quality validation integration
- Implementation: Airflow DAG + Spark
- Features:
- Automated feature extraction from Silver and Gold layers
- Customer behavior feature generation
- Product performance metrics
- Time-based features for recommendation
- Feature versioning and lineage tracking
The platform implements a hybrid data quality approach:
- Integrated into processing pipelines
- Schema validation and enforcement
- Business rule validation
- Completeness and consistency checks
- Automated monitoring and alerting
- Declarative data quality definitions
- Expectations defined for each data asset
- Scheduled validation with Airflow
- Detailed quality documentation
- Validation results and reporting
- See GREAT_EXPECTATIONS_INTEGRATION.md for details
- Real-time tracking of message delivery delays
- SLA monitoring for streaming data
- Configurable alerting thresholds
- Custom implementation for production monitoring
- See KAFKA_DELAY_IMPLEMENTATION.md for details
- Marketing channel optimization
- Customer segmentation analysis
- Campaign effectiveness tracking
- Sales performance metrics
- User activity patterns
- Customer-product interaction features
- Behavioral features (recency, frequency, monetary)
- Category affinity calculations
- Feature preparation for recommendation engine
- ML feature versioning
- Message broker for real-time data streams
- Topics organized by data domain
- Custom delay monitoring implementation
- Producer and consumer implementations
- Distributed data processing engine
- Supports both batch and streaming workloads
- Executes transformations between data layers
- Preconfigured for optimal performance
- Workflow orchestration platform
- DAGs for each major processing pipeline:
- Bronze to Silver ETL
- Silver to Gold ETL
- ML Feature Generation
- Data Quality Monitoring
- Great Expectations Validation
- S3-compatible object storage
- Stores data in all processing stages
- Organized by data layer (Bronze/Silver/Gold)
- Accessible via S3 API
- Container orchestration for all components
- Isolated environments for each service
- Simplified deployment and configuration
- Organized by functional area
# Check Docker resource allocation (memory, CPU)
docker stats
# Check logs for specific services
docker-compose logs -f <service_name>
# Restart specific service
docker-compose restart <service_name># Check Airflow logs
ls -la orchestration/logs/dag_id=*/
# Restart Airflow services
docker-compose -f orchestration/docker-compose.yml restart# Check Kafka logs
ls -la logs/kafka/
# Verify Kafka service is running
docker-compose -f streaming/docker-compose.yml ps# Check Spark logs
ls -la logs/spark/
# Access Spark UI for job details
open http://localhost:8081# Run the complete demo again to verify data flow
python run_complete_demo.py
# Check specific layer processing
python silver_transformations.py
python gold_aggregations.py# Generate sample data
python create_sample_data.py
# Generate current data
python generate_current_data.py# Run Silver layer transformations
python silver_transformations.py
# Run Gold layer aggregations
python gold_aggregations.py# Generate visualizations
python visualize_tables.py# Run data quality validation
python run_great_expectations_demo.pyThe start_platform_demo.sh script orchestrates the platform startup:
-
Service Initialization
- Starts all Docker services using multiple docker-compose files
- Sets up networking between components
- Initializes storage volumes
-
Service Readiness Check
- Waits for Kafka to be ready
- Waits for Spark to be ready
- Waits for Airflow to be ready
- Waits for MinIO to be ready
-
Data Preparation
- Generates sample data using
create_sample_data.py - Sets up Kafka topics and configurations
- Generates sample data using
-
Demo Preparation
- Sets up Kafka delay monitoring
- Prepares demo environment
- Initializes logging
The run_complete_demo.py script demonstrates the full platform capabilities:
-
Bronze Layer Processing
- Ingests raw data from sample sources
- Streams events through Kafka
- Performs minimal transformations
- Stores data in Bronze tables
-
Silver Layer Processing
- Reads data from Bronze layer
- Applies cleaning and standardization
- Validates data quality
- Stores results in Silver tables
-
Gold Layer Processing
- Reads data from Silver layer
- Creates business aggregations
- Generates ML features
- Stores results in Gold tables
-
Visualization Generation
- Creates analytical visualizations
- Generates performance dashboards
- Outputs results to visualizations directory
- Increase Spark executors for larger data volumes
- Add Kafka partitions for higher throughput
- Implement distributed MinIO for storage scaling
- Configure Airflow for parallel task execution
- Enable authentication for all services
- Implement encryption for data in transit and at rest
- Use proper network segmentation
- Implement role-based access control
- Deploy multiple instances of each service
- Configure proper replication factors
- Implement automated failover
- Set up comprehensive monitoring and alerting
For more detailed information about specific components, refer to:
- PROJECT_SUMMARY.md: Complete project overview and architecture
- KAFKA_DELAY_IMPLEMENTATION.md: Kafka delay monitoring details
- GREAT_EXPECTATIONS_INTEGRATION.md: Data quality framework details
- docs/data-model.md: Detailed data model documentation
- docs/setup-guide.md: Comprehensive setup instructions
- docs/user-guide.md: Platform usage documentation
# Stop all services
docker-compose down
# Stop specific component services
docker-compose -f streaming/docker-compose.yml down
docker-compose -f processing/docker-compose.yml down
docker-compose -f orchestration/docker-compose.yml down# Complete reset (WARNING: deletes all data)
docker-compose down -v
docker system prune -f
rm -rf logs/* storage/data/* spark-warehouse/*- Advanced Analytics: Enhanced dashboards and reporting
- Machine Learning: Production deployment of recommendation models
- Real-time Alerting: Enhanced monitoring and notification system
- Data Lineage: Implementation of data lineage tracking
- Extended Data Quality: Additional validation rules and checks
This implementation provides a solid foundation for a production-ready data engineering platform that can scale with business needs while maintaining data quality and operational excellence.