← back to Stayclaim
scripts/ingest-la-city-permits.ts
392 lines
/**
* ingest-la-city-permits.ts
*
* Pulls all 1,560,839 LA City Department of Building & Safety building permits
* from the three Socrata datasets on data.lacity.org:
*
* - pi9x-tg5x (2020→present, ~387K rows, daily refresh) ← canonical replacement for stale xnhu-aczu
* - dyxf-7hc4 (2010-2019, ~533K rows, static)
* - e67z-kt2n (<2010, ~640K rows, static)
*
* Verified field shape (consistent across all 3): permit_nbr, primary_address,
* zip_code, cd, apn, zone, permit_group, permit_type, permit_sub_type,
* use_code, use_desc, issue_date, status_desc, valuation, lat, lon, plus
* dataset-specific extras (square_footage on dyxf+e67z, construction on e67z).
*
* For each row:
* 1. Upsert a `listing` row (source='ladbs_permit', source_id=primary_address) —
* one listing per distinct LA-City address, geocoded.
* 2. Insert a `permit` row linked to that listing.
*
* Idempotency: ON CONFLICT (source_dataset, permit_number) DO UPDATE for permits;
* ON CONFLICT (source, source_id) DO UPDATE COALESCE for listings.
*
* Batched 1000 rows per upsert; ~50K-row Socrata pages. Resume-safe via
* --start-page and --dataset flags.
*
* Tier: A (LA City government primary record).
*
* Usage: npx tsx scripts/ingest-la-city-permits.ts
* [--dataset=pi9x-tg5x|dyxf-7hc4|e67z-kt2n|all] default: all
* [--start-page=N] default: 0
* [--page-size=50000] default: 50000
* [--app-token=XXXXXXXX] optional Socrata token
*/
import { Pool } from 'pg';
const pool = new Pool({
host: process.env.PGHOST ?? '/tmp',
database: process.env.PGDATABASE ?? 'stayclaim',
user: process.env.PGUSER ?? process.env.USER,
password: process.env.PGPASSWORD,
port: parseInt(process.env.PGPORT ?? '5432', 10),
max: 8,
});
const args = Object.fromEntries(
process.argv.slice(2).map(a => {
const m = a.match(/^--([^=]+)(?:=(.*))?$/);
return m ? [m[1], m[2] ?? 'true'] : [a, 'true'];
})
);
const ONLY_DATASET = args.dataset as string | undefined;
const START_PAGE = args['start-page'] ? parseInt(args['start-page'] as string, 10) : 0;
const PAGE_SIZE = args['page-size'] ? parseInt(args['page-size'] as string, 10) : 50000;
const APP_TOKEN = (args['app-token'] as string) || process.env.SOCRATA_APP_TOKEN || undefined;
const BATCH_INSERT_SIZE = 1000;
type Dataset = { slug: string; expectedRows: number };
const DATASETS: Dataset[] = [
{ slug: 'pi9x-tg5x', expectedRows: 387531 },
{ slug: 'dyxf-7hc4', expectedRows: 533366 },
{ slug: 'e67z-kt2n', expectedRows: 639942 },
];
type SocrataRow = {
permit_nbr?: string;
primary_address?: string;
zip_code?: string;
cd?: string;
apn?: string;
zone?: string;
permit_group?: string;
permit_type?: string;
permit_sub_type?: string;
use_code?: string;
use_desc?: string;
issue_date?: string;
status_desc?: string;
status_date?: string;
submitted_date?: string;
valuation?: string;
lat?: string;
lon?: string;
square_footage?: string;
work_desc?: string;
cnc?: string;
cpa?: string;
};
async function ensureSchema() {
await pool.query(`
CREATE TABLE IF NOT EXISTS permit (
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
source TEXT NOT NULL DEFAULT 'ladbs',
source_dataset TEXT NOT NULL,
permit_number TEXT NOT NULL,
listing_id UUID REFERENCES listing(id),
primary_address TEXT,
zip_code TEXT,
council_district TEXT,
apn TEXT,
zone TEXT,
permit_group TEXT,
permit_type TEXT,
permit_sub_type TEXT,
use_code TEXT,
use_desc TEXT,
issue_date DATE,
status_date DATE,
submitted_date DATE,
status_desc TEXT,
valuation NUMERIC(15,2),
square_footage INT,
work_desc TEXT,
neighborhood_council TEXT,
community_plan_area TEXT,
latitude NUMERIC(9,6),
longitude NUMERIC(9,6),
source_url TEXT,
source_tier CHAR(1) NOT NULL DEFAULT 'A',
source_label TEXT NOT NULL DEFAULT 'LA City Department of Building & Safety',
retrieved_at TIMESTAMPTZ DEFAULT now(),
UNIQUE (source_dataset, permit_number)
);
CREATE INDEX IF NOT EXISTS idx_permit_listing ON permit(listing_id);
CREATE INDEX IF NOT EXISTS idx_permit_address ON permit(primary_address);
CREATE INDEX IF NOT EXISTS idx_permit_apn ON permit(apn);
CREATE INDEX IF NOT EXISTS idx_permit_zip_date ON permit(zip_code, issue_date);
CREATE INDEX IF NOT EXISTS idx_permit_issue_date ON permit(issue_date);
CREATE INDEX IF NOT EXISTS idx_permit_council_district ON permit(council_district);
`);
}
function parseNum(v?: string): number | null {
if (!v) return null;
const n = parseFloat(v);
return isFinite(n) ? n : null;
}
function parseInt2(v?: string): number | null {
if (!v) return null;
const n = parseInt(v, 10);
return isFinite(n) ? n : null;
}
function parseDate(v?: string): string | null {
if (!v) return null;
const m = v.match(/^(\d{4}-\d{2}-\d{2})/);
return m ? m[1] : null;
}
function buildSourceUrl(slug: string, permitNbr: string): string {
// Direct Socrata row-filter URL: dataset page, filtered to that one permit.
// This lets a user click through and see the canonical record.
return `https://data.lacity.org/resource/${slug}.json?permit_nbr=${encodeURIComponent(permitNbr)}`;
}
/**
* Fetch one Socrata page. Returns rows or empty array if past the end.
*/
async function fetchPage(slug: string, offset: number, limit: number): Promise<SocrataRow[]> {
const url = new URL(`https://data.lacity.org/resource/${slug}.json`);
url.searchParams.set('$limit', String(limit));
url.searchParams.set('$offset', String(offset));
url.searchParams.set('$order', 'permit_nbr');
if (APP_TOKEN) url.searchParams.set('$$app_token', APP_TOKEN);
for (let attempt = 1; attempt <= 5; attempt++) {
try {
const res = await fetch(url, { headers: { 'User-Agent': 'stayclaim-ingest/1.0' } });
if (res.status === 429 || res.status === 503) {
await new Promise(r => setTimeout(r, 2000 * attempt));
continue;
}
if (!res.ok) throw new Error(`HTTP ${res.status} ${res.statusText}`);
return (await res.json()) as SocrataRow[];
} catch (e) {
if (attempt >= 5) throw e;
await new Promise(r => setTimeout(r, 1500 * attempt));
}
}
return [];
}
/**
* Bulk-upsert listings for a batch of (address, lat, lon, zip) tuples.
* Dedupes by primary_address (the source_id).
* Returns a map of primary_address → listing_id.
*/
function canonicalize(addr: string): string {
return addr
.toLowerCase()
.replace(/[^\w\s-]/g, '')
.replace(/\s+/g, '-')
.replace(/-+/g, '-')
.replace(/^-|-$/g, '')
.slice(0, 90);
}
async function upsertListings(rows: SocrataRow[]): Promise<Map<string, string>> {
// Dedup by canonical (slug-equivalent) form so within-batch collisions are
// collapsed before the INSERT. Map original raw address → canonical so
// upsertPermits can look up listing_id by primary_address.
const dedup = new Map<string, { addr: string; lat: number | null; lon: number | null; zip: string | null }>();
const rawToCanonical = new Map<string, string>();
for (const r of rows) {
if (!r.primary_address) continue;
const addr = r.primary_address.trim();
if (!addr) continue;
const canonical = canonicalize(addr);
if (!canonical) continue;
rawToCanonical.set(addr, canonical);
if (dedup.has(canonical)) {
// merge: prefer non-null lat/lon/zip
const ex = dedup.get(canonical)!;
const lat = parseNum(r.lat);
const lon = parseNum(r.lon);
if (ex.lat == null && lat != null) ex.lat = lat;
if (ex.lon == null && lon != null) ex.lon = lon;
if (!ex.zip && r.zip_code) ex.zip = r.zip_code;
continue;
}
dedup.set(canonical, {
addr,
lat: parseNum(r.lat),
lon: parseNum(r.lon),
zip: r.zip_code ?? null,
});
}
// canonical → listing.id
const canonicalToListing = new Map<string, string>();
if (dedup.size === 0) {
return new Map();
}
const arr = Array.from(dedup.entries());
for (let i = 0; i < arr.length; i += BATCH_INSERT_SIZE) {
const slice = arr.slice(i, i + BATCH_INSERT_SIZE);
const tuples: string[] = [];
const values: any[] = [];
slice.forEach(([canonical, r], idx) => {
const base = idx * 7;
tuples.push(`($${base+1},$${base+2},$${base+3},$${base+4},$${base+5},$${base+6},$${base+7})`);
const slug = `${canonical}-ladbs`;
const sourceId = `ladbs:${canonical}`;
values.push(
slug, // 1 slug
sourceId, // 2 source_id
r.addr, // 3 title
r.addr, // 4 address_line1
r.zip, // 5 postal_code
r.lat, // 6 latitude
r.lon, // 7 longitude
);
});
const sql = `
INSERT INTO listing
(slug, source, source_id, title, address_line1, city, state, country, postal_code, latitude, longitude, is_public, tier)
VALUES ${tuples.map((t, i) => {
const base = i * 7;
return `($${base+1},'ladbs_permit',$${base+2},$${base+3},$${base+4},'Los Angeles','CA','US',$${base+5},$${base+6},$${base+7},true,'free')`;
}).join(',')}
ON CONFLICT (source, source_id) DO UPDATE SET
latitude = COALESCE(listing.latitude, EXCLUDED.latitude),
longitude = COALESCE(listing.longitude, EXCLUDED.longitude),
postal_code = COALESCE(listing.postal_code, EXCLUDED.postal_code),
updated_at = now()
RETURNING id, source_id
`;
const r = await pool.query<{ id: string; source_id: string }>(sql, values);
for (const row of r.rows) {
const canonical = row.source_id.replace(/^ladbs:/, '');
canonicalToListing.set(canonical, row.id);
}
}
// Resolve raw → listing via canonical
const result = new Map<string, string>();
for (const [raw, canonical] of rawToCanonical) {
const lid = canonicalToListing.get(canonical);
if (lid) result.set(raw, lid);
}
return result;
}
async function upsertPermits(slug: string, rows: SocrataRow[], listingMap: Map<string, string>): Promise<number> {
const valid = rows.filter(r => r.permit_nbr && r.permit_nbr.trim());
if (valid.length === 0) return 0;
let inserted = 0;
for (let i = 0; i < valid.length; i += BATCH_INSERT_SIZE) {
const slice = valid.slice(i, i + BATCH_INSERT_SIZE);
const tuples: string[] = [];
const values: any[] = [];
slice.forEach((r, idx) => {
const base = idx * 21;
tuples.push(`(${Array.from({length: 21}, (_, k) => `$${base+k+1}`).join(',')})`);
values.push(
slug, // 1 source_dataset
r.permit_nbr!.trim(), // 2 permit_number
r.primary_address ? listingMap.get(r.primary_address.trim()) ?? null : null, // 3 listing_id
r.primary_address ?? null, // 4 primary_address
r.zip_code ?? null, // 5 zip_code
r.cd ?? null, // 6 council_district
r.apn ?? null, // 7 apn
r.zone ?? null, // 8 zone
r.permit_group ?? null, // 9 permit_group
r.permit_type ?? null, // 10 permit_type
r.permit_sub_type ?? null, // 11 permit_sub_type
r.use_code ?? null, // 12 use_code
r.use_desc ?? null, // 13 use_desc
parseDate(r.issue_date), // 14 issue_date
parseDate(r.status_date), // 15 status_date
parseDate(r.submitted_date), // 16 submitted_date
r.status_desc ?? null, // 17 status_desc
parseNum(r.valuation), // 18 valuation
parseInt2(r.square_footage), // 19 square_footage
r.work_desc ?? null, // 20 work_desc
buildSourceUrl(slug, r.permit_nbr!.trim()), // 21 source_url
);
});
const sql = `
INSERT INTO permit
(source_dataset, permit_number, listing_id, primary_address, zip_code, council_district,
apn, zone, permit_group, permit_type, permit_sub_type, use_code, use_desc,
issue_date, status_date, submitted_date, status_desc, valuation, square_footage,
work_desc, source_url)
VALUES ${tuples.join(',')}
ON CONFLICT (source_dataset, permit_number) DO UPDATE SET
listing_id = COALESCE(EXCLUDED.listing_id, permit.listing_id),
status_date = COALESCE(EXCLUDED.status_date, permit.status_date),
status_desc = COALESCE(EXCLUDED.status_desc, permit.status_desc),
retrieved_at = now()
`;
const r = await pool.query(sql, values);
inserted += r.rowCount ?? 0;
}
return inserted;
}
async function ingestDataset(ds: Dataset) {
console.log(`\n=== ${ds.slug} (~${ds.expectedRows.toLocaleString()} rows) ===`);
const t0 = Date.now();
let offset = START_PAGE * PAGE_SIZE;
let pulled = 0, listingsTouched = 0, permitsInserted = 0;
while (true) {
const tFetch = Date.now();
const rows = await fetchPage(ds.slug, offset, PAGE_SIZE);
if (rows.length === 0) break;
const listingMap = await upsertListings(rows);
const inserted = await upsertPermits(ds.slug, rows, listingMap);
pulled += rows.length;
listingsTouched += listingMap.size;
permitsInserted += inserted;
const dt = (Date.now() - t0) / 1000;
const dtFetch = (Date.now() - tFetch) / 1000;
console.log(` page off=${offset} rows=${rows.length} listings=${listingMap.size} permits=${inserted} (page=${dtFetch.toFixed(1)}s, total ${pulled.toLocaleString()}/${ds.expectedRows.toLocaleString()} @ ${(pulled/dt).toFixed(0)}/s)`);
if (rows.length < PAGE_SIZE) break;
offset += PAGE_SIZE;
}
const dt = (Date.now() - t0) / 1000;
console.log(` ✓ ${ds.slug}: pulled=${pulled.toLocaleString()} listings=${listingsTouched.toLocaleString()} permits=${permitsInserted.toLocaleString()} elapsed=${(dt/60).toFixed(1)}m`);
}
async function main() {
await ensureSchema();
const datasets = ONLY_DATASET && ONLY_DATASET !== 'all'
? DATASETS.filter(d => d.slug === ONLY_DATASET)
: DATASETS;
if (datasets.length === 0) {
console.error(`Unknown dataset: ${ONLY_DATASET}`);
process.exit(1);
}
const t0 = Date.now();
for (const ds of datasets) {
await ingestDataset(ds);
}
const dt = (Date.now() - t0) / 1000;
console.log(`\n=== ALL DONE in ${(dt/60).toFixed(1)}m ===`);
const { rows } = await pool.query<{ ds: string; n: string }>(
`SELECT source_dataset as ds, count(*)::text as n FROM permit GROUP BY source_dataset ORDER BY source_dataset`
);
for (const r of rows) console.log(` ${r.ds}: ${parseInt(r.n).toLocaleString()} permits`);
const { rows: lrows } = await pool.query<{ n: string }>(
`SELECT count(*)::text as n FROM listing WHERE source = 'ladbs_permit'`
);
console.log(` listings: ${parseInt(lrows[0].n).toLocaleString()}`);
await pool.end();
}
main().catch(e => {
console.error('FATAL', e);
process.exit(1);
});