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
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
- Source: 60 CSV files (2020-2025) in
data/directory - Loader:
batch_loader.py- loads data with ON CONFLICT upsert - Target: PostgreSQL
fact_tripstable - Frequency: One-time load + periodic full reloads
- 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
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
- 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
- Views:
popular_stations: Top 20 start stations with trip countstrip_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
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
| 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 |
| 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 |
✅ Requirements engineering
✅ Exploratory data analysis (EDA)
✅ Batch database loading
✅ Real-time streaming (Kafka & Spark)
✅ Dashboard visualization (Streamlit)
- Docker & Docker Compose installed
- 8GB+ available RAM
- Python 3.10+ (for local development)
-
Clone the repository:
git clone https://github.com/AngeloOttendorfer02/DEMAI-divvi-trips.git cd DEMAI-divvi-trips -
Start the pipeline:
cd docker docker-compose up --build -dThis 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)
-
Access the dashboard:
- Streamlit: http://localhost:8501
- Adminer: http://localhost:8080 (divvy/divvy)
- Spark UI: http://localhost:4040
-
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;"
-
Stop the pipeline:
docker-compose down -v # Remove volumes too
- 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
- Primary Key:
-
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
- Primary Key: (
-
agg_daily_usage: Daily aggregations with weekday indicators
- Primary Key:
date - Metrics: member ratio, trip counts, average metrics, weekend flags
- Primary Key:
-
agg_station_popularity: Station-level aggregations (start stations)
- Primary Key:
station_name - Metrics: total starts/ends, average duration, member ratio
- Primary Key:
-
summary_rider_behavior: User type aggregations
- Primary Key:
member_casual - Metrics: trip counts, duration, distance by user type
- Primary Key:
- popular_stations: Top 20 start stations with trip counts
- trip_patterns_by_hour: Hourly trip patterns by user type with null-aware averages
- 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
- 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
- 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
- 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
- Spark Checkpoints:
/chk/directory stores state for stateful operations - Kafka Offsets: Tracked automatically, enables recovery from failures
- Benefit: Exactly-once semantics, no data loss
Cause: Stale Spark checkpoints referencing deleted Kafka messages
Solution:
docker-compose down -v
docker-compose up --build -dThis removes all volumes including old checkpoints.
Cause: Update mode queries emitting windows multiple times
Solution: Ensure upsert function uses ON CONFLICT DO UPDATE, not INSERT alone
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)
ENDCause: Tables not populated or connection closed
Solution: Keep cached connection alive, don't close after each query
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
- Create aggregation logic in
spark_stream.py - Add checkpoint location (e.g.,
/chk/new_query) - Define output table with appropriate primary keys
- Create table in PostgreSQL with ON CONFLICT constraints
- Add upsert call with correct key columns
- Create query function in
streamlit/app.py - Add new tab or chart using Plotly
- Cache results with
@st.cache_data - Handle empty data gracefully with
st.info()messages
- 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
- 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
- Batch Data: ~2.5M trips (2020-2025)
- Streaming Rate: 1.73M trips/day (at 20 msg/sec)
- Database Size: ~500MB (fact table + aggregations)
- Apache Kafka Documentation
- Apache Spark Streaming Guide
- PostgreSQL UPSERT
- Streamlit Documentation
- Divvy Bike Share Data
This project is part of the FH Master AI Data Engineering course.
Author: Angelo Ottendorfer
Repository: https://github.com/AngeloOttendorfer02/DEMAI-divvi-trips
Last Updated: January 2026