A production-grade, distributed ETL pipeline for processing application log data at scale using PySpark and a Medallion (Bronze/Silver/Gold) architecture — compatible with Databricks, AWS EMR, Azure Synapse, and local Spark clusters.
Raw Logs (CSV/JSON/Parquet)
│
▼
┌─────────────┐
│ INGESTION │ Schema enforcement, multi-format reader, streaming watermark
└──────┬──────┘
│
▼
┌─────────────┐
│ BRONZE │ Type casting · Timestamp parsing (5 formats) · Partition derivation
└──────┬──────┘ Categorical normalisation · Pipeline audit columns
│
├──────► DQ Framework (6 checks, weighted score, critical gate)
│
▼
┌─────────────┐
│ SILVER │ Deduplication (Window) · Latency bucketing · Status code family
└──────┬──────┘ Error flags · Rolling 1h error rate · PII normalisation
│
├──────► DQ Framework (5 checks)
│
▼
┌─────────────┐
│ GOLD │ 5 aggregation tables written as Parquet
└─────────────┘ service_health_hourly · error_summary_daily · traffic_volume_hourly
sla_compliance_daily · top_slow_endpoints
| Feature | Detail |
|---|---|
| Multi-environment Spark | Local, Databricks, EMR, Synapse — one config factory |
| Medallion Architecture | Bronze / Silver / Gold with partition-aware writes |
| Data Quality Framework | 11 configurable checks, weighted scoring, critical gate (aborts pipeline) |
| Window Functions | Deduplication + rolling 1h error rate via PySpark Windows |
| Latency SLA Tracking | Bucketed p50/p95/p99 per service per hour |
| Multi-format Ingestion | CSV, JSON, Parquet + Structured Streaming watermark |
| Adaptive Query Execution | AQE + partition coalescing enabled for all environments |
| Observability | JSON metrics file per run with DQ scores and record counts |
| Test Suite | 17 pytest unit tests across all layers (Bronze, Silver, Gold, DQ) |
log_analytics_pipeline/
├── src/
│ ├── pipeline.py # Main orchestrator
│ ├── ingestion/
│ │ └── log_reader.py # Multi-format reader + streaming watermark
│ ├── transformation/
│ │ ├── bronze_layer.py # Schema enforcement, timestamp parsing
│ │ ├── silver_layer.py # Dedup, enrichment, window functions
│ │ └── gold_layer.py # 5 KPI aggregation tables
│ ├── quality/
│ │ └── dq_framework.py # Configurable DQ engine (11 checks)
│ └── utils/
│ ├── spark_utils.py # Environment-aware SparkSession factory
│ └── metrics_writer.py # JSON run metrics for observability
├── tests/
│ └── test_pipeline.py # 17 pytest unit tests
├── data/
│ └── generate_sample_data.py # Synthetic 10K-row log generator
├── dashboard/
│ └── pipeline_dashboard.html # Pipeline monitoring dashboard
├── requirements.txt
├── pytest.ini
└── README.md
pip install -r requirements.txt
cd data
python generate_sample_data.py
# → data/raw/sample_logs.csv (10,000 rows)
cd src
python pipeline.py \
--input ../data/raw/sample_logs.csv \
--output ../data/processed \
--env local
pytest tests/ -v --tb=short
# → 17 tests across Bronze, Silver, Gold, DQ layers
# In a Databricks notebook:
%run /path/to/src/pipeline
run_pipeline(
input_path = "dbfs:/mnt/raw/logs/",
output_path = "dbfs:/mnt/processed/",
env = "databricks"
)
The SparkSession factory automatically enables:
spark-submit \
--master yarn \
--deploy-mode cluster \
src/pipeline.py \
--input s3://your-bucket/raw/logs/ \
--output s3://your-bucket/processed/ \
--env emr
Each layer runs a configurable set of weighted checks. The pipeline aborts if:
critical=True check fails, OR| Check | Threshold | Weight | Critical | |—|—|—|—| | null_event_timestamp | 99.0% | 2.0 | ✓ | | null_service_name | 98.0% | 1.5 | | | null_log_level | 97.0% | 1.0 | | | null_host | 95.0% | 1.0 | | | valid_log_level | 95.0% | 1.0 | | | null_request_id | 90.0% | 0.8 | |
| Check | Threshold | Weight | Critical | |—|—|—|—| | null_event_timestamp | 100% | 2.0 | ✓ | | null_service_name | 99.0% | 1.5 | | | valid_status_code | 95.0% | 1.0 | | | valid_response_time | 95.0% | 1.0 | | | valid_bytes_sent | 97.0% | 0.8 | |
| Table | Description | Key Metrics |
|---|---|---|
service_health_hourly |
Per-service health per hour | error_rate_pct, p50/p95/p99_latency_ms |
error_summary_daily |
Error breakdown | by log_level, status_code_family, region |
traffic_volume_hourly |
Traffic volume | request_count, total_mb_sent, unique_users |
sla_compliance_daily |
SLA adherence | % requests per latency bucket per service |
top_slow_endpoints |
Ranked slow services | p95/p99 latency, slow_request_count |