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/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
new file mode 100755
index 00000000..d7ebdbcf
--- /dev/null
+++ b/data-science/beta_pipeline/README.md
@@ -0,0 +1,195 @@
+# LA Parking Citations DB
+
+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).
+
+Two storage backends share the same CSV parsing logic:
+
+| 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 |
+
+PostGIS adds a **load → clean pipeline** (`parking_pipeline.py`) that builds a slim
+`citations_clean` table after each full CSV load.
+
+For module-by-module code breakdown, schemas, and data-flow diagrams, see
+**[ARCHITECTURE.md](ARCHITECTURE.md)**.
+
+## Quick start
+
+### 1. Install uv and create the venv
+
+Requires [uv](https://docs.astral.sh/uv/) (manages Python 3.12 and dependencies).
+
+```bash
+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
+```
+
+Polars ≥1.40 needs **Python 3.10+**; the venv uses 3.12.
+
+### 2. Configure environment
+
+```bash
+cp .env.example .env
+```
+
+Edit `.env`:
+
+- `SOCRATA_APP_TOKEN` — optional; raises API rate limits for `parking_db.py sync`
+- `DATABASE_URL` — PostGIS connection string (matches `docker-compose.yml`)
+
+### 3. Choose a path
+
+**SQLite (simplest — no Docker):**
+
+```bash
+.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
+```
+
+**PostGIS (spatial + cleaned table):**
+
+```bash
+docker compose up -d # start PostGIS
+.venv/bin/python parking_pipeline.py run Parking_Citations_20250811.csv
+.venv/bin/python parking_clean.py stats
+```
+
+Replace `Parking_Citations_*.csv` with your actual download filename.
+
+## Project layout
+
+| 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) |
+
+Data files (not in repo): `Parking_Citations_*.csv`, `parking_citations.db`.
+
+## SQLite workflow
+
+### Bulk load (one-time)
+
+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
+```
+
+### Incremental sync (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
+```
+
+Automate with cron or launchd — each run is cheap when nothing is new.
+
+## PostGIS workflow
+
+### Start the database
+
+```bash
+docker compose up -d
+docker compose ps # wait for "healthy"
+```
+
+Default connection: `postgresql://parking:parking@localhost:5432/parking`
+
+### Full pipeline (recommended)
+
+Loads raw rows into `citations`, then rebuilds `citations_clean`:
+
+```bash
+.venv/bin/python parking_pipeline.py run Parking_Citations_20250811.csv
+```
+
+Or step by step:
+
+```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
+```
+
+### Cleaned table columns
+
+| 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 |
+
+Example query:
+
+```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;
+```
+
+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
+
+**SQLite 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,
+ )
+```
+
+**PostGIS via psql:**
+
+```bash
+docker compose exec db psql -U parking -d parking \
+ -c "SELECT COUNT(*), COUNT(geom) FROM citations_clean;"
+```
+
+## Caveats
+
+- **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
new file mode 100755
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_clean.py b/data-science/beta_pipeline/parking_clean.py
new file mode 100644
index 00000000..15a9d16c
--- /dev/null
+++ b/data-science/beta_pipeline/parking_clean.py
@@ -0,0 +1,206 @@
+"""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. ``drop_incomplete`` then removes
+rows that lack ``issue_datetime`` or source ``loc_lat``.
+
+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
+"""
+
+# 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)."""
+ 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 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`` 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:
+ 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} after datetime build "
+ 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)
+ drop_incomplete(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
new file mode 100755
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 100755
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_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 |
|---|
| str | datetime[μs] | str | str | str | str | str | str | str | str | str | str | str | i32 | str | str | f64 | str | str | str | f64 | f64 | str |
| "4602073232" | 2025-04-26 00:00:00 | "904" | null | "0000" | "CA" | "202512" | null | "FORD" | "PA" | "WT" | "1875 20TH ST W" | null | 55 | "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" | null | 53 | "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" | null | 56 | "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" | null | 55 | "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" | null | 53 | "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" | null | 56 | "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" | null | 56 | "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" | null | 56 | "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" | null | 56 | "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" | null | 56 | "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_pipeline.py b/data-science/beta_pipeline/parking_pipeline.py
new file mode 100644
index 00000000..9553e99e
--- /dev/null
+++ b/data-science/beta_pipeline/parking_pipeline.py
@@ -0,0 +1,85 @@
+"""Orchestrate PostGIS load → clean pipeline.
+
+After the raw ``citations`` table is updated from the city CSV, rebuild the
+cleaned ``citations_clean`` table, then drop rows missing datetime or loc_lat.
+
+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, drop_incomplete, 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), 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,
+ }
+
+
+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
new file mode 100755
index 00000000..7b560c42
--- /dev/null
+++ b/data-science/beta_pipeline/parking_postgis.py
@@ -0,0 +1,259 @@
+"""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)
+ 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
+
+
+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)}")
+ if args.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
+
+
+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 100755
index 00000000..265e6000
--- /dev/null
+++ b/data-science/beta_pipeline/requirements.txt
@@ -0,0 +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
+}