← back to Stayclaim
scripts/ingest-la-additional-permits.ts
222 lines
/**
* ingest-la-additional-permits.ts
*
* Pulls 5 additional LADBS permit datasets that weren't in the main building-permit
* ingest. Same Socrata pattern; each row creates/updates a listing and inserts
* a permit row.
*
* - 3f9m-afei Certificate of Occupancy (88K)
* - 67is-svtd Mechanical Permits 2020+ (279K)
* - 5m3t-xjex Mechanical Permits 2010-2019 (503K)
* - ysqd-apz7 Electrical Permits 2020+ (338K)
* - j7mw-thyc Bureau of Engineering Permits (626K)
*
* Tier: A (LA City government primary record).
*
* Usage: npx tsx scripts/ingest-la-additional-permits.ts [--dataset=SLUG]
*/
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;
type Dataset = {
slug: string;
source: string;
source_label: string;
expectedRows: number;
};
// Only datasets sharing the standard LADBS schema (permit_nbr, primary_address, etc.)
// C-of-O (3f9m-afei) and BOE (j7mw-thyc) have different schemas — separate ingesters.
const DATASETS: Dataset[] = [
{ slug: '67is-svtd', source: 'la_city_mechanical', source_label: 'LADBS Mechanical Permits 2020+', expectedRows: 279087 },
{ slug: '5m3t-xjex', source: 'la_city_mechanical', source_label: 'LADBS Mechanical Permits 2010-2019', expectedRows: 503281 },
{ slug: 'ysqd-apz7', source: 'la_city_electrical', source_label: 'LADBS Electrical Permits 2020+', expectedRows: 338674 },
];
const PAGE_SIZE = 50000;
const BATCH = 1000;
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;
valuation?: string;
square_footage?: string;
work_desc?: string;
lat?: string;
lon?: string;
};
function pn(v?: string): number | null { if (!v) return null; const n = parseFloat(v); return isFinite(n) ? n : null; }
function pi(v?: string): number | null { if (!v) return null; const n = parseInt(v, 10); return isFinite(n) ? n : null; }
function pd(v?: string): string | null { if (!v) return null; const m = v.match(/^(\d{4}-\d{2}-\d{2})/); return m ? m[1] : null; }
function canonicalize(addr: string): string {
return addr.toLowerCase().replace(/[^\w\s-]/g, '').replace(/\s+/g, '-').replace(/-+/g, '-').replace(/^-|-$/g, '').slice(0, 90);
}
async function fetchPage(slug: string, offset: number): Promise<SocrataRow[]> {
const u = `https://data.lacity.org/resource/${slug}.json?$limit=${PAGE_SIZE}&$offset=${offset}&$order=permit_nbr`;
for (let attempt = 1; attempt <= 5; attempt++) {
try {
const res = await fetch(u, { 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}`);
return await res.json() as SocrataRow[];
} catch (e) {
if (attempt >= 5) throw e;
await new Promise(r => setTimeout(r, 1500 * attempt));
}
}
return [];
}
async function upsertListings(rows: SocrataRow[], source: string): Promise<Map<string, string>> {
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();
const canon = canonicalize(addr);
if (!canon) continue;
rawToCanonical.set(addr, canon);
if (!dedup.has(canon)) {
dedup.set(canon, { addr, lat: pn(r.lat), lon: pn(r.lon), zip: r.zip_code ?? null });
}
}
const canonToLid = 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) {
const slice = arr.slice(i, i + BATCH);
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})`);
values.push(`${canonical}-ladbs`, `ladbs:${canonical}`, r.addr, r.addr, r.zip, r.lat, r.lon);
});
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((_, k) => {
const b = k * 7;
return `($${b+1},'ladbs_permit',$${b+2},$${b+3},$${b+4},'Los Angeles','CA','US',$${b+5},$${b+6},$${b+7},true,'free')`;
}).join(',')}
ON CONFLICT (slug) 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) canonToLid.set(row.source_id.replace(/^ladbs:/, ''), row.id);
}
const result = new Map<string, string>();
for (const [raw, canonical] of rawToCanonical) {
const lid = canonToLid.get(canonical);
if (lid) result.set(raw, lid);
}
return result;
}
async function upsertPermits(ds: Dataset, rows: SocrataRow[], listingMap: Map<string, string>): Promise<number> {
const seen = new Set<string>();
const valid = rows.filter(r => r.permit_nbr && r.permit_nbr.trim() && !seen.has(r.permit_nbr.trim()) && (seen.add(r.permit_nbr.trim()), true));
if (valid.length === 0) return 0;
let inserted = 0;
for (let i = 0; i < valid.length; i += BATCH) {
const slice = valid.slice(i, i + BATCH);
const tuples: string[] = [];
const values: any[] = [];
slice.forEach((r, idx) => {
const base = idx * 18;
tuples.push(`(${Array.from({length: 18}, (_, k) => `$${base+k+1}`).join(',')})`);
values.push(
ds.source, ds.slug, r.permit_nbr!.trim(),
r.primary_address ? listingMap.get(r.primary_address.trim()) ?? null : null,
r.primary_address ?? null, r.zip_code ?? null, r.cd ?? null, r.apn ?? null, r.zone ?? null,
r.permit_group ?? null, r.permit_type ?? null, r.permit_sub_type ?? null,
r.use_code ?? null, r.use_desc ?? null, pd(r.issue_date), pd(r.status_date),
r.status_desc ?? null, pn(r.valuation),
);
});
await pool.query(
`INSERT INTO permit
(source, 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, status_desc, valuation)
VALUES ${tuples.join(',')}
ON CONFLICT (source_dataset, permit_number) DO UPDATE SET
status_desc = COALESCE(EXCLUDED.status_desc, permit.status_desc),
retrieved_at = now()`,
values
);
inserted += slice.length;
}
// Set source_url + source_label after insert (single update, indexed by source_dataset)
return inserted;
}
async function ingestDataset(ds: Dataset) {
console.log(`\n=== ${ds.slug} ${ds.source_label} (${ds.expectedRows.toLocaleString()}) ===`);
const t0 = Date.now();
let offset = 0, pulled = 0, permitsIns = 0;
while (true) {
const rows = await fetchPage(ds.slug, offset);
if (rows.length === 0) break;
const listingMap = await upsertListings(rows, ds.source);
const ins = await upsertPermits(ds, rows, listingMap);
pulled += rows.length;
permitsIns += ins;
const dt = (Date.now() - t0) / 1000;
console.log(` ${ds.slug}: ${pulled.toLocaleString()}/${ds.expectedRows.toLocaleString()} (${(pulled/dt).toFixed(0)}/s)`);
if (rows.length < PAGE_SIZE) break;
offset += PAGE_SIZE;
}
// Set source_url + label for this dataset
await pool.query(
`UPDATE permit SET source_url = $1, source_label = $2 WHERE source_dataset = $3 AND source_url IS NULL`,
[`https://data.lacity.org/resource/${ds.slug}.json?permit_nbr=`, ds.source_label, ds.slug]
);
// Source URL with permit number suffix needs per-row interpolation
await pool.query(
`UPDATE permit SET source_url = source_url || permit_number WHERE source_dataset = $1 AND source_url LIKE '%permit_nbr=' AND source_url NOT LIKE '%permit_nbr=' || permit_number`,
[ds.slug]
);
console.log(` ✓ ${ds.slug}: ${pulled.toLocaleString()} pulled, ${permitsIns.toLocaleString()} permits inserted`);
}
async function main() {
const datasets = ONLY_DATASET ? DATASETS.filter(d => d.slug === ONLY_DATASET) : DATASETS;
for (const ds of datasets) await ingestDataset(ds);
await pool.end();
}
main().catch(e => { console.error('FATAL', e); process.exit(1); });