Skip to content

About

No description, website, or topics provided.

Resources

Stars

0 stars

Watchers

0 watching

Forks

Repository files navigation

Divvy Real-Time Analytics Pipeline

Project Overview

This project implements a Lambda Architecture data pipeline for Divvy bike-sharing data, combining batch and real-time processing to enable near-real-time analytics and insights.

Data Source: https://divvy-tripdata.s3.amazonaws.com/index.html

Architecture Stack:

  • Message Broker: Apache Kafka 7.5.0 (topic: divvy_trips, 6 partitions, 720-hour retention)
  • Stream Processing: Apache Spark 3.5.1 (3 streaming queries with 1-hour watermark)
  • Database: PostgreSQL 16 (7 tables + 2 views)
  • Dashboard: Streamlit (5 analytical views with auto-refresh)
  • Orchestration: Docker Compose 3.9

Architecture Overview

graph LR
    CSV["📊 Historical CSV Data<br/>(2020-2025)<br/>60 files"] 
    
    CSV -->|Batch Load| BL["🔧 Batch Loader<br/>(Python)"]
    CSV -->|Stream Producer| KP["📤 Kafka Producer<br/>(20 msg/sec)<br/>Haversine Distance"]
    
    KP -->|Kafka Topic| KAFKA["🔄 Apache Kafka<br/>divvy_trips<br/>6 partitions"]
    
    KAFKA -->|Structured Streaming| SPARK["⚡ Apache Spark<br/>3 Streaming Queries<br/>1-hour Watermark"]
    
    BL --> PG["🗄️ PostgreSQL 16<br/>7 Tables + 2 Views<br/>UPSERT Pattern"]
    SPARK --> PG
    
    PG -->|Query| ST["📈 Streamlit Dashboard<br/>5 Analytical Tabs<br/>Auto-Refresh 30s"]
    
    style CSV fill:#e3f2fd
    style KP fill:#fff3e0
    style KAFKA fill:#ffe0b2
    style SPARK fill:#f3e5f5
    style PG fill:#e8f5e9
    style ST fill:#fce4ec
    style BL fill:#f1f8e9
Loading

Data Flow Pipeline

1. Historical Data (Batch Path)

  • Source: 60 CSV files (2020-2025) in data/ directory
  • Loader: batch_loader.py - loads data with ON CONFLICT upsert
  • Target: PostgreSQL fact_trips table
  • Frequency: One-time load + periodic full reloads

2. Real-Time Data (Streaming Path)

  • Source: CSV data simulated at 20 msg/sec by producer.py
  • Enrichment: Haversine distance calculation, temporal features (hour, day_of_week)
  • Message Format: 15 fields (ride_id, timestamps, coordinates, station names, user types, bike types)
  • Broker: Kafka with 6 partitions for parallel consumption

3. Stream Processing (Spark)

Three parallel streaming queries process the Kafka stream:

Query Window Output Mode Target Table Purpose
fact_trips Event APPEND fact_trips Store all raw trip events
agg_hourly_kpis 1-hour, 30-min overlap UPDATE agg_hourly_kpis Windowed metrics with rolling window
agg_daily_usage Daily UPDATE agg_daily_usage Daily aggregations with weekday flags

Type Safety: All coordinate columns explicitly cast to DOUBLE PRECISION after Kafka JSON parsing

4. Data Persistence

  • Pattern: UPSERT (INSERT ... ON CONFLICT DO UPDATE)
  • Deduplication: Primary keys on window_start+window_end or date
  • Null Handling: Explicit filtering and CASE statements in aggregations

5. Analytics & Visualization

  • Views:
    • popular_stations: Top 20 start stations with trip counts
    • trip_patterns_by_hour: Hourly patterns by user type with null-aware averages
  • Dashboard: 5 Streamlit tabs showing daily trends, hourly patterns, user types, bike types, and station analytics

Stakeholders

The following stakeholders are considered:

  • Data Analyst: Explores historical usage patterns via Streamlit dashboard
  • Operations Manager: Monitors bike usage trends in near real-time
  • Product Manager: Evaluates member vs casual usage ratios and preferences
  • System Owner: Ensures pipeline reliability, scalability, and reproducibility

Functional Requirements

Requirement Status Component
FR1: Load historical Divvy data from CSV files ✅ batch_loader.py
FR2: Preprocess timestamps, compute trip durations & distances ✅ Producer, Spark
FR3: Provide descriptive statistics of bike trips ✅ Jupyter notebooks
FR4: Visualize usage patterns and distributions ✅ Streamlit Dashboard
FR5: Aggregate trips by user type (member vs casual) ✅ Spark queries
FR6: Stream real-time bike trip data via Kafka ✅ producer.py
FR7: Process streaming data with Apache Spark ✅ spark_stream.py
FR8: Persist results in PostgreSQL database ✅ Spark JDBC writer
FR9: Expose analytics via interactive dashboard ✅ Streamlit app

Non-Functional Requirements

Requirement Target Implementation
NFR1 (Performance) Streaming analytics update within seconds 30-second Streamlit auto-refresh, Spark microbatches
NFR2 (Scalability) Support increased data volumes 6 Kafka partitions, Spark parallel processing
NFR3 (Reliability) No data loss during processing Spark checkpoints, Kafka 720-hour retention, UPSERT pattern
NFR4 (Maintainability) Modular, loosely coupled components Docker containerization, separated producer/spark/loader
NFR5 (Reproducibility) Fully reproducible environment Docker Compose orchestration, versioned images

Current Project Stage

✅ Requirements engineering
✅ Exploratory data analysis (EDA)
✅ Batch database loading
✅ Real-time streaming (Kafka & Spark)
✅ Dashboard visualization (Streamlit)


Quick Start Guide

Prerequisites

  • Docker & Docker Compose installed
  • 8GB+ available RAM
  • Python 3.10+ (for local development)

Setup & Run

  1. Clone the repository:

    git clone https://github.com/AngeloOttendorfer02/DEMAI-divvi-trips.git
    cd DEMAI-divvi-trips
  2. Start the pipeline:

    cd docker
    docker-compose up --build -d

    This starts all 10 services:

    • Zookeeper, Kafka (message broker)
    • PostgreSQL (data warehouse)
    • Spark (stream processor)
    • Producer (data generator)
    • Batch Loader (historical data)
    • Streamlit (dashboard)
    • Adminer (database UI)
  3. Access the dashboard:

  4. Monitor data flow:

    # Check Spark logs
    docker logs docker-spark-1 -f
    
    # Check Kafka producer
    docker logs docker-producer-1 -f
    
    # Check database
    docker exec docker-postgres-1 psql -U divvy -d divvy -c "SELECT COUNT(*) FROM fact_trips;"
  5. Stop the pipeline:

    docker-compose down -v  # Remove volumes too

Database Schema

Fact Tables

  • fact_trips: Raw trip events (15 columns, ~2.5M records after 2025 data)
    • Primary Key: ride_id
    • Includes: timestamps, coordinates, station names, user type, bike type

Aggregation Tables

  • agg_hourly_kpis: 1-hour windowed metrics with 30-minute overlap

    • Primary Key: (window_start, window_end)
    • Metrics: trip counts by user/bike type, average duration/distance
  • agg_daily_usage: Daily aggregations with weekday indicators

    • Primary Key: date
    • Metrics: member ratio, trip counts, average metrics, weekend flags
  • agg_station_popularity: Station-level aggregations (start stations)

    • Primary Key: station_name
    • Metrics: total starts/ends, average duration, member ratio
  • summary_rider_behavior: User type aggregations

    • Primary Key: member_casual
    • Metrics: trip counts, duration, distance by user type

Views

  • popular_stations: Top 20 start stations with trip counts
  • trip_patterns_by_hour: Hourly trip patterns by user type with null-aware averages

Key Technical Decisions

1. UPSERT Pattern for Aggregations

  • Why: Spark update mode queries emit windows multiple times as data arrives
  • Solution: PostgreSQL ON CONFLICT DO UPDATE clause instead of DELETE+INSERT
  • Benefit: No temporary tables, atomic updates, data consistency

2. Explicit Type Casting in Spark

  • Why: Kafka JSON parsing converts all fields to strings
  • Solution: Cast coordinates to DOUBLE PRECISION immediately after JSON parsing
  • Benefit: Prevents database type errors, ensures numeric operations work correctly

3. Null-Aware Aggregations in SQL

  • Why: NULL trip distances (when coordinates missing) produce NaN in averages
  • Solution: CASE statements that only aggregate non-NULL values
  • Benefit: Accurate metrics, no NaN values in views

4. Lambda Architecture (Batch + Streaming)

  • Batch Path: Historical data via batch_loader.py (one-time load)
  • Streaming Path: Real-time data via Kafka producer (continuous)
  • Benefit: Handles both historical reprocessing and real-time updates

5. Checkpoint-Based Fault Tolerance

  • Spark Checkpoints: /chk/ directory stores state for stateful operations
  • Kafka Offsets: Tracked automatically, enables recovery from failures
  • Benefit: Exactly-once semantics, no data loss

Troubleshooting

Issue: OffsetOutOfRangeException

Cause: Stale Spark checkpoints referencing deleted Kafka messages
Solution:

docker-compose down -v
docker-compose up --build -d

This removes all volumes including old checkpoints.

Issue: Duplicate key errors on aggregation tables

Cause: Update mode queries emitting windows multiple times
Solution: Ensure upsert function uses ON CONFLICT DO UPDATE, not INSERT alone

Issue: NaN values in view averages

Cause: Aggregating NULL values without filtering
Solution: Use CASE statements with conditional NULL filtering:

CASE 
  WHEN COUNT(CASE WHEN trip_distance_km IS NOT NULL THEN 1 END) = 0 THEN NULL
  ELSE AVG(CASE WHEN trip_distance_km IS NOT NULL THEN trip_distance_km END)
END

Issue: Streamlit shows "No data available yet"

Cause: Tables not populated or connection closed
Solution: Keep cached connection alive, don't close after each query


Project Structure

DEMAI-divvi-trips/
├── data/                          # CSV data files (60 files, 2020-2025)
├── notebooks/                     # Jupyter notebooks (EDA analysis)
├── docker/                        # Docker orchestration
│   ├── docker-compose.yml         # Service definitions (10 services)
│   └── Dockerfile                 # Custom Spark/Kafka images
├── spark/                         # Spark streaming job
│   ├── spark_stream.py            # 3 streaming queries + upsert logic
│   └── requirements.txt           # Python dependencies
├── kafka/                         # Kafka producer
│   ├── producer.py                # Real-time data generation with Haversine distance
│   └── requirements.txt
├── batch_loader.py                # Historical data loading script
├── streamlit/                     # Dashboard application
│   ├── app.py                     # 5-tab Streamlit dashboard
│   └── requirements.txt
├── sql/                           # Database schema
│   └── init.sql                   # 7 tables + 2 views definition
├── README.md                      # This file
└── PIPELINE_DOCUMENTATION.md      # Detailed technical documentation

Development & Extension

Adding a New Streaming Query

  1. Create aggregation logic in spark_stream.py
  2. Add checkpoint location (e.g., /chk/new_query)
  3. Define output table with appropriate primary keys
  4. Create table in PostgreSQL with ON CONFLICT constraints
  5. Add upsert call with correct key columns

Adding a New Dashboard View

  1. Create query function in streamlit/app.py
  2. Add new tab or chart using Plotly
  3. Cache results with @st.cache_data
  4. Handle empty data gracefully with st.info() messages

Scaling the Pipeline

  • Increase Kafka partitions: Modify docker-compose.yml
  • Increase Spark parallelism: Set spark.sql.shuffle.partitions
  • Add more storage: Scale PostgreSQL with read replicas
  • Optimize queries: Add indexes on commonly filtered columns

Performance Metrics

Current Throughput

  • Producer: 20 messages/second (configurable)
  • Spark Batch Interval: Microbatches every few seconds
  • Kafka Retention: 720 hours (30 days) with 10GB size limit
  • Dashboard Refresh: 30-second auto-refresh

Data Volume

  • Batch Data: ~2.5M trips (2020-2025)
  • Streaming Rate: 1.73M trips/day (at 20 msg/sec)
  • Database Size: ~500MB (fact table + aggregations)

References & Resources


License

This project is part of the FH Master AI Data Engineering course.


Contact & Attribution

Author: Angelo Ottendorfer
Repository: https://github.com/AngeloOttendorfer02/DEMAI-divvi-trips
Last Updated: January 2026

About

No description, website, or topics provided.

Resources

Stars

0 stars

Watchers

0 watching

Forks

Releases

Packages

Contributors

Languages