Home

Data eng

NYC Taxi Lakehouse: Batch and Streaming Pipeline

GitHub ↗

NYC TLC trip records

Parquet9.5M rowsJan to Mar 2024

Batch path

Accuracy · Airflow DAG, daily

Apache AirflowDocker
  1. Bronze

    Apache SparkDelta Lake

    PySparkParquet → Delta

    Raw rows plus ingestion time and source file

  2. Quality gaterows · nulls · schema · ranges

  3. Silver

    Apache SparkDelta Lake

    PySparkDelta Lake

    Typed, de-duplicated, outliers removed, enriched

  4. Quality gate

  5. Gold

    Apache SparkDelta Lake

    Window functionsDelta + CSV

    Five aggregate tables, partitioned by date

  6. Quality gate

  7. Data lake and SQL

    TerraformAWS

    S3Glue CatalogAthena

    Provisioned with Terraform

Out: Gold tables queryable in SQL through Athena

Streaming path

Freshness · results within minutes

Docker
  1. Producer

    Python

    PythonJSON events

    Replays trips at about 97 events a second

  2. Kafka topic

    Apache Kafka

    3 partitions

    Keyed by pickup zone to keep order

  3. Stream processing

    Apache Spark

    Structured Streaming

    5-minute tumbling windows, 10-minute watermark, backpressure cap

  4. Streaming tables

    Delta Lake

    Delta Lake

    Checkpointed, so a restart resumes where it stopped

Out: trips and revenue per borough, every 5 minutes

A Lambda-style design: the batch path reprocesses everything for accuracy, while the streaming path gives a fresh view within 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

  1. 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.

  2. 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.

  3. 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.

  4. 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.

  5. 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.

  6. 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.

  7. 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.