Data engineeringEnd-to-end pipeline

Mile High Signal: Denver 311 Data Pipeline

A scheduled ETL pipeline over Denver’s 311 service-request feed, from ingestion through a warehouse to a live dashboard.

Denver 311 dashboard
8 / 8components built
69automated tests
80%coverage gate in CI
5layers, one DAG

The problem

Denver publishes 311 requests (potholes, graffiti, illegal dumping, snow removal) as a rolling 12-month feed. It's useful data that is awkward to work with.

The API silently truncates any query past its record cap. Timestamps arrive as epoch milliseconds. Records age out of the window, so there's no history for questions like "has pothole response time improved year over year?" And with no change-data-capture, a naive re-pull double-counts updated requests.

Architecture

Ingest

Paged extract with a watermark, landed to S3 as gzip NDJSON

Process

PySpark cleans, casts, dedups, and derives resolution hours and SLA flags

Validate

Null, range, and schema-drift checks gate the load

Warehouse

Postgres star schema with agency, service, neighborhood, and date dimensions

Serve

Dashboard on the fact table, all scheduled by Airflow

Dashboard with request volume and resolution hours by agency
Requests and average resolution time by agency.
Map of 311 request locations across Denver
Where requests come from across the city.

Key design decisions

Raw data is immutable NDJSON

Raw stays a faithful copy of the API. NDJSON absorbs new fields without breaking; Parquet starts at the Spark output.

The watermark lives in S3

It survives Airflow rebuilds and backfills run outside Airflow, at the cost of one small JSON object.

At-least-once delivery

The watermark advances only after every file is durably written. Dedup on case number makes replays safe.

The endpoint is discovered

A stable dataset ID resolves to the live query URL at runtime, so a source migration doesn’t break the pipeline.

Responsible AI check

Bad data never loads. Quality gates check nulls, ranges, and schema drift before the warehouse load, and a failed check fails the whole run rather than letting questionable numbers reach a dashboard.

Count the cost. At Denver’s data volume, DuckDB or pandas would be faster and cheaper than Spark. The README says so plainly: Spark is there because the pattern scales, not because this volume needs it.

Stated limits. A single timestamp watermark misses records updated after creation, and a hard quality gate would become an outage risk across many sources. The README documents what would replace each piece at scale.