From f3121aa78e77be46b73415966ae5c656a01a7855 Mon Sep 17 00:00:00 2001 From: gregpawin Date: Mon, 18 May 2026 17:48:38 -0700 Subject: [PATCH 1/4] Added early pipeline code --- data-science/beta_pipeline/README.md | 263 +++++++++ data-science/beta_pipeline/docker-compose.yml | 20 + data-science/beta_pipeline/parking_db.py | 522 ++++++++++++++++++ .../beta_pipeline/parking_db_explore.ipynb | 451 +++++++++++++++ data-science/beta_pipeline/parking_postgis.py | 249 +++++++++ data-science/beta_pipeline/requirements.txt | 3 + 6 files changed, 1508 insertions(+) create mode 100644 data-science/beta_pipeline/README.md create mode 100644 data-science/beta_pipeline/docker-compose.yml create mode 100644 data-science/beta_pipeline/parking_db.py create mode 100644 data-science/beta_pipeline/parking_db_explore.ipynb create mode 100644 data-science/beta_pipeline/parking_postgis.py create mode 100644 data-science/beta_pipeline/requirements.txt diff --git a/data-science/beta_pipeline/README.md b/data-science/beta_pipeline/README.md new file mode 100644 index 00000000..9706d86d --- /dev/null +++ b/data-science/beta_pipeline/README.md @@ -0,0 +1,263 @@ +# LA Parking Citations DB + +A lightweight SQLite-backed store for the City of Los Angeles +[parking citations dataset](https://data.lacity.org/Transportation/Parking-Citations/4f5p-udkv). + +The flat CSV download is bulk-loaded once, then kept fresh from the +[Socrata API](https://dev.socrata.com/foundry/data.lacity.org/4f5p-udkv) by +fetching newest-first and stopping as soon as a `ticket_number` we already +have appears. + +## Layout + +| File | Purpose | +| --- | --- | +| `parking_db.py` | SQLite logic + CLI (`init`, `load-csv`, `sync`, `stats`). | +| `parking_postgis.py` | PostGIS logic + CLI (`init`, `load-csv`, `stats`). | +| `docker-compose.yml` | Local PostGIS 16 (`postgis/postgis` image). | +| `parking_db_explore.ipynb` | Notebook walkthrough — schema, queries, geospatial demo. | +| `Parking_Citations_*.csv` | The flat-file dump from data.lacity.org. | +| `parking_citations.db` | SQLite database (created on first load). | +| `requirements.txt` | Python deps (`polars`, `python-dotenv`, `psycopg`). | +| `.env` / `.env.example` | Socrata token + `DATABASE_URL` for PostGIS. | + +## Setup + +The project venv is managed with [`uv`](https://docs.astral.sh/uv/) (no `pip` +inside the venv). + +```bash +uv pip install -r requirements.txt +``` + +Drop your Socrata app token into `.env` (get one at +[data.lacity.org/profile/app_tokens](https://data.lacity.org/profile/app_tokens)): + +``` +SOCRATA_APP_TOKEN=your_real_token_here +``` + +The token isn't strictly required — the API works anonymously — but it raises +the rate limit and is recommended for any scheduled sync. + +Copy `.env.example` to `.env` and adjust values as needed. + +## PostGIS (Docker) + +For a spatial database with a native `geometry` column (and a path toward +cloud-hosted Postgres later), use the bundled Docker stack. + +### 1. Start PostGIS + +Requires [Docker Desktop](https://www.docker.com/products/docker-desktop/). + +```bash +docker compose up -d +``` + +Wait until healthy (`docker compose ps` should show `healthy`). Default +connection (also in `.env.example`): + +``` +postgresql://parking:parking@localhost:5432/parking +``` + +### 2. Install Python deps + +```bash +uv pip install -r requirements.txt +``` + +On Windows with a `venv` folder: + +```powershell +.\venv\Scripts\python.exe -m pip install -r requirements.txt +``` + +Set `DATABASE_URL` in `.env` if you change credentials or port. + +### 3. Bulk load the CSV + +Same streaming pipeline as SQLite — batches via Polars, `issue_date` +normalization, `ON CONFLICT DO NOTHING` on `ticket_number`. Each row also +gets a `geom` column (`geometry(Point, 4326)`) from `geocodelocation` WKT, +falling back to `loc_long` / `loc_lat` when WKT is missing. + +```bash +# macOS / Linux +.venv/bin/python parking_postgis.py init +.venv/bin/python parking_postgis.py load-csv Parking_Citations_20250811.csv + +# Windows +.\venv\Scripts\python.exe parking_postgis.py init +.\venv\Scripts\python.exe parking_postgis.py load-csv Parking_Citations_20250811.csv +``` + +Use your actual `Parking_Citations_*.csv` filename. Tune `--batch-size` if +needed. Re-running `load-csv` is safe (duplicates are skipped). + +### 4. Verify + +```bash +python parking_postgis.py stats +``` + +Or in `psql` (via Docker): + +```bash +docker compose exec db psql -U parking -d parking -c "SELECT COUNT(*), COUNT(geom) FROM citations;" +``` + +Example spatial query: + +```sql +SELECT ticket_number, ST_AsText(geom) +FROM citations +WHERE geom IS NOT NULL +LIMIT 5; +``` + +API incremental sync for PostGIS is not wired yet — use `parking_db.py sync` +for SQLite today, or extend `parking_postgis.py` with the same Socrata loop. + +## Updating the DB (SQLite) + +There are two operations: a one-time **bulk load** from the CSV, and a +recurring **sync** from the API. + +### 1. Initial bulk load (one-time, ~minutes) + +```bash +.venv/bin/python parking_db.py load-csv Parking_Citations_20260426.csv +``` + +What it does: + +- Creates `parking_citations.db` if missing. +- Streams the CSV in 100k-row batches via Polars (memory-safe — the 6.2 GB + file is never fully loaded). +- Normalizes `issue_date` (`"2025 Apr 26 12:00:00 AM"` → ISO 8601). +- `INSERT OR IGNORE` keyed on `ticket_number`, so re-running is safe. +- Logs the run in the `sync_log` table. + +Adjust `--batch-size` if you want to tune memory vs. throughput. +Equivalent from Python: + +```python +import parking_db +parking_db.bulk_load_csv("Parking_Citations_20260426.csv", "parking_citations.db") +``` + +### 2. Incremental sync from the API (recurring) + +```bash +.venv/bin/python parking_db.py sync +``` + +What it does: + +- Queries `https://data.lacity.org/resource/4f5p-udkv.json` + ordered by `:updated_at DESC`, paginating with `$limit` / `$offset`. +- For each page, checks which `ticket_number`s are already in the DB. +- Inserts the new ones via `INSERT OR IGNORE`. +- **Stops as soon as any page contains a `ticket_number` we already have** + — that's the "we're caught up" signal you asked for. +- Reads `SOCRATA_APP_TOKEN` from `.env` automatically; override with + `--app-token YOUR_TOKEN` if needed. + +Useful flags: + +```bash +.venv/bin/python parking_db.py sync --page-size 1000 # default +.venv/bin/python parking_db.py sync --max-pages 5 # cap pages (debug) +.venv/bin/python parking_db.py sync --app-token TOKEN # override .env +.venv/bin/python parking_db.py sync --db /tmp/test.db # different DB +``` + +Equivalent from Python: + +```python +import parking_db +result = parking_db.update_from_api("parking_citations.db") +# {'inserted': 247, 'pages': 1, 'caught_up': 1} +``` + +### 3. Verify + +```bash +.venv/bin/python parking_db.py stats +``` + +Prints row count, the latest `issue_date` in the DB, and the most recent +`sync_log` entry. + +### Automate it + +To keep the DB fresh you can drop the sync into `cron` (or `launchd`): + +```cron +# Every hour at :05 +5 * * * * cd /Users/gregpawin/Downloads/parking_db && \ + .venv/bin/python parking_db.py sync >> sync.log 2>&1 +``` + +Each run is cheap when nothing's new — usually one API page and an early exit. + +## Schema + +A single `citations` table keyed on `ticket_number` (`WITHOUT ROWID`), plus +indexes on `issue_date`, `violation_code`, and `make`. See +[`SCHEMA_SQL`](parking_db.py) in `parking_db.py` for the full DDL. + +The `geocodelocation` column is stored as WKT (`POINT (lon lat)`) regardless +of source — the API's GeoJSON `Point` is converted on ingest so the column +shape matches the CSV's format. + +A `sync_log` audit table records every load/sync (start, end, source, rows +inserted, notes). + +## Querying + +Anything that talks to SQLite works. From Python with Polars: + +```python +import sqlite3, polars as pl + +with sqlite3.connect("parking_citations.db") as conn: + df = pl.read_database( + """ + SELECT violation_description, COUNT(*) AS n, + ROUND(AVG(fine_amount), 2) AS avg_fine + FROM citations + WHERE violation_description IS NOT NULL + GROUP BY violation_description + ORDER BY n DESC + LIMIT 15 + """, + connection=conn, + ) +``` + +For ad-hoc exploration, the SQLite CLI works fine too: + +```bash +sqlite3 parking_citations.db +sqlite> SELECT COUNT(*) FROM citations; +sqlite> SELECT * FROM sync_log ORDER BY id DESC LIMIT 5; +``` + +## Caveats + +- **`:updated_at` vs `:created_at`.** The sync orders by `:updated_at DESC`, + so it picks up both new records and edits to existing ones. If a back-dated + correction lands without any genuinely new records, the very first page + will already contain a known `ticket_number` and the loop exits — that's + correct "caught up" behavior, but it means corrections aren't applied + (the table uses `INSERT OR IGNORE`). Switch to `INSERT OR REPLACE` and + remove the early-exit-on-match if you want corrections to flow through. +- **Future-dated rows.** A handful of rows in the dataset have `issue_date` + values years in the future. They're stored as-is; filter them out in your + queries if they're a problem. +- **Geospatial ops.** Lat/long are stored as floats; the WKT column is a + string. For real geometry operations, install `shapely` / `geopandas`, + or migrate the DB to [SpatiaLite](https://www.gaia-gis.it/fossil/libspatialite). diff --git a/data-science/beta_pipeline/docker-compose.yml b/data-science/beta_pipeline/docker-compose.yml new file mode 100644 index 00000000..0bb0d323 --- /dev/null +++ b/data-science/beta_pipeline/docker-compose.yml @@ -0,0 +1,20 @@ +services: + db: + image: postgis/postgis:16-3.4 + environment: + POSTGRES_USER: parking + POSTGRES_PASSWORD: parking + POSTGRES_DB: parking + ports: + - "5432:5432" + volumes: + - postgis_data:/var/lib/postgresql/data + healthcheck: + test: ["CMD-SHELL", "pg_isready -U parking -d parking"] + interval: 5s + timeout: 5s + retries: 10 + start_period: 10s + +volumes: + postgis_data: diff --git a/data-science/beta_pipeline/parking_db.py b/data-science/beta_pipeline/parking_db.py new file mode 100644 index 00000000..459fd7cf --- /dev/null +++ b/data-science/beta_pipeline/parking_db.py @@ -0,0 +1,522 @@ +"""Lightweight SQLite-backed store for LA parking citations. + +Bulk-loads from the city's CSV dump and incrementally syncs the newest records +from the Socrata API at https://dev.socrata.com/foundry/data.lacity.org/4f5p-udkv. + +CLI: + python parking_db.py init + python parking_db.py load-csv Parking_Citations_20260426.csv + python parking_db.py sync [--app-token TOKEN] + python parking_db.py stats +""" +from __future__ import annotations + +import argparse +import json +import os +import sqlite3 +import sys +import time +import urllib.error +import urllib.parse +import urllib.request +from collections.abc import Iterator +from pathlib import Path +from typing import Any + +import polars as pl + +# Load .env (sitting next to this file) into os.environ if python-dotenv is +# available. Falls through silently if it isn't — env vars set in the shell +# still work. +try: + from dotenv import load_dotenv as _load_dotenv + + _load_dotenv(Path(__file__).with_name(".env")) +except ImportError: + pass + +DATASET_ID = "4f5p-udkv" +API_URL = f"https://data.lacity.org/resource/{DATASET_ID}.json" +DB_FILENAME = "parking_citations.db" +APP_TOKEN_ENV = "SOCRATA_APP_TOKEN" + +# Canonical column order — used everywhere we INSERT. +COLUMNS: tuple[str, ...] = ( + "ticket_number", "issue_date", "issue_time", "meter_id", "marked_time", + "rp_state_plate", "plate_expiry_date", "vin", "make", "body_style", + "color", "location", "route", "agency", "violation_code", + "violation_description", "fine_amount", "agency_desc", "color_desc", + "body_style_desc", "loc_lat", "loc_long", "geocodelocation", +) + +POLARS_SCHEMA: dict[str, pl.DataType] = { + "ticket_number": pl.String, + "issue_date": pl.String, + "issue_time": pl.String, + "meter_id": pl.String, + "marked_time": pl.String, + "rp_state_plate": pl.String, + "plate_expiry_date": pl.String, + "vin": pl.String, + "make": pl.String, + "body_style": pl.String, + "color": pl.String, + "location": pl.String, + "route": pl.String, + "agency": pl.Int32, + "violation_code": pl.String, + "violation_description": pl.String, + "fine_amount": pl.Float64, + "agency_desc": pl.String, + "color_desc": pl.String, + "body_style_desc": pl.String, + "loc_lat": pl.Float64, + "loc_long": pl.Float64, + "geocodelocation": pl.String, +} + +CSV_DATE_FORMAT = "%Y %b %d %I:%M:%S %p" + +SCHEMA_SQL = """ +CREATE TABLE IF NOT EXISTS citations ( + ticket_number TEXT PRIMARY KEY, + issue_date TEXT, + issue_time TEXT, + meter_id TEXT, + marked_time TEXT, + rp_state_plate TEXT, + plate_expiry_date TEXT, + vin TEXT, + make TEXT, + body_style TEXT, + color TEXT, + location TEXT, + route TEXT, + agency INTEGER, + violation_code TEXT, + violation_description TEXT, + fine_amount REAL, + agency_desc TEXT, + color_desc TEXT, + body_style_desc TEXT, + loc_lat REAL, + loc_long REAL, + geocodelocation TEXT +) WITHOUT ROWID; + +CREATE TABLE IF NOT EXISTS sync_log ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + started_at TEXT NOT NULL, + finished_at TEXT, + source TEXT NOT NULL, + rows_inserted INTEGER DEFAULT 0, + notes TEXT +); +""" + +INDEX_SQL = """ +CREATE INDEX IF NOT EXISTS idx_citations_issue_date ON citations(issue_date); +CREATE INDEX IF NOT EXISTS idx_citations_violation_code ON citations(violation_code); +CREATE INDEX IF NOT EXISTS idx_citations_make ON citations(make); +""" + +INSERT_SQL = ( + f"INSERT OR IGNORE INTO citations ({', '.join(COLUMNS)}) " + f"VALUES ({', '.join('?' * len(COLUMNS))})" +) + + +def _connect(db_path: str | Path, *, fast: bool = False) -> sqlite3.Connection: + conn = sqlite3.connect(str(db_path)) + if fast: + # Trade durability for speed — only safe during bulk load. + conn.execute("PRAGMA journal_mode = OFF") + conn.execute("PRAGMA synchronous = OFF") + conn.execute("PRAGMA temp_store = MEMORY") + conn.execute("PRAGMA cache_size = -200000") + return conn + + +def init_db(db_path: str | Path = DB_FILENAME) -> None: + """Create schema + indexes (idempotent).""" + conn = _connect(db_path) + try: + conn.executescript(SCHEMA_SQL) + conn.executescript(INDEX_SQL) + conn.commit() + finally: + conn.close() + + +# ---------------------------------------------------------------------------- +# CSV bulk load +# ---------------------------------------------------------------------------- + +def _csv_batches(csv_path: str | Path, batch_size: int) -> Iterator[pl.DataFrame]: + reader = pl.read_csv_batched( + str(csv_path), + schema_overrides=POLARS_SCHEMA, + null_values=["", "NA", "N/A"], + ignore_errors=True, + batch_size=batch_size, + ) + while True: + batches = reader.next_batches(1) + if not batches: + return + yield batches[0] + + +def _normalize_csv_batch(df: pl.DataFrame) -> pl.DataFrame: + return df.with_columns( + pl.col("issue_date") + .str.strptime(pl.Datetime, format=CSV_DATE_FORMAT, strict=False) + .dt.strftime("%Y-%m-%dT%H:%M:%S.000") + ).select(list(COLUMNS)) + + +def bulk_load_csv( + csv_path: str | Path, + db_path: str | Path = DB_FILENAME, + *, + batch_size: int = 100_000, + progress: bool = True, +) -> int: + """Stream a Parking_Citations CSV into SQLite. Returns rows inserted.""" + init_db(db_path) + started = time.time() + started_iso = time.strftime("%Y-%m-%dT%H:%M:%SZ", time.gmtime(started)) + + total = 0 + conn = _connect(db_path, fast=True) + try: + log_id = conn.execute( + "INSERT INTO sync_log(started_at, source) VALUES(?, 'csv')", + (started_iso,), + ).lastrowid + conn.commit() + + for i, batch in enumerate(_csv_batches(csv_path, batch_size), start=1): + rows = list(_normalize_csv_batch(batch).iter_rows()) + with conn: + conn.executemany(INSERT_SQL, rows) + total += len(rows) + if progress: + elapsed = time.time() - started + rate = total / elapsed if elapsed else 0 + print( + f" batch {i}: +{len(rows):,} total={total:,} " + f"elapsed={elapsed:.1f}s rate={rate:,.0f} rows/s", + flush=True, + ) + + conn.execute( + "UPDATE sync_log SET finished_at=?, rows_inserted=? WHERE id=?", + (time.strftime("%Y-%m-%dT%H:%M:%SZ"), total, log_id), + ) + conn.commit() + finally: + conn.close() + return total + + +# ---------------------------------------------------------------------------- +# Socrata API sync +# ---------------------------------------------------------------------------- + +def _geojson_point_to_wkt(geom: dict | None) -> str | None: + if not geom or geom.get("type") != "Point": + return None + coords = geom.get("coordinates") or [] + if len(coords) < 2: + return None + return f"POINT ({coords[0]} {coords[1]})" + + +def _coerce(value: Any, kind: type) -> Any: + if value is None or value == "": + return None + try: + return kind(value) + except (TypeError, ValueError): + return None + + +def _api_record_to_row(record: dict) -> tuple[Any, ...]: + """Map a Socrata JSON record to the canonical column tuple.""" + geo = record.get("geocodelocation") + if isinstance(geo, dict): + geo = _geojson_point_to_wkt(geo) + return ( + record.get("ticket_number"), + record.get("issue_date"), + record.get("issue_time"), + record.get("meter_id"), + record.get("marked_time"), + record.get("rp_state_plate"), + record.get("plate_expiry_date"), + record.get("vin"), + record.get("make"), + record.get("body_style"), + record.get("color"), + record.get("location"), + record.get("route"), + _coerce(record.get("agency"), int), + record.get("violation_code"), + record.get("violation_description"), + _coerce(record.get("fine_amount"), float), + record.get("agency_desc"), + record.get("color_desc"), + record.get("body_style_desc"), + _coerce(record.get("loc_lat"), float), + _coerce(record.get("loc_long"), float), + geo, + ) + + +def _fetch_page( + *, + offset: int, + limit: int, + app_token: str | None, + timeout: float, + retries: int = 3, +) -> list[dict]: + params = { + "$limit": limit, + "$offset": offset, + "$order": ":updated_at DESC", + } + url = f"{API_URL}?{urllib.parse.urlencode(params)}" + req = urllib.request.Request(url) + if app_token: + req.add_header("X-App-Token", app_token) + + last_err: Exception | None = None + for attempt in range(1, retries + 1): + try: + with urllib.request.urlopen(req, timeout=timeout) as resp: + return json.loads(resp.read().decode("utf-8")) + except urllib.error.HTTPError as e: + # Surface Socrata's JSON error body — way more useful than a + # bare stack trace. 4xx responses (e.g. bad token) shouldn't + # be retried; only retry 5xx / network blips. + body = e.read().decode("utf-8", errors="replace") if e.fp else "" + hint = "" + if e.code == 403 and "Invalid app_token" in body: + hint = ( + "\n -> The Socrata App Token is being rejected. Verify" + " it at https://data.lacity.org/profile/app_tokens" + " (App Token, not API Key), or unset SOCRATA_APP_TOKEN" + " in .env to fall back to anonymous access." + ) + msg = f"HTTP {e.code} from {API_URL}: {body.strip()}{hint}" + if 400 <= e.code < 500: + raise RuntimeError(msg) from e + last_err = RuntimeError(msg) + time.sleep(min(2 ** attempt, 10)) + except (urllib.error.URLError, TimeoutError) as e: + last_err = e + time.sleep(min(2 ** attempt, 10)) + assert last_err is not None + raise last_err + + +def _existing_ticket_numbers( + conn: sqlite3.Connection, ticket_numbers: list[str] +) -> set[str]: + if not ticket_numbers: + return set() + placeholders = ",".join("?" * len(ticket_numbers)) + cur = conn.execute( + f"SELECT ticket_number FROM citations WHERE ticket_number IN ({placeholders})", + ticket_numbers, + ) + return {row[0] for row in cur} + + +def update_from_api( + db_path: str | Path = DB_FILENAME, + *, + app_token: str | None = None, + page_size: int = 1000, + max_pages: int | None = None, + timeout: float = 60.0, + progress: bool = True, +) -> dict[str, int]: + """Pull newest-first from the Socrata API and insert into the DB. + + Stops as soon as any page contains a ``ticket_number`` already present + locally — at that point the DB is considered caught up. + Returns ``{"inserted": N, "pages": P, "caught_up": 0|1}``. + + If ``app_token`` is None, falls back to the ``SOCRATA_APP_TOKEN`` env var + (loaded from ``.env`` at import time). + """ + if app_token is None: + app_token = os.getenv(APP_TOKEN_ENV) or None + # Treat the placeholder shipped in the example .env as "no token". + if app_token == "PASTE_YOUR_TOKEN_HERE": + app_token = None + + init_db(db_path) + started = time.strftime("%Y-%m-%dT%H:%M:%SZ", time.gmtime()) + + conn = _connect(db_path) + inserted = 0 + pages = 0 + caught_up = False + try: + log_id = conn.execute( + "INSERT INTO sync_log(started_at, source) VALUES(?, 'api')", + (started,), + ).lastrowid + conn.commit() + + offset = 0 + while True: + if max_pages is not None and pages >= max_pages: + break + + records = _fetch_page( + offset=offset, limit=page_size, + app_token=app_token, timeout=timeout, + ) + pages += 1 + if not records: + # Reached the end of the dataset without ever matching. + caught_up = True + break + + ticket_numbers = [ + r["ticket_number"] for r in records if r.get("ticket_number") + ] + existing = _existing_ticket_numbers(conn, ticket_numbers) + new_records = [ + r for r in records if r.get("ticket_number") not in existing + ] + rows = [_api_record_to_row(r) for r in new_records] + if rows: + with conn: + conn.executemany(INSERT_SQL, rows) + inserted += len(rows) + + if progress: + print( + f" page {pages}: fetched={len(records)} " + f"new={len(rows)} matched={len(existing)} " + f"total_new={inserted}", + flush=True, + ) + + if existing: + caught_up = True + break + offset += page_size + + conn.execute( + "UPDATE sync_log SET finished_at=?, rows_inserted=?, notes=? WHERE id=?", + ( + time.strftime("%Y-%m-%dT%H:%M:%SZ"), + inserted, + "caught_up" if caught_up else "page_limit_reached", + log_id, + ), + ) + conn.commit() + finally: + conn.close() + + return {"inserted": inserted, "pages": pages, "caught_up": int(caught_up)} + + +# ---------------------------------------------------------------------------- +# Stats / CLI +# ---------------------------------------------------------------------------- + +def db_stats(db_path: str | Path = DB_FILENAME) -> dict[str, Any]: + conn = _connect(db_path) + try: + cur = conn.cursor() + count = cur.execute("SELECT COUNT(*) FROM citations").fetchone()[0] + max_date = cur.execute( + "SELECT MAX(issue_date) FROM citations" + ).fetchone()[0] + last_sync = cur.execute( + "SELECT source, started_at, finished_at, rows_inserted, notes " + "FROM sync_log ORDER BY id DESC LIMIT 1" + ).fetchone() + finally: + conn.close() + return { + "row_count": count, + "max_issue_date": max_date, + "last_sync": ( + None if last_sync is None + else dict(zip( + ("source", "started_at", "finished_at", + "rows_inserted", "notes"), + last_sync, + )) + ), + } + + +def _build_parser() -> argparse.ArgumentParser: + p = argparse.ArgumentParser(description=__doc__.splitlines()[0]) + sub = p.add_subparsers(dest="cmd", required=True) + + p_init = sub.add_parser("init", help="Create the SQLite DB schema.") + p_init.add_argument("--db", default=DB_FILENAME) + + p_load = sub.add_parser("load-csv", help="Bulk-load CSV into SQLite.") + p_load.add_argument("csv", help="Path to Parking_Citations CSV") + p_load.add_argument("--db", default=DB_FILENAME) + p_load.add_argument("--batch-size", type=int, default=100_000) + + p_sync = sub.add_parser("sync", help="Pull newest records from Socrata API.") + p_sync.add_argument("--db", default=DB_FILENAME) + p_sync.add_argument( + "--app-token", + default=None, + help=( + "Socrata app token (raises rate limits). Defaults to the " + f"{APP_TOKEN_ENV} env var (auto-loaded from .env)." + ), + ) + p_sync.add_argument("--page-size", type=int, default=1000) + p_sync.add_argument("--max-pages", type=int, default=None, + help="Cap the number of pages fetched (debug).") + + p_stats = sub.add_parser("stats", help="Show DB stats.") + p_stats.add_argument("--db", default=DB_FILENAME) + return p + + +def main(argv: list[str] | None = None) -> int: + args = _build_parser().parse_args(argv) + if args.cmd == "init": + init_db(args.db) + print(f"Initialized {args.db}") + elif args.cmd == "load-csv": + n = bulk_load_csv(args.csv, args.db, batch_size=args.batch_size) + print(f"Loaded {n:,} rows into {args.db}") + elif args.cmd == "sync": + result = update_from_api( + args.db, + app_token=args.app_token, + page_size=args.page_size, + max_pages=args.max_pages, + ) + print( + f"Inserted {result['inserted']:,} new rows over {result['pages']} " + f"pages (caught_up={bool(result['caught_up'])})" + ) + elif args.cmd == "stats": + print(json.dumps(db_stats(args.db), indent=2, default=str)) + return 0 + + +if __name__ == "__main__": + sys.exit(main()) diff --git a/data-science/beta_pipeline/parking_db_explore.ipynb b/data-science/beta_pipeline/parking_db_explore.ipynb new file mode 100644 index 00000000..9582a72e --- /dev/null +++ b/data-science/beta_pipeline/parking_db_explore.ipynb @@ -0,0 +1,451 @@ +{ + "cells": [ + { + "cell_type": "code", + "metadata": {}, + "source": [ + "# LA Parking Citations — Exploration\n", + "\n", + "Loading `Parking_Citations_20260426.csv` (~6.2 GB) with [Polars](https://pola.rs/).\n", + "\n", + "The dataset has mixed types:\n", + "- **Strings**: plate state, make, color, violation description, location, etc.\n", + "- **Integers**: fine amount, agency code\n", + "- **Floats**: latitude, longitude\n", + "- **Datetimes**: issue date\n", + "- **Coordinates / geospatial**: `loc_lat` + `loc_long`, plus a WKT `POINT` in `geocodelocation`\n", + "\n", + "We use `pl.scan_csv` (lazy) so we only materialize what we ask for." + ], + "execution_count": null, + "outputs": [], + "id": "e28bdb40" + }, + { + "cell_type": "code", + "metadata": {}, + "source": [ + "from pathlib import Path\n", + "\n", + "import polars as pl\n", + "\n", + "CSV_PATH = Path(\"Parking_Citations_20260426.csv\")\n", + "print(f\"polars {pl.__version__} — file size: {CSV_PATH.stat().st_size / 1e9:.2f} GB\")" + ], + "execution_count": 1, + "outputs": [ + { + "output_type": "stream", + "text": [ + "polars 1.40.1 — file size: 6.22 GB\n" + ] + } + ], + "id": "7186a504" + }, + { + "cell_type": "code", + "metadata": {}, + "source": [ + "SCHEMA: dict[str, pl.DataType] = {\n", + " \"ticket_number\": pl.String,\n", + " \"issue_date\": pl.String,\n", + " \"issue_time\": pl.String,\n", + " \"meter_id\": pl.String,\n", + " \"marked_time\": pl.String,\n", + " \"rp_state_plate\": pl.String,\n", + " \"plate_expiry_date\": pl.String,\n", + " \"vin\": pl.String,\n", + " \"make\": pl.String,\n", + " \"body_style\": pl.String,\n", + " \"color\": pl.String,\n", + " \"location\": pl.String,\n", + " \"route\": pl.String,\n", + " \"agency\": pl.Int32,\n", + " \"violation_code\": pl.String,\n", + " \"violation_description\": pl.String,\n", + " \"fine_amount\": pl.Float64,\n", + " \"agency_desc\": pl.String,\n", + " \"color_desc\": pl.String,\n", + " \"body_style_desc\": pl.String,\n", + " \"loc_lat\": pl.Float64,\n", + " \"loc_long\": pl.Float64,\n", + " \"geocodelocation\": pl.String,\n", + "}\n", + "\n", + "lf = pl.scan_csv(\n", + " CSV_PATH,\n", + " schema_overrides=SCHEMA,\n", + " null_values=[\"\", \"NA\", \"N/A\"],\n", + " ignore_errors=True,\n", + ").with_columns(\n", + " pl.col(\"issue_date\").str.strptime(\n", + " pl.Datetime, format=\"%Y %b %d %I:%M:%S %p\", strict=False\n", + " ),\n", + ")\n", + "\n", + "lf.collect_schema()" + ], + "execution_count": 2, + "outputs": [ + { + "output_type": "execute_result", + "data": { + "text/plain": [ + "Schema([('ticket_number', String),\n", + " ('issue_date', Datetime(time_unit='us', time_zone=None)),\n", + " ('issue_time', String),\n", + " ('meter_id', String),\n", + " ('marked_time', String),\n", + " ('rp_state_plate', String),\n", + " ('plate_expiry_date', String),\n", + " ('vin', String),\n", + " ('make', String),\n", + " ('body_style', String),\n", + " ('color', String),\n", + " ('location', String),\n", + " ('route', String),\n", + " ('agency', Int32),\n", + " ('violation_code', String),\n", + " ('violation_description', String),\n", + " ('fine_amount', Float64),\n", + " ('agency_desc', String),\n", + " ('color_desc', String),\n", + " ('body_style_desc', String),\n", + " ('loc_lat', Float64),\n", + " ('loc_long', Float64),\n", + " ('geocodelocation', String)])" + ] + } + } + ], + "id": "0075d78b" + }, + { + "cell_type": "markdown", + "metadata": {}, + "source": [ + "## Preview\n", + "\n", + "`scan_csv` returns a `LazyFrame` — nothing has been loaded yet. We call `.head().collect()` to materialize just the first few rows." + ], + "id": "6b04e971" + }, + { + "cell_type": "code", + "metadata": {}, + "source": [ + "df_head = lf.head(10).collect()\n", + "df_head" + ], + "execution_count": 3, + "outputs": [ + { + "output_type": "execute_result", + "data": { + "text/html": [ + "
\n", + "shape: (10, 23)
ticket_numberissue_dateissue_timemeter_idmarked_timerp_state_plateplate_expiry_datevinmakebody_stylecolorlocationrouteagencyviolation_codeviolation_descriptionfine_amountagency_desccolor_descbody_style_descloc_latloc_longgeocodelocation
strdatetime[μs]strstrstrstrstrstrstrstrstrstrstri32strstrf64strstrstrf64f64str
"4602073232"2025-04-26 00:00:00"904"null"0000""CA""202512"null"FORD""PA""WT""1875 20TH ST W"null55"22500H""DOUBLE PARKING"68.0"55 - DOT - SOUTHERN""WHITE""PASSENGER CAR"34.038267-118.299593"POINT (-118.29959251 34.038266…
"4601302834"2025-04-26 00:00:00"830"null"0000""CA""202510"null"CHEV""PA""SL""5100 WOODMAN AVE"null53"22514""FIRE HYDRANT"68.0"53 - DOT - VALLEY""SILVER""PASSENGER CAR"34.16337-118.431059"POINT (-118.4310588 34.16337)"
"4601978496"2025-04-26 00:00:00"825""WU988""0000""CA""202508"null"TOYT""PA""SL""601 SAINT PAUL AV"null56"88.13B+""METER EXP."63.0"56 - DOT - CENTRAL""SILVER""PASSENGER CAR"34.052868-118.260923"POINT (-118.260923 34.05286752…
"4602159520"2025-04-26 00:00:00"935"null"0000""OR""202606"null"CHEV""PU""BK""25828 PRESIDENT AVE"null55"80.61""STANDNG IN ALLEY"68.0"55 - DOT - SOUTHERN""BLACK""PICK-UP TRUCK"33.788801-118.304095"POINT (-118.30409538 33.788800…
"4602061811"2025-04-26 00:00:00"1,255"null"0000""CA""202406"null"NISS""PA"null"7620 VARIEL AVE"null53"80.73.2""EXCEED 72HRS-ST"68.0"53 - DOT - VALLEY"null"PASSENGER CAR"34.208873-118.592801"POINT (-118.59280069 34.208873…
"4601664760"2025-04-26 00:00:00"942""WA126""0000""CA""202504"null"DODG""PU""GY""608 WESTLAKE AV S"null56"88.13B+""METER EXP."63.0"56 - DOT - CENTRAL""GREY""PICK-UP TRUCK"34.058402-118.274051"POINT (-118.27405053 34.058401…
"4602043423"2025-04-26 00:00:00"1,118"null"0000""CA""202605"null"FORD""PA""RD""2303 CHARLOTTE ST"null56"80.69AP+""NO STOP/STANDING"93.0"56 - DOT - CENTRAL""RED""PASSENGER CAR"34.056803-118.202765"POINT (-118.20276478 34.056803…
"4602639440"2025-04-26 00:00:00"1,244"null"0000""CA""202504"null"TOYT""PA""SL""1700 BARNETT ROAD"null56"80.61""STANDNG IN ALLEY"68.0"56 - DOT - CENTRAL""SILVER""PASSENGER CAR"34.061847-118.175331"POINT (-118.17533122 34.061846…
"4602094615"2025-04-26 00:00:00"920"null"0000""CA""0"null"BUIC""PA""BK""325 6TH ST W"null56"80.56E2""YELLOW ZONE"58.0"56 - DOT - CENTRAL""BLACK""PASSENGER CAR"34.04697-118.252961"POINT (-118.25296108 34.046970…
"4601664771"2025-04-26 00:00:00"959""WA1563""0000""CA""202507"null"TOYT""PA""BL""1242 VALENCIA ST"null56"88.13B+""METER EXP."63.0"56 - DOT - CENTRAL""BLUE""PASSENGER CAR"34.044013-118.275095"POINT (-118.27509548 34.044012…
" + ], + "text/plain": [ + "shape: (10, 23)\n", + "┌───────────┬───────────┬───────────┬──────────┬───┬───────────┬───────────┬───────────┬───────────┐\n", + "│ ticket_nu ┆ issue_dat ┆ issue_tim ┆ meter_id ┆ … ┆ body_styl ┆ loc_lat ┆ loc_long ┆ geocodelo │\n", + "│ mber ┆ e ┆ e ┆ --- ┆ ┆ e_desc ┆ --- ┆ --- ┆ cation │\n", + "│ --- ┆ --- ┆ --- ┆ str ┆ ┆ --- ┆ f64 ┆ f64 ┆ --- │\n", + "│ str ┆ datetime[ ┆ str ┆ ┆ ┆ str ┆ ┆ ┆ str │\n", + "│ ┆ μs] ┆ ┆ ┆ ┆ ┆ ┆ ┆ │\n", + "╞═══════════╪═══════════╪═══════════╪══════════╪═══╪═══════════╪═══════════╪═══════════╪═══════════╡\n", + "│ 460207323 ┆ 2025-04-2 ┆ 904 ┆ null ┆ … ┆ PASSENGER ┆ 34.038267 ┆ -118.2995 ┆ POINT (-1 │\n", + "│ 2 ┆ 6 ┆ ┆ ┆ ┆ CAR ┆ ┆ 93 ┆ 18.299592 │\n", + "│ ┆ 00:00:00 ┆ ┆ ┆ ┆ ┆ ┆ ┆ 51 34.038 │\n", + "│ ┆ ┆ ┆ ┆ ┆ ┆ ┆ ┆ 266… │\n", + "│ 460130283 ┆ 2025-04-2 ┆ 830 ┆ null ┆ … ┆ PASSENGER ┆ 34.16337 ┆ -118.4310 ┆ POINT (-1 │\n", + "│ 4 ┆ 6 ┆ ┆ ┆ ┆ CAR ┆ ┆ 59 ┆ 18.431058 │\n", + "│ ┆ 00:00:00 ┆ ┆ ┆ ┆ ┆ ┆ ┆ 8 │\n", + "│ ┆ ┆ ┆ ┆ ┆ ┆ ┆ ┆ 34.16337) │\n", + "│ 460197849 ┆ 2025-04-2 ┆ 825 ┆ WU988 ┆ … ┆ PASSENGER ┆ 34.052868 ┆ -118.2609 ┆ POINT (-1 │\n", + "│ 6 ┆ 6 ┆ ┆ ┆ ┆ CAR ┆ ┆ 23 ┆ 18.260923 │\n", + "│ ┆ 00:00:00 ┆ ┆ ┆ ┆ ┆ ┆ ┆ 34.052867 │\n", + "│ ┆ ┆ ┆ ┆ ┆ ┆ ┆ ┆ 52… │\n", + "│ 460215952 ┆ 2025-04-2 ┆ 935 ┆ null ┆ … ┆ PICK-UP ┆ 33.788801 ┆ -118.3040 ┆ POINT (-1 │\n", + "│ 0 ┆ 6 ┆ ┆ ┆ ┆ TRUCK ┆ ┆ 95 ┆ 18.304095 │\n", + "│ ┆ 00:00:00 ┆ ┆ ┆ ┆ ┆ ┆ ┆ 38 33.788 │\n", + "│ ┆ ┆ ┆ ┆ ┆ ┆ ┆ ┆ 800… │\n", + "│ 460206181 ┆ 2025-04-2 ┆ 1,255 ┆ null ┆ … ┆ PASSENGER ┆ 34.208873 ┆ -118.5928 ┆ POINT (-1 │\n", + "│ 1 ┆ 6 ┆ ┆ ┆ ┆ CAR ┆ ┆ 01 ┆ 18.592800 │\n", + "│ ┆ 00:00:00 ┆ ┆ ┆ ┆ ┆ ┆ ┆ 69 34.208 │\n", + "│ ┆ ┆ ┆ ┆ ┆ ┆ ┆ ┆ 873… │\n", + "│ 460166476 ┆ 2025-04-2 ┆ 942 ┆ WA126 ┆ … ┆ PICK-UP ┆ 34.058402 ┆ -118.2740 ┆ POINT (-1 │\n", + "│ 0 ┆ 6 ┆ ┆ ┆ ┆ TRUCK ┆ ┆ 51 ┆ 18.274050 │\n", + "│ ┆ 00:00:00 ┆ ┆ ┆ ┆ ┆ ┆ ┆ 53 34.058 │\n", + "│ ┆ ┆ ┆ ┆ ┆ ┆ ┆ ┆ 401… │\n", + "│ 460204342 ┆ 2025-04-2 ┆ 1,118 ┆ null ┆ … ┆ PASSENGER ┆ 34.056803 ┆ -118.2027 ┆ POINT (-1 │\n", + "│ 3 ┆ 6 ┆ ┆ ┆ ┆ CAR ┆ ┆ 65 ┆ 18.202764 │\n", + "│ ┆ 00:00:00 ┆ ┆ ┆ ┆ ┆ ┆ ┆ 78 34.056 │\n", + "│ ┆ ┆ ┆ ┆ ┆ ┆ ┆ ┆ 803… │\n", + "│ 460263944 ┆ 2025-04-2 ┆ 1,244 ┆ null ┆ … ┆ PASSENGER ┆ 34.061847 ┆ -118.1753 ┆ POINT (-1 │\n", + "│ 0 ┆ 6 ┆ ┆ ┆ ┆ CAR ┆ ┆ 31 ┆ 18.175331 │\n", + "│ ┆ 00:00:00 ┆ ┆ ┆ ┆ ┆ ┆ ┆ 22 34.061 │\n", + "│ ┆ ┆ ┆ ┆ ┆ ┆ ┆ ┆ 846… │\n", + "│ 460209461 ┆ 2025-04-2 ┆ 920 ┆ null ┆ … ┆ PASSENGER ┆ 34.04697 ┆ -118.2529 ┆ POINT (-1 │\n", + "│ 5 ┆ 6 ┆ ┆ ┆ ┆ CAR ┆ ┆ 61 ┆ 18.252961 │\n", + "│ ┆ 00:00:00 ┆ ┆ ┆ ┆ ┆ ┆ ┆ 08 34.046 │\n", + "│ ┆ ┆ ┆ ┆ ┆ ┆ ┆ ┆ 970… │\n", + "│ 460166477 ┆ 2025-04-2 ┆ 959 ┆ WA1563 ┆ … ┆ PASSENGER ┆ 34.044013 ┆ -118.2750 ┆ POINT (-1 │\n", + "│ 1 ┆ 6 ┆ ┆ ┆ ┆ CAR ┆ ┆ 95 ┆ 18.275095 │\n", + "│ ┆ 00:00:00 ┆ ┆ ┆ ┆ ┆ ┆ ┆ 48 34.044 │\n", + "│ ┆ ┆ ┆ ┆ ┆ ┆ ┆ ┆ 012… │\n", + "└───────────┴───────────┴───────────┴──────────┴───┴───────────┴───────────┴───────────┴───────────┘" + ] + } + } + ], + "id": "5dc7c790" + }, + { + "cell_type": "code", + "metadata": {}, + "source": [ + "df_head.schema" + ], + "execution_count": 4, + "outputs": [ + { + "output_type": "execute_result", + "data": { + "text/plain": [ + "Schema([('ticket_number', String),\n", + " ('issue_date', Datetime(time_unit='us', time_zone=None)),\n", + " ('issue_time', String),\n", + " ('meter_id', String),\n", + " ('marked_time', String),\n", + " ('rp_state_plate', String),\n", + " ('plate_expiry_date', String),\n", + " ('vin', String),\n", + " ('make', String),\n", + " ('body_style', String),\n", + " ('color', String),\n", + " ('location', String),\n", + " ('route', String),\n", + " ('agency', Int32),\n", + " ('violation_code', String),\n", + " ('violation_description', String),\n", + " ('fine_amount', Float64),\n", + " ('agency_desc', String),\n", + " ('color_desc', String),\n", + " ('body_style_desc', String),\n", + " ('loc_lat', Float64),\n", + " ('loc_long', Float64),\n", + " ('geocodelocation', String)])" + ] + } + } + ], + "id": "97d5ada3" + }, + { + "cell_type": "markdown", + "metadata": {}, + "source": [ + "## Geospatial columns\n", + "\n", + "The CSV has two flavors of location data:\n", + "\n", + "1. `loc_lat` / `loc_long` — already typed as `Float64`. Use these for any numeric/spatial filtering inside Polars.\n", + "2. `geocodelocation` — a [WKT](https://en.wikipedia.org/wiki/Well-known_text_representation_of_geometry) string like `POINT (-118.29959251 34.0382668)`. Polars keeps it as text; if you need real geometry ops (intersections, buffers, projections), pair it with `shapely` / `geopandas` (`pip install shapely geopandas`).\n", + "\n", + "Quick demo: count citations issued inside a rough bounding box around downtown LA." + ], + "id": "e73a9551" + }, + { + "cell_type": "code", + "metadata": {}, + "source": [ + "downtown_la = (\n", + " lf.filter(\n", + " pl.col(\"loc_lat\").is_between(34.03, 34.07)\n", + " & pl.col(\"loc_long\").is_between(-118.27, -118.23)\n", + " )\n", + " .select(pl.len().alias(\"citations_in_dtla\"))\n", + " .collect()\n", + ")\n", + "downtown_la" + ], + "execution_count": 5, + "outputs": [ + { + "output_type": "execute_result", + "data": { + "text/html": [ + "
\n", + "shape: (1, 1)
citations_in_dtla
u32
3174199
" + ], + "text/plain": [ + "shape: (1, 1)\n", + "┌───────────────────┐\n", + "│ citations_in_dtla │\n", + "│ --- │\n", + "│ u32 │\n", + "╞═══════════════════╡\n", + "│ 3174199 │\n", + "└───────────────────┘" + ] + } + } + ], + "id": "4740870d" + }, + { + "cell_type": "markdown", + "metadata": {}, + "source": [ + "## Building a queryable SQLite DB\n", + "\n", + "The full lazy scan above is great for one-off analyses but it has to re-read the 6.2 GB CSV every time. For day-to-day querying — and to keep the dataset fresh from the [Socrata API](https://dev.socrata.com/foundry/data.lacity.org/4f5p-udkv) — we ingest the CSV once into a local SQLite file, then top it up with newest-first API calls.\n", + "\n", + "The logic lives in [`parking_db.py`](parking_db.py); the cells below show end-to-end usage.\n", + "\n", + "**Workflow**\n", + "\n", + "1. `parking_db.init_db()` — create the schema (idempotent).\n", + "2. `parking_db.bulk_load_csv(...)` — stream the CSV in 100k-row batches; uses `INSERT OR IGNORE` so it's safe to re-run.\n", + "3. `parking_db.update_from_api(...)` — paginate Socrata sorted by `:updated_at DESC` and stop as soon as a page contains a `ticket_number` we already have.\n", + "4. `parking_db.db_stats(...)` — quick health check.\n", + "\n", + "> The bulk load is the slow step (≈25 M rows, expect a few minutes). After that, `update_from_api` is fast — usually only a page or two." + ], + "id": "179e9006" + }, + { + "cell_type": "code", + "metadata": {}, + "source": [ + "import parking_db\n", + "\n", + "DB_PATH = \"parking_citations.db\"\n", + "parking_db.init_db(DB_PATH)\n", + "parking_db.db_stats(DB_PATH)" + ], + "execution_count": null, + "outputs": [], + "id": "548918b6" + }, + { + "cell_type": "markdown", + "metadata": {}, + "source": [ + "### One-time bulk load\n", + "\n", + "⚠️ This will take a few minutes (≈25 M rows). Re-running is safe — the primary key on `ticket_number` plus `INSERT OR IGNORE` makes it idempotent." + ], + "id": "8e80e531" + }, + { + "cell_type": "code", + "metadata": {}, + "source": [ + "parking_db.bulk_load_csv(CSV_PATH, DB_PATH, batch_size=100_000)\n", + "parking_db.db_stats(DB_PATH)" + ], + "execution_count": null, + "outputs": [], + "id": "f3e0e705" + }, + { + "cell_type": "markdown", + "metadata": {}, + "source": [ + "### Incremental sync from the API\n", + "\n", + "Walks the dataset newest-first (`:updated_at DESC`). On each page it checks whether any `ticket_number` is already in the DB; the first such hit means we're caught up and we stop.\n", + "\n", + "The Socrata app token is read automatically from `.env` (env var `SOCRATA_APP_TOKEN`) — drop your token into the `.env` file in this folder and you're done. You can still pass `app_token=\"...\"` explicitly to override." + ], + "id": "6339d82f" + }, + { + "cell_type": "code", + "metadata": {}, + "source": [ + "import os\n", + "\n", + "print(\"token loaded:\", bool(os.getenv(\"SOCRATA_APP_TOKEN\")))\n", + "\n", + "parking_db.update_from_api(DB_PATH, page_size=1000)" + ], + "execution_count": null, + "outputs": [], + "id": "bf5c4de2" + }, + { + "cell_type": "markdown", + "metadata": {}, + "source": [ + "### Querying the DB with Polars\n", + "\n", + "Once the data is in SQLite you can pull arbitrary slices straight into a Polars DataFrame via `pl.read_database_uri`. SQLite returns results in milliseconds for indexed lookups." + ], + "id": "3047df23" + }, + { + "cell_type": "code", + "metadata": {}, + "source": [ + "import sqlite3\n", + "\n", + "with sqlite3.connect(DB_PATH) as conn:\n", + " top_violations = pl.read_database(\n", + " \"\"\"\n", + " SELECT violation_description, COUNT(*) AS n,\n", + " ROUND(AVG(fine_amount), 2) AS avg_fine\n", + " FROM citations\n", + " WHERE violation_description IS NOT NULL\n", + " GROUP BY violation_description\n", + " ORDER BY n DESC\n", + " LIMIT 15\n", + " \"\"\",\n", + " connection=conn,\n", + " )\n", + "top_violations" + ], + "execution_count": null, + "outputs": [], + "id": "ba414131" + } + ], + "metadata": { + "kernelspec": { + "display_name": ".venv (3.13.2)", + "language": "python", + "name": "python3" + }, + "language_info": { + "codemirror_mode": { + "name": "ipython", + "version": 3 + }, + "file_extension": ".py", + "mimetype": "text/x-python", + "name": "python", + "nbconvert_exporter": "python", + "pygments_lexer": "ipython3", + "version": "3.13.2" + } + }, + "nbformat": 4, + "nbformat_minor": 5 +} \ No newline at end of file diff --git a/data-science/beta_pipeline/parking_postgis.py b/data-science/beta_pipeline/parking_postgis.py new file mode 100644 index 00000000..3e01db21 --- /dev/null +++ b/data-science/beta_pipeline/parking_postgis.py @@ -0,0 +1,249 @@ +"""PostGIS-backed store for LA parking citations (Docker-friendly). + +Bulk-loads from the city's CSV dump into PostgreSQL + PostGIS. Reuses the CSV +streaming and normalization logic from ``parking_db``. + +CLI: + python parking_postgis.py init + python parking_postgis.py load-csv Parking_Citations_20250811.csv + python parking_postgis.py stats + +Requires a running PostGIS instance — see ``docker-compose.yml``. +""" +from __future__ import annotations + +import argparse +import json +import os +import sys +import time +from pathlib import Path +from typing import Any + +import psycopg +from psycopg.rows import dict_row + +from parking_db import COLUMNS, _csv_batches, _normalize_csv_batch + +try: + from dotenv import load_dotenv as _load_dotenv + + _load_dotenv(Path(__file__).with_name(".env")) +except ImportError: + pass + +DEFAULT_DATABASE_URL = "postgresql://parking:parking@localhost:5432/parking" +DATABASE_URL_ENV = "DATABASE_URL" + +ATTR_COLUMNS = list(COLUMNS) + +SCHEMA_SQL = """ +CREATE EXTENSION IF NOT EXISTS postgis; + +CREATE TABLE IF NOT EXISTS citations ( + ticket_number TEXT PRIMARY KEY, + issue_date TEXT, + issue_time TEXT, + meter_id TEXT, + marked_time TEXT, + rp_state_plate TEXT, + plate_expiry_date TEXT, + vin TEXT, + make TEXT, + body_style TEXT, + color TEXT, + location TEXT, + route TEXT, + agency INTEGER, + violation_code TEXT, + violation_description TEXT, + fine_amount DOUBLE PRECISION, + agency_desc TEXT, + color_desc TEXT, + body_style_desc TEXT, + loc_lat DOUBLE PRECISION, + loc_long DOUBLE PRECISION, + geocodelocation TEXT, + geom geometry(Point, 4326) +); + +CREATE TABLE IF NOT EXISTS sync_log ( + id SERIAL PRIMARY KEY, + started_at TIMESTAMPTZ NOT NULL, + finished_at TIMESTAMPTZ, + source TEXT NOT NULL, + rows_inserted INTEGER DEFAULT 0, + notes TEXT +); +""" + +INDEX_SQL = """ +CREATE INDEX IF NOT EXISTS idx_citations_issue_date + ON citations (issue_date); +CREATE INDEX IF NOT EXISTS idx_citations_violation_code + ON citations (violation_code); +CREATE INDEX IF NOT EXISTS idx_citations_make ON citations (make); +CREATE INDEX IF NOT EXISTS idx_citations_geom + ON citations USING GIST (geom); +""" + +INSERT_SQL = """ +INSERT INTO citations ( + ticket_number, issue_date, issue_time, meter_id, marked_time, + rp_state_plate, plate_expiry_date, vin, make, body_style, + color, location, route, agency, violation_code, + violation_description, fine_amount, agency_desc, color_desc, + body_style_desc, loc_lat, loc_long, geocodelocation, geom +) +VALUES ( + %(ticket_number)s, %(issue_date)s, %(issue_time)s, %(meter_id)s, + %(marked_time)s, %(rp_state_plate)s, %(plate_expiry_date)s, %(vin)s, + %(make)s, %(body_style)s, %(color)s, %(location)s, %(route)s, + %(agency)s, %(violation_code)s, %(violation_description)s, + %(fine_amount)s, %(agency_desc)s, %(color_desc)s, %(body_style_desc)s, + %(loc_lat)s, %(loc_long)s, %(geocodelocation)s, + COALESCE( + ST_GeomFromText(NULLIF(%(geocodelocation)s, ''), 4326), + CASE + WHEN %(loc_lat)s IS NOT NULL AND %(loc_long)s IS NOT NULL + THEN ST_SetSRID(ST_MakePoint(%(loc_long)s, %(loc_lat)s), 4326) + END + ) +) +ON CONFLICT (ticket_number) DO NOTHING +""" + + +def database_url(explicit: str | None = None) -> str: + return explicit or os.getenv(DATABASE_URL_ENV) or DEFAULT_DATABASE_URL + + +def connect(dsn: str | None = None, *, autocommit: bool = False) -> psycopg.Connection: + return psycopg.connect(database_url(dsn), autocommit=autocommit) + + +def init_db(dsn: str | None = None) -> None: + """Create PostGIS extension, tables, and indexes (idempotent).""" + with connect(dsn) as conn: + with conn.cursor() as cur: + cur.execute(SCHEMA_SQL) + cur.execute(INDEX_SQL) + conn.commit() + + +def _row_to_params(row: tuple[Any, ...]) -> dict[str, Any]: + return dict(zip(ATTR_COLUMNS, row, strict=True)) + + +def bulk_load_csv( + csv_path: str | Path, + dsn: str | None = None, + *, + batch_size: int = 100_000, + progress: bool = True, +) -> int: + """Stream a Parking_Citations CSV into PostGIS. Returns rows processed.""" + init_db(dsn) + url = database_url(dsn) + started = time.time() + started_iso = time.strftime("%Y-%m-%dT%H:%M:%SZ", time.gmtime(started)) + + total = 0 + with psycopg.connect(url) as conn: + with conn.cursor() as cur: + cur.execute( + "INSERT INTO sync_log (started_at, source) VALUES (%s, 'csv') " + "RETURNING id", + (started_iso,), + ) + log_id = cur.fetchone()[0] + conn.commit() + + # Faster bulk ingest; safe because we only append new batches. + cur.execute("SET synchronous_commit = OFF") + + for i, batch in enumerate(_csv_batches(csv_path, batch_size), start=1): + params = [ + _row_to_params(row) + for row in _normalize_csv_batch(batch).iter_rows() + ] + cur.executemany(INSERT_SQL, params) + conn.commit() + total += len(params) + if progress: + elapsed = time.time() - started + rate = total / elapsed if elapsed else 0 + print( + f" batch {i}: +{len(params):,} total={total:,} " + f"elapsed={elapsed:.1f}s rate={rate:,.0f} rows/s", + flush=True, + ) + + cur.execute( + "UPDATE sync_log SET finished_at = %s, rows_inserted = %s " + "WHERE id = %s", + (time.strftime("%Y-%m-%dT%H:%M:%SZ"), total, log_id), + ) + conn.commit() + return total + + +def db_stats(dsn: str | None = None) -> dict[str, Any]: + with connect(dsn) as conn: + with conn.cursor(row_factory=dict_row) as cur: + cur.execute("SELECT COUNT(*) AS n FROM citations") + count = cur.fetchone()["n"] + cur.execute( + "SELECT COUNT(*) AS n FROM citations WHERE geom IS NOT NULL" + ) + with_geom = cur.fetchone()["n"] + cur.execute("SELECT MAX(issue_date) AS d FROM citations") + max_date = cur.fetchone()["d"] + cur.execute( + "SELECT source, started_at, finished_at, rows_inserted, notes " + "FROM sync_log ORDER BY id DESC LIMIT 1" + ) + last_sync = cur.fetchone() + return { + "row_count": count, + "rows_with_geom": with_geom, + "max_issue_date": max_date, + "last_sync": last_sync, + } + + +def _build_parser() -> argparse.ArgumentParser: + p = argparse.ArgumentParser(description=__doc__.splitlines()[0]) + p.add_argument( + "--database-url", + default=None, + help=f"Postgres DSN (default: {DATABASE_URL_ENV} env or docker default)", + ) + sub = p.add_subparsers(dest="cmd", required=True) + + sub.add_parser("init", help="Create PostGIS schema and indexes.") + + p_load = sub.add_parser("load-csv", help="Bulk-load CSV into PostGIS.") + p_load.add_argument("csv", help="Path to Parking_Citations CSV") + p_load.add_argument("--batch-size", type=int, default=100_000) + + sub.add_parser("stats", help="Show DB stats.") + return p + + +def main(argv: list[str] | None = None) -> int: + args = _build_parser().parse_args(argv) + dsn = args.database_url + if args.cmd == "init": + init_db(dsn) + print(f"Initialized PostGIS at {database_url(dsn)}") + elif args.cmd == "load-csv": + n = bulk_load_csv(args.csv, dsn, batch_size=args.batch_size) + print(f"Loaded {n:,} rows into {database_url(dsn)}") + elif args.cmd == "stats": + print(json.dumps(db_stats(dsn), indent=2, default=str)) + return 0 + + +if __name__ == "__main__": + sys.exit(main()) diff --git a/data-science/beta_pipeline/requirements.txt b/data-science/beta_pipeline/requirements.txt new file mode 100644 index 00000000..af834c2d --- /dev/null +++ b/data-science/beta_pipeline/requirements.txt @@ -0,0 +1,3 @@ +polars>=1.40,<2 +python-dotenv>=1.0 +psycopg[binary]>=3.2,<4 From 281bd691a1c576e129e1f4b26234df9b3daf6db5 Mon Sep 17 00:00:00 2001 From: gregpawin Date: Mon, 8 Jun 2026 17:58:30 -0700 Subject: [PATCH 2/4] Implement architecture documentation and add pipeline orchestration for parking citations - Introduced `ARCHITECTURE.md` to detail project structure, data flow, and key functions. - Added `parking_pipeline.py` to orchestrate the loading of raw citations and rebuilding of the cleaned table. - Created `parking_clean.py` for managing the cleaned citations table. - Updated permissions for several scripts and configuration files. - Enhanced `parking_postgis.py` to support rebuilding the cleaned table after loading CSV data. --- data-science/beta_pipeline/ARCHITECTURE.md | 449 ++++++++++++++++++ data-science/beta_pipeline/README.md | 272 ++++------- data-science/beta_pipeline/docker-compose.yml | 0 data-science/beta_pipeline/parking_clean.py | 171 +++++++ data-science/beta_pipeline/parking_db.py | 0 .../beta_pipeline/parking_db_explore.ipynb | 0 .../beta_pipeline/parking_pipeline.py | 81 ++++ data-science/beta_pipeline/parking_postgis.py | 9 + data-science/beta_pipeline/requirements.txt | 0 9 files changed, 812 insertions(+), 170 deletions(-) create mode 100644 data-science/beta_pipeline/ARCHITECTURE.md mode change 100644 => 100755 data-science/beta_pipeline/README.md mode change 100644 => 100755 data-science/beta_pipeline/docker-compose.yml create mode 100644 data-science/beta_pipeline/parking_clean.py mode change 100644 => 100755 data-science/beta_pipeline/parking_db.py mode change 100644 => 100755 data-science/beta_pipeline/parking_db_explore.ipynb create mode 100644 data-science/beta_pipeline/parking_pipeline.py mode change 100644 => 100755 data-science/beta_pipeline/parking_postgis.py mode change 100644 => 100755 data-science/beta_pipeline/requirements.txt diff --git a/data-science/beta_pipeline/ARCHITECTURE.md b/data-science/beta_pipeline/ARCHITECTURE.md new file mode 100644 index 00000000..3d4b3b12 --- /dev/null +++ b/data-science/beta_pipeline/ARCHITECTURE.md @@ -0,0 +1,449 @@ +# Architecture & Code Reference + +This document describes how the LA Parking Citations project is structured: +modules, data flow, schemas, and the key functions in each file. + +## Overview + +``` +Parking_Citations_*.csv ──► parking_db.py ──► parking_citations.db (SQLite) + │ + └──► parking_postgis.py ──► citations (PostGIS raw) + │ + └──► parking_clean.py ──► citations_clean + +Socrata API ──► parking_db.py sync ──► parking_citations.db +``` + +The city publishes two access paths for the same dataset: + +| Source | URL | Used by | +| --- | --- | --- | +| CSV export | data.lacity.org (Transportation → Parking Citations) | `load-csv` in both backends | +| Socrata API | `https://data.lacity.org/resource/4f5p-udkv.json` | `parking_db.py sync` only | + +--- + +## Shared foundations + +### Canonical column set + +Both backends store the same 23 citation attributes, defined once in +`parking_db.py` as `COLUMNS`: + +``` +ticket_number, issue_date, issue_time, meter_id, marked_time, +rp_state_plate, plate_expiry_date, vin, make, body_style, +color, location, route, agency, violation_code, +violation_description, fine_amount, agency_desc, color_desc, +body_style_desc, loc_lat, loc_long, geocodelocation +``` + +`ticket_number` is the primary key everywhere. + +### CSV streaming (`parking_db.py`) + +The 6+ GB CSV is never loaded whole into memory. Two functions handle ingestion: + +#### `_csv_batches(csv_path, batch_size)` + +Uses Polars `read_csv_batched` with: + +- `POLARS_SCHEMA` — explicit types per column (strings for most fields, `Int32` for `agency`, `Float64` for coordinates and fines) +- `null_values=["", "NA", "N/A"]` +- `ignore_errors=True` — skip malformed rows rather than abort +- Default batch size: **100,000 rows** + +Yields one `pl.DataFrame` per batch. + +#### `_normalize_csv_batch(df)` + +Transforms each batch before insert: + +1. Parses `issue_date` from the CSV format `"2025 Apr 26 12:00:00 AM"` using + `CSV_DATE_FORMAT = "%Y %b %d %I:%M:%S %p"` +2. Writes it back as ISO text: `"2025-04-26T00:00:00.000"` +3. Selects columns in canonical `COLUMNS` order + +`issue_time` is **not** combined at CSV load time. It stays as a separate string +(e.g. `"904"`, `"1430"`) and is merged into a timestamp only in the PostGIS clean +step. + +### Environment loading + +All modules optionally load `.env` from the project directory via `python-dotenv`. +If the package is missing, env vars set in the shell still work. + +| Variable | Used by | Default | +| --- | --- | --- | +| `SOCRATA_APP_TOKEN` | `parking_db.py sync` | none (anonymous API access) | +| `DATABASE_URL` | PostGIS modules | `postgresql://parking:parking@localhost:5432/parking` | + +--- + +## `parking_db.py` — SQLite backend + +Single-file SQLite database (`parking_citations.db`) with bulk CSV load and +incremental API sync. + +### Schema + +**`citations`** — all 23 columns, `ticket_number TEXT PRIMARY KEY`, `WITHOUT ROWID`. + +**`sync_log`** — audit trail for every load/sync: + +| Column | Type | Notes | +| --- | --- | --- | +| `id` | INTEGER PK | Auto-increment | +| `started_at` | TEXT | ISO UTC timestamp | +| `finished_at` | TEXT | Set when run completes | +| `source` | TEXT | `'csv'` or `'api'` | +| `rows_inserted` | INTEGER | Rows processed this run | +| `notes` | TEXT | e.g. `'caught_up'` for API sync | + +**Indexes:** `issue_date`, `violation_code`, `make`. + +### Key functions + +#### `init_db(db_path)` + +Idempotent DDL: creates tables and indexes. + +#### `bulk_load_csv(csv_path, db_path, *, batch_size, progress)` + +1. Calls `init_db` +2. Opens a **fast** connection (`journal_mode=OFF`, `synchronous=OFF`) for throughput +3. Inserts a `sync_log` row with `source='csv'` +4. For each CSV batch: normalize → `executemany(INSERT OR IGNORE)` +5. Updates `sync_log` with row count and finish time +6. Returns total rows processed (including duplicates attempted) + +#### `update_from_api(db_path, *, app_token, page_size, max_pages, timeout, progress)` + +Incremental sync loop: + +``` +offset = 0 +loop: + fetch page ordered by :updated_at DESC + if empty → caught up, break + find ticket_numbers already in DB + insert new records (INSERT OR IGNORE) + if any ticket_number on this page already exists → caught up, break + offset += page_size +``` + +**API record mapping** (`_api_record_to_row`): + +- Converts Socrata GeoJSON `Point` → WKT via `_geojson_point_to_wkt` +- Coerces `agency`, `fine_amount`, `loc_lat`, `loc_long` with `_coerce` + +**HTTP handling** (`_fetch_page`): + +- Retries 5xx and network errors with exponential backoff +- Surfaces 4xx errors immediately (e.g. invalid app token with a helpful hint) +- Sends `X-App-Token` header when configured + +Returns `{"inserted": N, "pages": P, "caught_up": 0|1}`. + +#### `db_stats(db_path)` + +Returns row count, max `issue_date`, and the most recent `sync_log` entry. + +### CLI + +``` +python parking_db.py init +python parking_db.py load-csv FILE [--db PATH] [--batch-size N] +python parking_db.py sync [--db PATH] [--app-token TOKEN] [--page-size N] [--max-pages N] +python parking_db.py stats [--db PATH] +``` + +--- + +## `parking_postgis.py` — PostGIS raw store + +PostgreSQL + PostGIS backend. Reuses `_csv_batches` and `_normalize_csv_batch` +from `parking_db.py` so CSV handling is identical. + +### Schema + +**`citations`** — same 23 columns as SQLite, plus: + +```sql +geom geometry(Point, 4326) +``` + +**`sync_log`** — same purpose as SQLite, with native `TIMESTAMPTZ` types. + +**Indexes:** `issue_date`, `violation_code`, `make`, and a **GiST index on `geom`**. + +### Geometry construction (insert time) + +Each row's `geom` is built in SQL during insert: + +```sql +COALESCE( + ST_GeomFromText(NULLIF(geocodelocation, ''), 4326), -- WKT from CSV + ST_SetSRID(ST_MakePoint(loc_long, loc_lat), 4326) -- fallback to lat/long +) +``` + +Priority: WKT string first, then coordinate pair. Rows with neither get `geom = NULL`. + +### Key functions + +#### `connect(dsn, *, autocommit)` + +Wraps `psycopg.connect` with DSN resolution via `database_url()`. + +#### `init_db(dsn)` + +Creates PostGIS extension, `citations`, `sync_log`, and indexes. + +#### `bulk_load_csv(csv_path, dsn, *, batch_size, progress)` + +1. Calls `init_db` +2. Logs start in `sync_log` +3. Sets `synchronous_commit = OFF` for faster bulk ingest +4. For each batch: normalize → convert rows to dicts → `executemany(INSERT ... ON CONFLICT DO NOTHING)` +5. Commits per batch, prints progress (batch number, total, rate) +6. Updates `sync_log` on completion + +Returns total rows processed. + +#### `db_stats(dsn)` + +Returns row count, rows with non-null `geom`, max `issue_date`, and last sync log entry. + +### CLI + +``` +python parking_postgis.py init [--database-url DSN] +python parking_postgis.py load-csv FILE [--batch-size N] [--clean] +python parking_postgis.py stats [--database-url DSN] +``` + +The `--clean` flag calls `parking_clean.rebuild_clean()` after the load finishes. + +--- + +## `parking_clean.py` — Cleaned analytics table + +Transforms the raw `citations` table into a slim `citations_clean` table +suited for analysis and mapping. + +### Schema + +```sql +CREATE TABLE citations_clean ( + ticket_number TEXT PRIMARY KEY, + issue_datetime TIMESTAMPTZ NOT NULL, + violation_code TEXT, + violation_description TEXT, + fine_amount DOUBLE PRECISION, + geom geometry(Point, 4326) +); +``` + +**Indexes:** `issue_datetime`, `violation_code`, GiST on `geom`. + +### Datetime combination + +Raw data splits date and time across two columns: + +| Column | Example | Meaning | +| --- | --- | --- | +| `issue_date` | `2025-04-26T00:00:00.000` | Date (time portion is always midnight from CSV) | +| `issue_time` | `904`, `1430` | Time as HHMM without leading zeros | + +The rebuild SQL: + +1. Takes the **date** from `issue_date` +2. Parses `issue_time` by left-padding to 4 digits (`904` → `0904`) +3. Extracts hours (first 2 digits) and minutes (last 2 digits) +4. Adds an interval to the date → `TIMESTAMPTZ` + +If `issue_time` is missing or non-numeric, the datetime defaults to midnight on +the issue date. + +### Text cleaning + +- `violation_code` and `violation_description`: `BTRIM`, then `NULLIF` empty strings +- `fine_amount` and `geom`: copied as-is from raw row + +### Key functions + +#### `init_clean(dsn)` + +Calls `init_db` (raw schema) then creates `citations_clean` and its indexes. + +#### `rebuild_clean(dsn, *, progress)` + +Full refresh strategy: + +1. `TRUNCATE citations_clean` +2. `INSERT INTO citations_clean SELECT ... FROM citations WHERE ticket_number IS NOT NULL AND issue_date IS NOT NULL` +3. Returns row count + +This is idempotent and safe to re-run after every raw load. + +#### `clean_stats(dsn)` + +Returns row count, geom coverage, and min/max `issue_datetime`. + +### CLI + +``` +python parking_clean.py init [--database-url DSN] +python parking_clean.py rebuild [--database-url DSN] +python parking_clean.py stats [--database-url DSN] +``` + +--- + +## `parking_pipeline.py` — Orchestration + +Thin wrapper that chains raw load and clean rebuild. + +### `run_pipeline(csv_path, dsn, *, batch_size, skip_load, progress)` + +``` +if not skip_load: + bulk_load_csv(csv_path, dsn) # parking_postgis +rebuild_clean(dsn) # parking_clean +return { database_url, raw: db_stats(), clean: clean_stats() } +``` + +### CLI + +``` +python parking_pipeline.py run FILE [--batch-size N] [--database-url DSN] +python parking_pipeline.py clean [--database-url DSN] +``` + +The `clean` subcommand skips CSV load and only rebuilds `citations_clean`. + +--- + +## `docker-compose.yml` + +Runs PostGIS 16 on port 5432: + +| Setting | Value | +| --- | --- | +| Image | `postgis/postgis:16-3.4` | +| User / password / database | `parking` / `parking` / `parking` | +| Volume | `postgis_data` (persists between restarts) | +| Healthcheck | `pg_isready` every 5s | + +--- + +## Data flow diagrams + +### PostGIS full pipeline + +```mermaid +flowchart LR + CSV[Parking_Citations CSV] + PG[parking_postgis.py] + RAW[(citations)] + CL[parking_clean.py] + CLEAN[(citations_clean)] + + CSV -->|Polars batches| PG + PG -->|INSERT + geom| RAW + RAW -->|TRUNCATE + SELECT| CL + CL --> CLEAN +``` + +### SQLite sync loop + +```mermaid +flowchart TD + API[Socrata API] + SYNC[parking_db.py sync] + DB[(parking_citations.db)] + + API -->|pages newest-first| SYNC + SYNC -->|INSERT OR IGNORE new rows| DB + SYNC -->|stop when ticket_number match| SYNC +``` + +--- + +## Module dependency graph + +``` +parking_db.py (standalone — CSV + SQLite + API) + ↑ +parking_postgis.py (imports COLUMNS, _csv_batches, _normalize_csv_batch) + ↑ +parking_clean.py (imports connect, database_url, init_db) + ↑ +parking_pipeline.py (imports bulk_load_csv, rebuild_clean, stats helpers) +``` + +`parking_db.py` has no imports from the other modules. The PostGIS stack depends +on it only for CSV parsing shared code. + +--- + +## Design decisions + +### Why two backends? + +- **SQLite** is zero-infra, supports API sync today, and works well for ad-hoc + analysis with Polars or the sqlite3 CLI. +- **PostGIS** adds native geometry, GiST spatial indexes, and a cleaned table + for downstream analytics or cloud migration. + +### Why a separate clean table? + +Raw `citations` preserves every column from the city export. `citations_clean` +narrows to the fields most useful for analysis, combines split date/time fields +into one timestamp, trims text, and drops rows missing a ticket or date. Rebuilding +via `TRUNCATE + INSERT` keeps the clean step simple and deterministic. + +### Why `INSERT OR IGNORE` / `ON CONFLICT DO NOTHING`? + +Re-running bulk loads is safe — duplicate `ticket_number`s are skipped. The tradeoff +is that **corrections to existing tickets are not applied** during API sync. To +accept corrections, switch to `INSERT OR REPLACE` / `ON CONFLICT DO UPDATE` and +adjust the sync early-exit logic. + +### Performance choices + +| Location | Optimization | Risk | +| --- | --- | --- | +| SQLite bulk load | `journal_mode=OFF`, `synchronous=OFF` | Less crash safety during load only | +| PostGIS bulk load | `synchronous_commit=OFF` | Recent commits may be lost on crash | +| CSV reading | 100k-row Polars batches | Tunable via `--batch-size` | + +Both bulk loaders commit per batch so progress is not lost mid-run. + +--- + +## Known data quirks + +1. **Future-dated `issue_date` values** — a small number of rows have dates years + ahead. Stored as-is; filter in queries if needed. +2. **`issue_time` format** — HHMM without leading zeros. Values like `"0"` or + `"0000"` parse as midnight. +3. **Missing geometry** — not every row has coordinates. Check `geom IS NOT NULL` + for spatial analysis. +4. **API vs CSV shape** — the API returns GeoJSON for location; the CSV has WKT. + SQLite sync converts GeoJSON → WKT so the column shape is consistent. + +--- + +## Extending the project + +Common next steps: + +| Goal | Starting point | +| --- | --- | +| PostGIS API sync | Copy `update_from_api` from `parking_db.py`, adapt inserts to include `geom` | +| Incremental clean | Replace `TRUNCATE` with upsert on `ticket_number` for new rows only | +| Scheduled pipeline | cron/launchd calling `parking_pipeline.py run` after each CSV download | +| Cloud Postgres | Point `DATABASE_URL` at RDS/Supabase; same code, no Docker | diff --git a/data-science/beta_pipeline/README.md b/data-science/beta_pipeline/README.md old mode 100644 new mode 100755 index 9706d86d..d7ebdbcf --- a/data-science/beta_pipeline/README.md +++ b/data-science/beta_pipeline/README.md @@ -1,224 +1,167 @@ # LA Parking Citations DB -A lightweight SQLite-backed store for the City of Los Angeles -[parking citations dataset](https://data.lacity.org/Transportation/Parking-Citations/4f5p-udkv). +Tools for loading, syncing, and analyzing the City of Los Angeles +[parking citations dataset](https://data.lacity.org/Transportation/Parking-Citations/4f5p-udkv) +(~6 GB CSV, millions of rows). -The flat CSV download is bulk-loaded once, then kept fresh from the -[Socrata API](https://dev.socrata.com/foundry/data.lacity.org/4f5p-udkv) by -fetching newest-first and stopping as soon as a `ticket_number` we already -have appears. +Two storage backends share the same CSV parsing logic: -## Layout +| Backend | Module | Best for | +| --- | --- | --- | +| **SQLite** | `parking_db.py` | Local file DB, incremental API sync | +| **PostGIS** | `parking_postgis.py` + `parking_clean.py` | Spatial queries, cleaned analytics table | -| File | Purpose | -| --- | --- | -| `parking_db.py` | SQLite logic + CLI (`init`, `load-csv`, `sync`, `stats`). | -| `parking_postgis.py` | PostGIS logic + CLI (`init`, `load-csv`, `stats`). | -| `docker-compose.yml` | Local PostGIS 16 (`postgis/postgis` image). | -| `parking_db_explore.ipynb` | Notebook walkthrough — schema, queries, geospatial demo. | -| `Parking_Citations_*.csv` | The flat-file dump from data.lacity.org. | -| `parking_citations.db` | SQLite database (created on first load). | -| `requirements.txt` | Python deps (`polars`, `python-dotenv`, `psycopg`). | -| `.env` / `.env.example` | Socrata token + `DATABASE_URL` for PostGIS. | - -## Setup - -The project venv is managed with [`uv`](https://docs.astral.sh/uv/) (no `pip` -inside the venv). - -```bash -uv pip install -r requirements.txt -``` - -Drop your Socrata app token into `.env` (get one at -[data.lacity.org/profile/app_tokens](https://data.lacity.org/profile/app_tokens)): - -``` -SOCRATA_APP_TOKEN=your_real_token_here -``` - -The token isn't strictly required — the API works anonymously — but it raises -the rate limit and is recommended for any scheduled sync. - -Copy `.env.example` to `.env` and adjust values as needed. +PostGIS adds a **load → clean pipeline** (`parking_pipeline.py`) that builds a slim +`citations_clean` table after each full CSV load. -## PostGIS (Docker) +For module-by-module code breakdown, schemas, and data-flow diagrams, see +**[ARCHITECTURE.md](ARCHITECTURE.md)**. -For a spatial database with a native `geometry` column (and a path toward -cloud-hosted Postgres later), use the bundled Docker stack. +## Quick start -### 1. Start PostGIS +### 1. Install uv and create the venv -Requires [Docker Desktop](https://www.docker.com/products/docker-desktop/). +Requires [uv](https://docs.astral.sh/uv/) (manages Python 3.12 and dependencies). ```bash -docker compose up -d -``` - -Wait until healthy (`docker compose ps` should show `healthy`). Default -connection (also in `.env.example`): +curl -LsSf https://astral.sh/uv/install.sh | sh +source $HOME/.local/bin/env # add uv to PATH (once per shell) +cd parking +uv venv --python 3.12 .venv +uv pip install -r requirements.txt --python .venv/bin/python ``` -postgresql://parking:parking@localhost:5432/parking -``` -### 2. Install Python deps +Polars ≥1.40 needs **Python 3.10+**; the venv uses 3.12. + +### 2. Configure environment ```bash -uv pip install -r requirements.txt +cp .env.example .env ``` -On Windows with a `venv` folder: - -```powershell -.\venv\Scripts\python.exe -m pip install -r requirements.txt -``` +Edit `.env`: -Set `DATABASE_URL` in `.env` if you change credentials or port. +- `SOCRATA_APP_TOKEN` — optional; raises API rate limits for `parking_db.py sync` +- `DATABASE_URL` — PostGIS connection string (matches `docker-compose.yml`) -### 3. Bulk load the CSV +### 3. Choose a path -Same streaming pipeline as SQLite — batches via Polars, `issue_date` -normalization, `ON CONFLICT DO NOTHING` on `ticket_number`. Each row also -gets a `geom` column (`geometry(Point, 4326)`) from `geocodelocation` WKT, -falling back to `loc_long` / `loc_lat` when WKT is missing. +**SQLite (simplest — no Docker):** ```bash -# macOS / Linux -.venv/bin/python parking_postgis.py init -.venv/bin/python parking_postgis.py load-csv Parking_Citations_20250811.csv - -# Windows -.\venv\Scripts\python.exe parking_postgis.py init -.\venv\Scripts\python.exe parking_postgis.py load-csv Parking_Citations_20250811.csv +.venv/bin/python parking_db.py load-csv Parking_Citations_20260426.csv +.venv/bin/python parking_db.py sync # incremental updates +.venv/bin/python parking_db.py stats ``` -Use your actual `Parking_Citations_*.csv` filename. Tune `--batch-size` if -needed. Re-running `load-csv` is safe (duplicates are skipped). - -### 4. Verify +**PostGIS (spatial + cleaned table):** ```bash -python parking_postgis.py stats +docker compose up -d # start PostGIS +.venv/bin/python parking_pipeline.py run Parking_Citations_20250811.csv +.venv/bin/python parking_clean.py stats ``` -Or in `psql` (via Docker): +Replace `Parking_Citations_*.csv` with your actual download filename. -```bash -docker compose exec db psql -U parking -d parking -c "SELECT COUNT(*), COUNT(geom) FROM citations;" -``` +## Project layout -Example spatial query: - -```sql -SELECT ticket_number, ST_AsText(geom) -FROM citations -WHERE geom IS NOT NULL -LIMIT 5; -``` +| File | Purpose | +| --- | --- | +| `parking_db.py` | SQLite store + CLI (`init`, `load-csv`, `sync`, `stats`) | +| `parking_postgis.py` | PostGIS raw store + CLI (`init`, `load-csv`, `stats`) | +| `parking_clean.py` | Cleaned table builder (`init`, `rebuild`, `stats`) | +| `parking_pipeline.py` | Orchestrates load → clean (`run`, `clean`) | +| `docker-compose.yml` | Local PostGIS 16 | +| `parking_db_explore.ipynb` | Notebook walkthrough | +| `ARCHITECTURE.md` | Detailed code and schema reference | +| `requirements.txt` | `polars`, `python-dotenv`, `psycopg` | +| `.env` / `.env.example` | Secrets and connection strings | +| `.venv/` | Python 3.12 virtual environment (created by uv) | -API incremental sync for PostGIS is not wired yet — use `parking_db.py sync` -for SQLite today, or extend `parking_postgis.py` with the same Socrata loop. +Data files (not in repo): `Parking_Citations_*.csv`, `parking_citations.db`. -## Updating the DB (SQLite) +## SQLite workflow -There are two operations: a one-time **bulk load** from the CSV, and a -recurring **sync** from the API. +### Bulk load (one-time) -### 1. Initial bulk load (one-time, ~minutes) +Streams the CSV in 100k-row Polars batches — the full file is never loaded into +memory. Normalizes `issue_date` to ISO 8601. Safe to re-run (`INSERT OR IGNORE`). ```bash +.venv/bin/python parking_db.py init .venv/bin/python parking_db.py load-csv Parking_Citations_20260426.csv +.venv/bin/python parking_db.py load-csv Parking_Citations_20260426.csv --batch-size 50000 ``` -What it does: - -- Creates `parking_citations.db` if missing. -- Streams the CSV in 100k-row batches via Polars (memory-safe — the 6.2 GB - file is never fully loaded). -- Normalizes `issue_date` (`"2025 Apr 26 12:00:00 AM"` → ISO 8601). -- `INSERT OR IGNORE` keyed on `ticket_number`, so re-running is safe. -- Logs the run in the `sync_log` table. - -Adjust `--batch-size` if you want to tune memory vs. throughput. -Equivalent from Python: - -```python -import parking_db -parking_db.bulk_load_csv("Parking_Citations_20260426.csv", "parking_citations.db") -``` +### Incremental sync (recurring) -### 2. Incremental sync from the API (recurring) +Fetches from the Socrata API newest-first and stops when a page contains a +`ticket_number` already in the DB. ```bash .venv/bin/python parking_db.py sync +.venv/bin/python parking_db.py sync --page-size 1000 --max-pages 5 # debug ``` -What it does: +Automate with cron or launchd — each run is cheap when nothing is new. -- Queries `https://data.lacity.org/resource/4f5p-udkv.json` - ordered by `:updated_at DESC`, paginating with `$limit` / `$offset`. -- For each page, checks which `ticket_number`s are already in the DB. -- Inserts the new ones via `INSERT OR IGNORE`. -- **Stops as soon as any page contains a `ticket_number` we already have** - — that's the "we're caught up" signal you asked for. -- Reads `SOCRATA_APP_TOKEN` from `.env` automatically; override with - `--app-token YOUR_TOKEN` if needed. +## PostGIS workflow -Useful flags: +### Start the database ```bash -.venv/bin/python parking_db.py sync --page-size 1000 # default -.venv/bin/python parking_db.py sync --max-pages 5 # cap pages (debug) -.venv/bin/python parking_db.py sync --app-token TOKEN # override .env -.venv/bin/python parking_db.py sync --db /tmp/test.db # different DB +docker compose up -d +docker compose ps # wait for "healthy" ``` -Equivalent from Python: +Default connection: `postgresql://parking:parking@localhost:5432/parking` -```python -import parking_db -result = parking_db.update_from_api("parking_citations.db") -# {'inserted': 247, 'pages': 1, 'caught_up': 1} -``` +### Full pipeline (recommended) -### 3. Verify +Loads raw rows into `citations`, then rebuilds `citations_clean`: ```bash -.venv/bin/python parking_db.py stats +.venv/bin/python parking_pipeline.py run Parking_Citations_20250811.csv ``` -Prints row count, the latest `issue_date` in the DB, and the most recent -`sync_log` entry. - -### Automate it - -To keep the DB fresh you can drop the sync into `cron` (or `launchd`): +Or step by step: -```cron -# Every hour at :05 -5 * * * * cd /Users/gregpawin/Downloads/parking_db && \ - .venv/bin/python parking_db.py sync >> sync.log 2>&1 +```bash +.venv/bin/python parking_postgis.py init +.venv/bin/python parking_postgis.py load-csv Parking_Citations_20250811.csv --clean +.venv/bin/python parking_postgis.py stats +.venv/bin/python parking_clean.py rebuild # refresh clean table only ``` -Each run is cheap when nothing's new — usually one API page and an early exit. +### Cleaned table columns -## Schema +| Column | Description | +| --- | --- | +| `ticket_number` | Primary key | +| `issue_datetime` | Combined date + time (HHMM, e.g. `"904"` → 09:04) | +| `violation_code` | Trimmed text | +| `violation_description` | Trimmed text | +| `fine_amount` | Numeric fine | +| `geom` | `geometry(Point, 4326)` copied from raw row | -A single `citations` table keyed on `ticket_number` (`WITHOUT ROWID`), plus -indexes on `issue_date`, `violation_code`, and `make`. See -[`SCHEMA_SQL`](parking_db.py) in `parking_db.py` for the full DDL. +Example query: -The `geocodelocation` column is stored as WKT (`POINT (lon lat)`) regardless -of source — the API's GeoJSON `Point` is converted on ingest so the column -shape matches the CSV's format. +```sql +SELECT ticket_number, issue_datetime, violation_description, ST_AsText(geom) +FROM citations_clean +WHERE geom IS NOT NULL +ORDER BY issue_datetime DESC +LIMIT 10; +``` -A `sync_log` audit table records every load/sync (start, end, source, rows -inserted, notes). +PostGIS API sync is not implemented yet — use SQLite `sync` today, or extend +`parking_postgis.py` with the same Socrata loop from `parking_db.py`. ## Querying -Anything that talks to SQLite works. From Python with Polars: +**SQLite with Polars:** ```python import sqlite3, polars as pl @@ -238,26 +181,15 @@ with sqlite3.connect("parking_citations.db") as conn: ) ``` -For ad-hoc exploration, the SQLite CLI works fine too: +**PostGIS via psql:** ```bash -sqlite3 parking_citations.db -sqlite> SELECT COUNT(*) FROM citations; -sqlite> SELECT * FROM sync_log ORDER BY id DESC LIMIT 5; +docker compose exec db psql -U parking -d parking \ + -c "SELECT COUNT(*), COUNT(geom) FROM citations_clean;" ``` ## Caveats -- **`:updated_at` vs `:created_at`.** The sync orders by `:updated_at DESC`, - so it picks up both new records and edits to existing ones. If a back-dated - correction lands without any genuinely new records, the very first page - will already contain a known `ticket_number` and the loop exits — that's - correct "caught up" behavior, but it means corrections aren't applied - (the table uses `INSERT OR IGNORE`). Switch to `INSERT OR REPLACE` and - remove the early-exit-on-match if you want corrections to flow through. -- **Future-dated rows.** A handful of rows in the dataset have `issue_date` - values years in the future. They're stored as-is; filter them out in your - queries if they're a problem. -- **Geospatial ops.** Lat/long are stored as floats; the WKT column is a - string. For real geometry operations, install `shapely` / `geopandas`, - or migrate the DB to [SpatiaLite](https://www.gaia-gis.it/fossil/libspatialite). +- **Sync uses `INSERT OR IGNORE`** — corrections to existing tickets are not applied. See [ARCHITECTURE.md](ARCHITECTURE.md) for details. +- **Future-dated rows** exist in the source data; filter in queries if needed. +- **SQLite geospatial** stores lat/long as floats and WKT as text. Use PostGIS for native geometry. diff --git a/data-science/beta_pipeline/docker-compose.yml b/data-science/beta_pipeline/docker-compose.yml old mode 100644 new mode 100755 diff --git a/data-science/beta_pipeline/parking_clean.py b/data-science/beta_pipeline/parking_clean.py new file mode 100644 index 00000000..4d6c985c --- /dev/null +++ b/data-science/beta_pipeline/parking_clean.py @@ -0,0 +1,171 @@ +"""Build a cleaned citations table from raw PostGIS ``citations`` rows. + +Reads the full raw table populated by ``parking_postgis`` and writes +``citations_clean`` with a combined issue timestamp, trimmed text fields, +and geometry copied from the source row. + +CLI: + python parking_clean.py init + python parking_clean.py rebuild + python parking_clean.py stats +""" +from __future__ import annotations + +import argparse +import json +import sys +import time +from typing import Any + +from parking_postgis import connect, database_url, init_db + +CLEAN_TABLE = "citations_clean" + +CLEAN_SCHEMA_SQL = f""" +CREATE TABLE IF NOT EXISTS {CLEAN_TABLE} ( + ticket_number TEXT PRIMARY KEY, + issue_datetime TIMESTAMPTZ NOT NULL, + violation_code TEXT, + violation_description TEXT, + fine_amount DOUBLE PRECISION, + geom geometry(Point, 4326) +); +""" + +CLEAN_INDEX_SQL = f""" +CREATE INDEX IF NOT EXISTS idx_{CLEAN_TABLE}_issue_datetime + ON {CLEAN_TABLE} (issue_datetime); +CREATE INDEX IF NOT EXISTS idx_{CLEAN_TABLE}_violation_code + ON {CLEAN_TABLE} (violation_code); +CREATE INDEX IF NOT EXISTS idx_{CLEAN_TABLE}_geom + ON {CLEAN_TABLE} USING GIST (geom); +""" + +# issue_time is HHMM without leading zeros (e.g. "904" -> 09:04, "1430" -> 14:30). +# issue_date is ISO text from the CSV loader (e.g. "2025-04-26T00:00:00.000"). +REBUILD_SQL = f""" +INSERT INTO {CLEAN_TABLE} ( + ticket_number, + issue_datetime, + violation_code, + violation_description, + fine_amount, + geom +) +SELECT + ticket_number, + ( + CAST(issue_date AS TIMESTAMP)::DATE + + CASE + WHEN issue_time IS NOT NULL + AND BTRIM(issue_time) ~ '^\\d{{1,4}}$' + THEN MAKE_INTERVAL( + hours => CAST( + SUBSTRING(LPAD(BTRIM(issue_time), 4, '0') FROM 1 FOR 2) AS INTEGER + ), + mins => CAST( + SUBSTRING(LPAD(BTRIM(issue_time), 4, '0') FROM 3 FOR 2) AS INTEGER + ) + ) + ELSE INTERVAL '0' + END + )::TIMESTAMPTZ AS issue_datetime, + NULLIF(BTRIM(violation_code), '') AS violation_code, + NULLIF(BTRIM(violation_description), '') AS violation_description, + fine_amount, + geom +FROM citations +WHERE ticket_number IS NOT NULL + AND issue_date IS NOT NULL +""" + + +def init_clean(dsn: str | None = None) -> None: + """Ensure raw + cleaned schemas exist (idempotent).""" + init_db(dsn) + with connect(dsn) as conn: + with conn.cursor() as cur: + cur.execute(CLEAN_SCHEMA_SQL) + cur.execute(CLEAN_INDEX_SQL) + conn.commit() + + +def rebuild_clean(dsn: str | None = None, *, progress: bool = True) -> int: + """Replace ``citations_clean`` from ``citations``. Returns row count.""" + init_clean(dsn) + started = time.time() + if progress: + print(f"Rebuilding {CLEAN_TABLE} at {database_url(dsn)} …", flush=True) + + with connect(dsn) as conn: + with conn.cursor() as cur: + cur.execute(f"TRUNCATE {CLEAN_TABLE}") + cur.execute(REBUILD_SQL) + cur.execute(f"SELECT COUNT(*) FROM {CLEAN_TABLE}") + count = cur.fetchone()[0] + conn.commit() + + if progress: + elapsed = time.time() - started + rate = count / elapsed if elapsed else 0 + print( + f" {count:,} rows in {CLEAN_TABLE} " + f"elapsed={elapsed:.1f}s rate={rate:,.0f} rows/s", + flush=True, + ) + return count + + +def clean_stats(dsn: str | None = None) -> dict[str, Any]: + with connect(dsn) as conn: + with conn.cursor() as cur: + cur.execute(f"SELECT COUNT(*) FROM {CLEAN_TABLE}") + count = cur.fetchone()[0] + cur.execute( + f"SELECT COUNT(*) FROM {CLEAN_TABLE} WHERE geom IS NOT NULL" + ) + with_geom = cur.fetchone()[0] + cur.execute(f"SELECT MAX(issue_datetime) FROM {CLEAN_TABLE}") + max_dt = cur.fetchone()[0] + cur.execute(f"SELECT MIN(issue_datetime) FROM {CLEAN_TABLE}") + min_dt = cur.fetchone()[0] + return { + "row_count": count, + "rows_with_geom": with_geom, + "min_issue_datetime": min_dt, + "max_issue_datetime": max_dt, + } + + +def _build_parser() -> argparse.ArgumentParser: + p = argparse.ArgumentParser(description=__doc__.splitlines()[0]) + p.add_argument( + "--database-url", + default=None, + help="Postgres DSN (default: DATABASE_URL env or docker default)", + ) + sub = p.add_subparsers(dest="cmd", required=True) + + sub.add_parser("init", help="Create cleaned-table schema and indexes.") + + sub.add_parser("rebuild", help="Rebuild citations_clean from citations.") + + sub.add_parser("stats", help="Show cleaned-table stats.") + return p + + +def main(argv: list[str] | None = None) -> int: + args = _build_parser().parse_args(argv) + dsn = args.database_url + if args.cmd == "init": + init_clean(dsn) + print(f"Initialized {CLEAN_TABLE} at {database_url(dsn)}") + elif args.cmd == "rebuild": + rebuild_clean(dsn) + elif args.cmd == "stats": + print(json.dumps(clean_stats(dsn), indent=2, default=str)) + return 0 + + +if __name__ == "__main__": + sys.exit(main()) diff --git a/data-science/beta_pipeline/parking_db.py b/data-science/beta_pipeline/parking_db.py old mode 100644 new mode 100755 diff --git a/data-science/beta_pipeline/parking_db_explore.ipynb b/data-science/beta_pipeline/parking_db_explore.ipynb old mode 100644 new mode 100755 diff --git a/data-science/beta_pipeline/parking_pipeline.py b/data-science/beta_pipeline/parking_pipeline.py new file mode 100644 index 00000000..44786e60 --- /dev/null +++ b/data-science/beta_pipeline/parking_pipeline.py @@ -0,0 +1,81 @@ +"""Orchestrate PostGIS load → clean pipeline. + +After the raw ``citations`` table is updated from the city CSV, rebuild the +cleaned ``citations_clean`` table. + +CLI: + python parking_pipeline.py run Parking_Citations_20250811.csv + python parking_pipeline.py clean # skip load, rebuild clean only +""" +from __future__ import annotations + +import argparse +import json +import sys + +from parking_clean import clean_stats, rebuild_clean +from parking_postgis import bulk_load_csv, database_url, db_stats + + +def run_pipeline( + csv_path: str | None = None, + dsn: str | None = None, + *, + batch_size: int = 100_000, + skip_load: bool = False, + progress: bool = True, +) -> dict[str, object]: + """Load raw citations (optional) then rebuild citations_clean.""" + if not skip_load: + if not csv_path: + raise ValueError("csv_path is required unless skip_load=True") + bulk_load_csv(csv_path, dsn, batch_size=batch_size, progress=progress) + rebuild_clean(dsn, progress=progress) + return { + "database_url": database_url(dsn), + "raw": db_stats(dsn), + "clean": clean_stats(dsn), + } + + +def _build_parser() -> argparse.ArgumentParser: + p = argparse.ArgumentParser(description=__doc__.splitlines()[0]) + p.add_argument( + "--database-url", + default=None, + help="Postgres DSN (default: DATABASE_URL env or docker default)", + ) + p.add_argument( + "--batch-size", + type=int, + default=100_000, + help="CSV batch size for load-csv (run subcommand only)", + ) + sub = p.add_subparsers(dest="cmd", required=True) + + p_run = sub.add_parser( + "run", + help="Bulk-load CSV into citations, then rebuild citations_clean.", + ) + p_run.add_argument("csv", help="Path to Parking_Citations CSV") + + sub.add_parser( + "clean", + help="Rebuild citations_clean from existing citations (no CSV load).", + ) + return p + + +def main(argv: list[str] | None = None) -> int: + args = _build_parser().parse_args(argv) + dsn = args.database_url + if args.cmd == "run": + summary = run_pipeline(args.csv, dsn, batch_size=args.batch_size) + else: + summary = run_pipeline(dsn=dsn, skip_load=True) + print(json.dumps(summary, indent=2, default=str)) + return 0 + + +if __name__ == "__main__": + sys.exit(main()) diff --git a/data-science/beta_pipeline/parking_postgis.py b/data-science/beta_pipeline/parking_postgis.py old mode 100644 new mode 100755 index 3e01db21..67e22454 --- a/data-science/beta_pipeline/parking_postgis.py +++ b/data-science/beta_pipeline/parking_postgis.py @@ -226,6 +226,11 @@ def _build_parser() -> argparse.ArgumentParser: p_load = sub.add_parser("load-csv", help="Bulk-load CSV into PostGIS.") p_load.add_argument("csv", help="Path to Parking_Citations CSV") p_load.add_argument("--batch-size", type=int, default=100_000) + p_load.add_argument( + "--clean", + action="store_true", + help="Rebuild citations_clean after load (see parking_pipeline.py).", + ) sub.add_parser("stats", help="Show DB stats.") return p @@ -240,6 +245,10 @@ def main(argv: list[str] | None = None) -> int: elif args.cmd == "load-csv": n = bulk_load_csv(args.csv, dsn, batch_size=args.batch_size) print(f"Loaded {n:,} rows into {database_url(dsn)}") + if args.clean: + from parking_clean import rebuild_clean + + rebuild_clean(dsn) elif args.cmd == "stats": print(json.dumps(db_stats(dsn), indent=2, default=str)) return 0 diff --git a/data-science/beta_pipeline/requirements.txt b/data-science/beta_pipeline/requirements.txt old mode 100644 new mode 100755 From c5b1491f7c5caee0ab5c47dad2126947f7e20206 Mon Sep 17 00:00:00 2001 From: gregpawin Date: Mon, 14 Sep 2026 16:47:36 -0700 Subject: [PATCH 3/4] Add pipeline dependencies and violation analysis notebook Pin the beta_pipeline requirements and add the violation-code analysis notebook used to derive the reference mappings. Ignore .venv/ and raw_data/ at the repo root: the pipeline expects the raw citations CSV to be staged locally and it is multi-GB, so it must never be committable. Co-authored-by: Cursor --- .gitignore | 6 +++ data-science/beta_pipeline/requirements.txt | 7 +++ .../beta_pipeline/violation_analysis.ipynb | 51 +++++++++++++++++++ 3 files changed, 64 insertions(+) create mode 100644 data-science/beta_pipeline/violation_analysis.ipynb diff --git a/.gitignore b/.gitignore index 92ebeaba..440e9039 100644 --- a/.gitignore +++ b/.gitignore @@ -12,6 +12,9 @@ storybook-static/ .pnp.js node_modules/ +# Python virtual environments +.venv/ + # Deployment .vercel/ @@ -38,3 +41,6 @@ coverage/ # TypeScript **/*.tsbuildinfo + +# Local datasets (the raw citations CSV is multi-GB) +raw_data/ diff --git a/data-science/beta_pipeline/requirements.txt b/data-science/beta_pipeline/requirements.txt index af834c2d..265e6000 100755 --- a/data-science/beta_pipeline/requirements.txt +++ b/data-science/beta_pipeline/requirements.txt @@ -1,3 +1,10 @@ polars>=1.40,<2 python-dotenv>=1.0 psycopg[binary]>=3.2,<4 + +# notebooks (violation_analysis.ipynb) +pandas>=2.2 +numpy>=2.0 +matplotlib>=3.9 +seaborn>=0.13 +geopandas>=1.0 diff --git a/data-science/beta_pipeline/violation_analysis.ipynb b/data-science/beta_pipeline/violation_analysis.ipynb new file mode 100644 index 00000000..967d1537 --- /dev/null +++ b/data-science/beta_pipeline/violation_analysis.ipynb @@ -0,0 +1,51 @@ +{ + "cells": [ + { + "cell_type": "code", + "execution_count": 1, + "id": "b24fcc80", + "metadata": {}, + "outputs": [ + { + "ename": "ModuleNotFoundError", + "evalue": "No module named 'pandas'", + "output_type": "error", + "traceback": [ + "\u001b[31m---------------------------------------------------------------------------\u001b[39m", + "\u001b[31mModuleNotFoundError\u001b[39m Traceback (most recent call last)", + "\u001b[36mCell\u001b[39m\u001b[36m \u001b[39m\u001b[32mIn[1]\u001b[39m\u001b[32m, line 1\u001b[39m\n\u001b[32m----> \u001b[39m\u001b[32m1\u001b[39m \u001b[38;5;28;01mimport\u001b[39;00m pandas \u001b[38;5;28;01mas\u001b[39;00m pd\n\u001b[32m 2\u001b[39m \u001b[38;5;28;01mimport\u001b[39;00m numpy \u001b[38;5;28;01mas\u001b[39;00m np\n\u001b[32m 3\u001b[39m \u001b[38;5;28;01mimport\u001b[39;00m matplotlib.pyplot \u001b[38;5;28;01mas\u001b[39;00m plt\n\u001b[32m 4\u001b[39m \u001b[38;5;28;01mimport\u001b[39;00m seaborn \u001b[38;5;28;01mas\u001b[39;00m sns\n", + "\u001b[31mModuleNotFoundError\u001b[39m: No module named 'pandas'" + ] + } + ], + "source": [ + "import pandas as pd\n", + "import numpy as np\n", + "import matplotlib.pyplot as plt\n", + "import seaborn as sns\n", + "import geopandas as gpd" + ] + } + ], + "metadata": { + "kernelspec": { + "display_name": ".venv", + "language": "python", + "name": "python3" + }, + "language_info": { + "codemirror_mode": { + "name": "ipython", + "version": 3 + }, + "file_extension": ".py", + "mimetype": "text/x-python", + "name": "python", + "nbconvert_exporter": "python", + "pygments_lexer": "ipython3", + "version": "3.12.13" + } + }, + "nbformat": 4, + "nbformat_minor": 5 +} From fb7cd21d38925faa5d35c25c111b074dac0f0daf Mon Sep 17 00:00:00 2001 From: gregpawin Date: Mon, 13 Jul 2026 17:25:58 -0700 Subject: [PATCH 4/4] Enhance parking data pipeline by adding drop_incomplete function - Implemented `drop_incomplete` function in `parking_clean.py` to remove rows missing `issue_datetime` or `loc_lat`. - Updated `rebuild_clean` to call `drop_incomplete` after datetime creation. - Modified `parking_pipeline.py` to include dropping incomplete rows in the pipeline process. - Adjusted `parking_postgis.py` to ensure `drop_incomplete` is executed after rebuilding the cleaned table. --- data-science/beta_pipeline/parking_clean.py | 41 +++++++++++++++++-- .../beta_pipeline/parking_pipeline.py | 10 +++-- data-science/beta_pipeline/parking_postgis.py | 3 +- 3 files changed, 47 insertions(+), 7 deletions(-) diff --git a/data-science/beta_pipeline/parking_clean.py b/data-science/beta_pipeline/parking_clean.py index 4d6c985c..15a9d16c 100644 --- a/data-science/beta_pipeline/parking_clean.py +++ b/data-science/beta_pipeline/parking_clean.py @@ -2,7 +2,8 @@ Reads the full raw table populated by ``parking_postgis`` and writes ``citations_clean`` with a combined issue timestamp, trimmed text fields, -and geometry copied from the source row. +and geometry copied from the source row. ``drop_incomplete`` then removes +rows that lack ``issue_datetime`` or source ``loc_lat``. CLI: python parking_clean.py init @@ -79,6 +80,15 @@ AND issue_date IS NOT NULL """ +# Runs after datetime creation: keep only rows that have both a datetime and +# a source latitude (loc_lat on the raw citations row). +DROP_INCOMPLETE_SQL = f""" +DELETE FROM {CLEAN_TABLE} AS c +USING citations AS r +WHERE c.ticket_number = r.ticket_number + AND (c.issue_datetime IS NULL OR r.loc_lat IS NULL) +""" + def init_clean(dsn: str | None = None) -> None: """Ensure raw + cleaned schemas exist (idempotent).""" @@ -90,8 +100,32 @@ def init_clean(dsn: str | None = None) -> None: conn.commit() +def drop_incomplete(dsn: str | None = None, *, progress: bool = True) -> int: + """Remove clean rows missing ``issue_datetime`` or source ``loc_lat``. + + Intended to run after datetime creation in ``rebuild_clean``. + Returns the number of rows deleted. + """ + with connect(dsn) as conn: + with conn.cursor() as cur: + cur.execute(DROP_INCOMPLETE_SQL) + deleted = cur.rowcount if cur.rowcount is not None and cur.rowcount >= 0 else 0 + conn.commit() + + if progress: + print( + f" dropped {deleted:,} rows missing issue_datetime or loc_lat", + flush=True, + ) + return deleted + + def rebuild_clean(dsn: str | None = None, *, progress: bool = True) -> int: - """Replace ``citations_clean`` from ``citations``. Returns row count.""" + """Replace ``citations_clean`` from ``citations`` with combined datetimes. + + Does not filter on ``loc_lat`` — call ``drop_incomplete`` afterward. + Returns row count after the datetime INSERT. + """ init_clean(dsn) started = time.time() if progress: @@ -109,7 +143,7 @@ def rebuild_clean(dsn: str | None = None, *, progress: bool = True) -> int: elapsed = time.time() - started rate = count / elapsed if elapsed else 0 print( - f" {count:,} rows in {CLEAN_TABLE} " + f" {count:,} rows in {CLEAN_TABLE} after datetime build " f"elapsed={elapsed:.1f}s rate={rate:,.0f} rows/s", flush=True, ) @@ -162,6 +196,7 @@ def main(argv: list[str] | None = None) -> int: print(f"Initialized {CLEAN_TABLE} at {database_url(dsn)}") elif args.cmd == "rebuild": rebuild_clean(dsn) + drop_incomplete(dsn) elif args.cmd == "stats": print(json.dumps(clean_stats(dsn), indent=2, default=str)) return 0 diff --git a/data-science/beta_pipeline/parking_pipeline.py b/data-science/beta_pipeline/parking_pipeline.py index 44786e60..9553e99e 100644 --- a/data-science/beta_pipeline/parking_pipeline.py +++ b/data-science/beta_pipeline/parking_pipeline.py @@ -1,7 +1,7 @@ """Orchestrate PostGIS load → clean pipeline. After the raw ``citations`` table is updated from the city CSV, rebuild the -cleaned ``citations_clean`` table. +cleaned ``citations_clean`` table, then drop rows missing datetime or loc_lat. CLI: python parking_pipeline.py run Parking_Citations_20250811.csv @@ -13,7 +13,7 @@ import json import sys -from parking_clean import clean_stats, rebuild_clean +from parking_clean import clean_stats, drop_incomplete, rebuild_clean from parking_postgis import bulk_load_csv, database_url, db_stats @@ -25,16 +25,20 @@ def run_pipeline( skip_load: bool = False, progress: bool = True, ) -> dict[str, object]: - """Load raw citations (optional) then rebuild citations_clean.""" + """Load raw citations (optional), rebuild clean datetimes, drop incomplete.""" if not skip_load: if not csv_path: raise ValueError("csv_path is required unless skip_load=True") bulk_load_csv(csv_path, dsn, batch_size=batch_size, progress=progress) + # 1) Combined issue_datetime from issue_date + issue_time rebuild_clean(dsn, progress=progress) + # 2) Require both datetime and loc_lat + dropped = drop_incomplete(dsn, progress=progress) return { "database_url": database_url(dsn), "raw": db_stats(dsn), "clean": clean_stats(dsn), + "dropped_incomplete": dropped, } diff --git a/data-science/beta_pipeline/parking_postgis.py b/data-science/beta_pipeline/parking_postgis.py index 67e22454..7b560c42 100755 --- a/data-science/beta_pipeline/parking_postgis.py +++ b/data-science/beta_pipeline/parking_postgis.py @@ -246,9 +246,10 @@ def main(argv: list[str] | None = None) -> int: n = bulk_load_csv(args.csv, dsn, batch_size=args.batch_size) print(f"Loaded {n:,} rows into {database_url(dsn)}") if args.clean: - from parking_clean import rebuild_clean + from parking_clean import drop_incomplete, rebuild_clean rebuild_clean(dsn) + drop_incomplete(dsn) elif args.cmd == "stats": print(json.dumps(db_stats(dsn), indent=2, default=str)) return 0