Data eng
NYC Taxi Lakehouse: Batch and Streaming Pipeline
NYC TLC trip records
Parquet9.5M rowsJan to Mar 2024
Batch path
Accuracy · Airflow DAG, daily
Bronze

PySparkParquet → Delta
Raw rows plus ingestion time and source file
Quality gaterows · nulls · schema · ranges
Silver

PySparkDelta Lake
Typed, de-duplicated, outliers removed, enriched
Quality gate
Gold

Window functionsDelta + CSV
Five aggregate tables, partitioned by date
Quality gate
Data lake and SQL
S3Glue CatalogAthena
Provisioned with Terraform
Out: Gold tables queryable in SQL through Athena
Streaming path
Freshness · results within minutes
Producer
PythonJSON events
Replays trips at about 97 events a second
Kafka topic
3 partitions
Keyed by pickup zone to keep order
Stream processing
Structured Streaming
5-minute tumbling windows, 10-minute watermark, backpressure cap
Streaming tables

Delta Lake
Checkpointed, so a restart resumes where it stopped
Out: trips and revenue per borough, every 5 minutes
The problem
Raw event data is only useful once people can trust it and reach it. Records arrive in bulk and in real time, with duplicates, gaps and impossible values mixed in, and the people who depend on them want two things at once: a complete, correct history and a view of what is happening right now. This project builds both on 9.5M New York taxi trips, using PySpark and Delta Lake for a Bronze, Silver and Gold batch path and Kafka with Spark Structured Streaming for a real-time path, with checks that stop bad data before it reaches anyone.
How I approached it
- 01
Ingest. Downloaded three months of yellow-taxi trip records from the NYC TLC as Parquet files, read them with PySpark and wrote them to a Bronze Delta Lake table partitioned by pickup month. The rows are untouched apart from two audit columns: ingestion time and source file.
- 02
Clean. A PySpark job turns Bronze into Silver. It enforces snake_case names and explicit types, drops rows with nulls in critical fields, removes duplicates on a composite key (pickup and drop-off time, both locations, fare) and filters outliers such as fares over $500 or speeds over 100 mph. It then derives trip duration, average speed, time of day, a weekend flag and tip percentage, and writes Delta partitioned by pickup date.
- 03
Aggregate. Built five Gold Delta tables with Spark aggregations: hourly trip metrics, a daily borough summary with day-over-day change (a lag window function per borough), peak against off-peak, payment-type breakdown and trip-distance buckets. Each is also exported as CSV for querying in Athena.
- 04
Validate. Wrote a small data-quality framework that runs after every layer: row-count and drop-percentage checks, per-column null thresholds, schema-drift detection and statistical bounds. A failed check raises an error and halts the run.
- 05
Stream. A Python producer replays trips as JSON events into a three-partition Kafka topic, keyed by pickup zone, at about 97 events a second. Spark Structured Streaming consumes them with a cap on offsets per micro-batch for backpressure, a 10-minute watermark for late events and 5-minute tumbling windows, writing checkpointed results to Delta.
- 06
Orchestrate. An Airflow DAG, running in Docker on a custom image with PySpark and Delta Lake, chains download, Bronze, Silver, Gold, the three validations and the S3 upload. It runs daily, retries a failed task twice with a 5-minute delay, and passes row counts and quality reports between tasks through XCom.
- 07
Serve. Terraform provisions a versioned S3 data lake with public access blocked, an Athena workgroup and a Glue database. A script uploads the Gold exports and registers them as Glue tables, so they can be queried in SQL from Athena with no server to manage.
The hard part
Fitting a big-data stack on one laptop. Spark had 2 GB of memory to process 9.5M rows, alongside Kafka and Airflow running in Docker. The usual approach, keeping the cleaned dataset in memory and building every summary table from it, does not fit in that space. So each Gold table is built one at a time from a fresh read of Silver: slower, but it finishes every time. The same constraint drove the smaller choices too, down to the lightest Airflow setup that would run the workflow.
Results
9.5 million trips, raw to query-ready, in one run
A single workflow takes 9,554,778 raw records through cleaning and into five tables an analyst can query straight away.
Bad data stops at the gate
Checks run between every layer, and a failure stops the pipeline before the problem can reach a report. Twelve unit tests cover the cleaning rules and the checks themselves.
Fresh numbers within minutes
The streaming path turns a live feed into per-borough trip and revenue figures every five minutes, without waiting for the daily batch.
Rebuildable from scratch
The cloud resources are defined in Terraform and the services in Docker Compose, so the whole setup can be torn down and recreated with a few commands.
Future improvements
The two paths write to separate tables. The next step is a serving layer that merges them into one view, which is what makes a Lambda architecture complete. Athena currently reads CSV exports; moving those to date-partitioned Parquet would cut how much data each query scans.
- PySpark
- Delta Lake
- Kafka
- Spark Structured Streaming
- Airflow
- AWS S3
- Glue
- Athena
- Terraform
- Docker