Skip to content

About

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.

Topics

Resources

Stars

1 star

Watchers

0 watching

Forks

Latest commit

 

History

2 Commits

Folders and files

NameName
Last commit message
Last commit date
 
 
 
 
 
 
 
 
 
 
 
 
 
 

Repository files navigation

Log Analytics ETL Pipeline

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.


Architecture Overview

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

Key Features

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)

Project Structure

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

Quickstart

1. Install dependencies

pip install -r requirements.txt

2. Generate sample data

cd data
python generate_sample_data.py
# → data/raw/sample_logs.csv (10,000 rows)

3. Run the pipeline

cd src
python pipeline.py \
  --input ../data/raw/sample_logs.csv \
  --output ../data/processed \
  --env local

4. Run tests

pytest tests/ -v --tb=short
# → 17 tests across Bronze, Silver, Gold, DQ layers

Running on Databricks

# 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:

  • Delta preview, Databricks IO caching (SSD)
  • Adaptive Query Execution + skew join handling

Running on AWS EMR

spark-submit \
  --master yarn \
  --deploy-mode cluster \
  src/pipeline.py \
  --input s3://your-bucket/raw/logs/ \
  --output s3://your-bucket/processed/ \
  --env emr

Data Quality Framework

Each layer runs a configurable set of weighted checks. The pipeline aborts if:

  • A critical=True check fails, OR
  • The weighted DQ score falls below 80%

Bronze Checks (6)

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

Silver Checks (5)

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

Gold Tables

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

Tech Stack

  • PySpark 3.5 — distributed transformation engine
  • Databricks / EMR / Synapse / HDInsight — cloud execution targets
  • Python 3.10+ — modular, typed codebase
  • pytest — 17-test unit suite
  • Parquet — columnar output format with partition pruning

About

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.

Topics

Resources

Stars

1 star

Watchers

0 watching

Forks

Releases

Packages

Contributors

Languages