End-to-end traffic analytics pipeline built with Kafka, Spark Structured Streaming, Delta Lake, dbt, and a Streamlit dashboard — from raw streaming ingestion through a tested dimensional warehouse to an operational analytics UI.
Urban traffic data is noisy, fast-moving, and difficult to use directly for analytics. Raw streaming events are not enough for decision-making because they need validation, structure, and business-friendly models before they can support monitoring or dashboards.
I built this project to solve that problem end to end:
- ingest live traffic events in real time
- process them through a Lakehouse architecture
- transform them into clean analytical tables
- validate data quality automatically with dbt tests
- expose them in a dashboard that is easy to explore
The goal was to show how raw streaming data can become something useful — and trustworthy — for operational traffic analysis.
This project simulates live traffic data, processes it through Bronze, Silver, and Gold layers, builds a tested dimensional warehouse with dbt, and serves a dashboard for operational monitoring and analytics.
Pipeline flow:
Kafka -> Bronze Delta -> Silver Delta -> Gold Delta -> Postgres staging -> dbt warehouse -> Streamlit Dashboard
- Real-time traffic ingestion with Kafka
- Bronze layer for raw streaming capture
- Silver layer for validation, typing, deduplication, and feature engineering
- Gold layer for star-schema style analytics tables (Delta Lake)
- Warehouse layer with dbt:
- dimensional model:
dim_locationsandfct_traffic_events - 28 data tests: uniqueness, not_null, accepted_values, foreign key relationships
- GDPR PII classification meta tags on every column
- Elementary observability for test monitoring and reporting
- dimensional model:
- Streamlit dashboard reading from the tested dbt warehouse:
- KPI overview
- zone activity map
- congestion and road analysis
- weather and speed insights
- zone and road drill-down page
- CI/CD with GitHub Actions:
- dbt tests + Elementary report on every pull request
- Terraform plan validation on infrastructure changes
- Infrastructure as Code with Terraform:
- S3 bucket for CI artifacts
- IAM roles with GitHub Actions OIDC (no static credentials)
apps/ Spark jobs for bronze, silver, and gold layers
dashboard/ Streamlit app and dashboard pages
dbt/ dbt project: models, tests, seeds, Elementary
models/marts/ Dimensional models (dim_locations, fct_traffic_events)
models/views/ Compatibility views for the dashboard
models/marts/schema.yml Tests and GDPR PII meta tags
seeds/ CI fixture data for automated testing
packages.yml Elementary dbt package
infra/terraform/ Terraform: S3 bucket + IAM OIDC for GitHub Actions
scripts/ Pipeline runner, loader, and Elementary report scripts
.github/workflows/ CI: dbt-ci.yml and terraform-plan.yml
docs/images/ Screenshots used in the GitHub README
hive-conf/ Hive metastore configuration
producer/ Kafka traffic producer
sql/ SQL setup scripts
docker-compose.yaml Local infrastructure: Spark (2 workers), Kafka, Hive, Postgres x2
The project follows a streaming Lakehouse architecture from ingestion to a tested warehouse and analytics layer.
flowchart LR
subgraph INGESTION["Ingestion"]
P["Producer\nPython + Faker"] --> K["Kafka\nKRaft mode"]
end
subgraph LAKEHOUSE["Lakehouse · Medallion Architecture"]
direction LR
B["Bronze\nRaw ingestion"] --> S["Silver\nValidation · Dedup\nFeature engineering"]
S --> G["Gold\nStar schema\ndim + fact tables"]
end
subgraph WAREHOUSE["Warehouse"]
direction TB
PG["PostgreSQL 16\nStaging tables"]
DBT["dbt Core\ndim_locations\nfct_traffic_events\n28 tests · GDPR tags"]
EL["Elementary\nObservability\nHTML reports"]
PG --> DBT
DBT --> EL
end
subgraph SERVING["Serving"]
ST["Streamlit Dashboard\nKPIs · Maps · Charts\nPlotly · PyDeck"]
end
subgraph CICD["CI / CD & Infrastructure"]
direction LR
GHA["GitHub Actions"]
TF["Terraform"]
AWS["AWS S3 + IAM\nOIDC"]
end
K --> B
G -->|"Loader\nDelta → Postgres"| PG
DBT --> ST
style INGESTION fill:#fef3c7,stroke:#f59e0b,color:#92400e
style LAKEHOUSE fill:#eff6ff,stroke:#93c5fd,color:#1e40af
style WAREHOUSE fill:#ecfdf5,stroke:#34d399,color:#065f46
style SERVING fill:#fef2f2,stroke:#ef4444,color:#991b1b
style CICD fill:#f0f9ff,stroke:#38bdf8,color:#0c4a6e
style P fill:#fef9c3,stroke:#ca8a04,color:#854d0e
style K fill:#f3e8ff,stroke:#a78bfa,color:#5b21b6
style B fill:#fef9c3,stroke:#ca8a04,color:#854d0e
style S fill:#e0e7ff,stroke:#6366f1,color:#3730a3
style G fill:#fef3c7,stroke:#d97706,color:#92400e
style PG fill:#ecfdf5,stroke:#34d399,color:#065f46
style DBT fill:#fff1f2,stroke:#f43f5e,color:#9f1239
style EL fill:#f0f9ff,stroke:#38bdf8,color:#0c4a6e
style ST fill:#fef2f2,stroke:#ef4444,color:#991b1b
style GHA fill:#dbeafe,stroke:#93c5fd,color:#1e40af
style TF fill:#dbeafe,stroke:#93c5fd,color:#1e40af
style AWS fill:#dbeafe,stroke:#93c5fd,color:#1e40af
Infrastructure: Docker Compose — Spark Master + 2 Workers (4 cores / 4 GB each) · Kafka (KRaft) · Hive Metastore · PostgreSQL x2 · Kafka UI
producer/traffic_data_producer.pysends simulated traffic events to Kafka.apps/traffic_bronze.pyconsumes Kafka and stores raw Delta data.apps/traffic_silver.pycleans and enriches Bronze data into Silver.apps/traffic_gold.pybuilds Gold analytical tables:fact_traffic,dim_zone,dim_road.scripts/load_gold_delta_to_psql.pyloads Gold Delta tables into Postgres staging tables.dbt runbuilds the dimensional warehouse:dim_locations,fct_traffic_events, and compatibility views.dbt testvalidates all 28 data quality tests.dashboard/app.pyreads the tested dbt warehouse tables and serves the UI.
After Gold tables land in Delta, dbt builds a proper dimensional model on top of them in Postgres.
The Gold Delta tables are good for the streaming pipeline, but they have no built-in way to validate the data. If a bad batch lands, nothing catches it. dbt adds a layer of automated testing so the data that reaches the dashboard is always validated.
dim_locations— one row per unique (city_zone, road_id) pair with zone and road attributes joined in. 20 rows.fct_traffic_events— one row per traffic event with alocation_idforeign key. 512 rows in the latest run.- Compatibility views (
fact_traffic,dim_zone,dim_road) so the dashboard reads from the warehouse without breaking its existing column contract.
- Uniqueness:
location_id,event_id - Not null: every column in both models
- Accepted values:
city_zone,road_id,traffic_risk,weather,speed_band - Relationships:
fct_traffic_events.location_idreferencesdim_locations.location_id
Every column has a gdpr_pii_classification meta tag. The only column marked as PII is vehicle_id (pseudonymous identifier). Elementary uses these tags to automatically suppress sample data from test failure reports for PII columns.
Elementary runs on top of dbt test results and generates an HTML observability report. In CI, this report is uploaded as a GitHub Actions artifact on every pull request.
The final result is a working local traffic analytics platform with:
- streaming ingestion from Kafka
- Lakehouse processing across Bronze, Silver, and Gold
- a tested dimensional warehouse built with dbt
- 28 automated data quality tests passing on every run
- GDPR-aware column metadata
- a Streamlit dashboard reading from the validated warehouse
- CI/CD that validates every change before merge
- infrastructure provisioned with Terraform
The dashboard is built with Streamlit and reads from the dbt warehouse tables in Postgres.
Overview: KPI cards, zone map, hourly pulse, top roads, weather impactZone and Road Explorer: zone drill-down, road pressure, heatmap, leaderboard
docker compose up -d
CLEAN=1 ./scripts/run_pipeline.shThis runs the full pipeline end to end: producer, Bronze, Silver, Gold, Delta-to-Postgres loader, dbt build + test, and launches the Streamlit dashboard at http://localhost:8501.
The script automatically polls each streaming layer for committed data and kills it before moving to the next step.
Click to expand individual steps
1. Start infrastructure
docker compose up -dThis starts Kafka (KRaft), Spark (master + 2 workers at 4 cores / 4 GB each), Hive metastore, Kafka UI, and two Postgres databases (Hive metadata + dbt warehouse).
2. Start Bronze stream, then produce data
# Terminal 1 — start Bronze (reads "latest" from Kafka, so start this first)
docker exec spark-master /opt/spark/bin/spark-submit \
--master spark://spark-master:7077 \
--packages io.delta:delta-spark_2.12:3.2.0,org.apache.spark:spark-sql-kafka-0-10_2.12:3.5.1 \
--conf spark.jars.ivy=/tmp/.ivy \
--conf spark.sql.extensions=io.delta.sql.DeltaSparkSessionExtension \
--conf spark.sql.catalog.spark_catalog=org.apache.spark.sql.delta.catalog.DeltaCatalog \
/opt/spark-apps/traffic_bronze.py
# Terminal 2 — produce events (Ctrl+C to stop after ~30s)
python producer/traffic_data_producer.pyKill the Bronze stream once warehouse/traffic_bronze/_delta_log/ appears.
3. Run Silver
docker exec spark-master /opt/spark/bin/spark-submit \
--master spark://spark-master:7077 \
--packages io.delta:delta-spark_2.12:3.2.0 \
--conf spark.jars.ivy=/tmp/.ivy \
--conf spark.sql.extensions=io.delta.sql.DeltaSparkSessionExtension \
--conf spark.sql.catalog.spark_catalog=org.apache.spark.sql.delta.catalog.DeltaCatalog \
/opt/spark-apps/traffic_silver.pyKill once warehouse/traffic_silver/_delta_log/ appears.
4. Run Gold (auto-terminates)
docker exec spark-master /opt/spark/bin/spark-submit \
--master spark://spark-master:7077 \
--packages io.delta:delta-spark_2.12:3.2.0 \
--conf spark.jars.ivy=/tmp/.ivy \
--conf spark.sql.extensions=io.delta.sql.DeltaSparkSessionExtension \
--conf spark.sql.catalog.spark_catalog=org.apache.spark.sql.delta.catalog.DeltaCatalog \
/opt/spark-apps/traffic_gold.py5. Load Gold into Postgres
python scripts/load_gold_delta_to_psql.py6. Run dbt
cd dbt && dbt deps && dbt run && dbt test7. Run dashboard
streamlit run dashboard/app.pyOpen http://localhost:8501.
./scripts/run_elementary_report.shThis generates an HTML observability report at edr_target/report.html.
GitHub Actions automatically:
- Starts a Postgres service container
- Seeds fixture data (no Spark needed in CI)
- Runs
dbt runanddbt test - Generates an Elementary report and uploads it as a PR artifact
When files under infra/terraform/ change, GitHub Actions runs terraform fmt, validate, and plan using OIDC-based AWS authentication.
The nexus-dbt-artifacts bucket in eu-central-1 (Frankfurt) is provisioned by Terraform and stores Elementary observability reports uploaded by CI. Public access is blocked and versioning is enabled.
Delta (streaming pipeline):
warehouse/traffic_bronzewarehouse/traffic_silverwarehouse/fact_trafficwarehouse/dim_zonewarehouse/dim_road
Postgres (dbt warehouse):
analytics.dim_locationsanalytics.fct_traffic_eventsanalytics.fact_traffic(view)analytics.dim_zone(view)analytics.dim_road(view)
- Add real geographic coordinates instead of representative zone centroids
- Support near real-time dashboard refresh
- Add anomaly detection tests with Elementary (volume, freshness)
- Add a DataHub catalogue for automated lineage visualization
- Deploy to AWS with Terraform (MSK for Kafka, EMR Serverless for Spark)
| Layer | Technology |
|---|---|
| Ingestion | Apache Kafka (KRaft mode), Python + Faker |
| Stream Processing | Apache Spark 3.5.1 Structured Streaming (1 master, 2 workers — 4c/4g each) |
| Storage | Delta Lake 3.2, Hive Metastore |
| Warehouse | PostgreSQL 16, dbt Core |
| Data Quality | dbt tests (28 total), GDPR PII meta tags |
| Observability | Elementary |
| Infrastructure | Terraform (AWS S3 + IAM OIDC) |
| CI/CD | GitHub Actions |
| Dashboard | Streamlit, Plotly, PyDeck |
| Orchestration | Docker Compose |
| Languages | Python, SQL |
- Generated data in
warehouse/is intentionally ignored from Git. - Connection env vars for the dbt Postgres are documented in
.env.example. - Run the full pipeline before opening the dashboard:
CLEAN=1 ./scripts/run_pipeline.sh



