← back to La Socrata Ingester
src/ingest.js
77 lines
import { upsert, getState, setState, pool } from './db.js';
import { socrataPages } from './adapters/socrata.js';
import { arcgisPages } from './adapters/arcgis.js';
import { sleep } from './adapters/http.js';
import { SOURCES, resolveMap } from './sources.js';
// Tables that carry a dataset_id column (so we stamp it onto each row).
const HAS_DATASET_ID = new Set(['la_building_permits_raw', 'la_code_enforcement_raw']);
const log = (...a) => console.log(new Date().toISOString(), ...a);
// Run one source end-to-end.
// opts: { full, maxRows, pageSize }
export async function ingestSource(name, opts = {}) {
const src = SOURCES[name];
if (!src) throw new Error(`unknown source: ${name}`);
const map = resolveMap(name);
const datasetLabel = src.datasetId || name;
const appToken = process.env.SOCRATA_APP_TOKEN || undefined;
const pageDelay = Number(process.env.INGEST_PAGE_DELAY_MS || 120);
const prev = await getState(name);
const since = !opts.full && src.cursorField ? prev?.last_cursor || null : null;
const pagerOpts = { since, appToken, maxRows: opts.maxRows ?? Infinity, fullScan: !!opts.full };
if (opts.pageSize) pagerOpts.pageSize = opts.pageSize;
// --year=YYYY (ArcGIS only): pull exactly one roll year, resumable/idempotent.
if (opts.year && src.platform === 'arcgis' && src.cursorField) {
pagerOpts.whereOverride = `${src.cursorField}='${opts.year}'`;
}
const pager = src.platform === 'socrata' ? socrataPages(src, pagerOpts) : arcgisPages(src, pagerOpts);
const mode = opts.year ? `year=${opts.year}` : opts.full ? 'FULL' : since ? `since ${since}` : 'initial';
log(`▶ ${name} [${src.platform} ${datasetLabel}] ${mode}${opts.maxRows ? ` (max ${opts.maxRows})` : ''}`);
let total = 0;
let maxCursor = prev?.last_cursor || null;
let status = 'ok';
try {
for await (const { rows, url } of pager) {
const mapped = rows.map((r) => {
const row = { ...map(r), raw: r, source_url: url };
if (HAS_DATASET_ID.has(src.table)) row.dataset_id = datasetLabel;
return row;
});
// drop rows with a null primary key component (can't upsert those)
const clean = mapped.filter((row) => src.conflict.every((k) => row[k] != null));
const skipped = mapped.length - clean.length;
await upsert(src.table, src.conflict, clean);
total += clean.length;
// advance high-water mark
if (src.cursorField) {
for (const r of rows) {
const c = r[src.cursorField];
if (c != null && (maxCursor == null || String(c) > maxCursor)) maxCursor = String(c);
}
}
log(` +${clean.length}${skipped ? ` (skipped ${skipped} keyless)` : ''} total=${total}`);
if (pageDelay) await sleep(pageDelay);
}
} catch (err) {
status = 'error: ' + err.message;
log(`✖ ${name} failed: ${err.message}`);
await setState(name, { dataset_id: datasetLabel, last_cursor: maxCursor, rows_upserted: total, last_status: status });
throw err;
}
await setState(name, { dataset_id: datasetLabel, last_cursor: maxCursor, rows_upserted: total, last_status: status });
log(`✔ ${name} done — ${total} rows upserted this run${maxCursor ? `, cursor→${maxCursor}` : ''}`);
return { name, total, maxCursor };
}
export async function closePool() {
await pool.end();
}