spb/countryatlas
Public
TypeScript 57%
Python 38.6%
JavaScript 3.6%
CSS 0.6%
1#!/usr/bin/env python2"""Run one or more connectors end to end (fetch → store_raw → normalize → validate → staging parquet) and print a3per-spec summary: rows, countries, period range, forecast rows, status/message.45 CA_DATA_DIR=~/countryatlas-data .venv/bin/python scripts/run_connector.py who fred bis ilo [--indicator SLUG]6 [--mode fetch|normalize] [--concurrency N]78Uses the pipeline's `run_fetch` (same code path as `ca fetch`), so the staging files land in9`staging/<connector>/<dataset>__<code>__<indicator>.parquet` with the .run.json / .meta.json / .issues.json sidecars.10FRED calls are serialised by the connector itself (≥ 1.1 s spacing), whatever the concurrency.11"""12from __future__ import annotations1314import argparse15import logging16import sys17from pathlib import Path1819sys.path.insert(0, str(Path(__file__).resolve().parents[1] / "src"))2021import polars as pl2223from countryatlas.config import settings24from countryatlas.pipeline.fetch import run_fetch25from countryatlas.pipeline.staging import spec_paths26from countryatlas.registry import source_specs272829def summarize(connector: str, indicator: str | None) -> list[str]:30 lines = []31 for spec in source_specs(connector=connector, indicator=indicator):32 p = spec_paths(spec)["parquet"]33 label = f"{spec.dataset}:{spec.code} → {spec.indicator_id}"34 if not p.exists():35 lines.append(f" ✗ {label}: no staging file")36 continue37 df = pl.read_parquet(p)38 if df.is_empty():39 lines.append(f" ✗ {label}: empty")40 continue41 n_c = df["country_id"].n_unique()42 pmin, pmax = df["period"].min(), df["period"].max()43 n_fc = int(df["is_forecast"].sum())44 n_q = int((df["status"] == "quarantined").sum())45 n_w = int((df["status"] == "warning").sum())46 freq = ",".join(sorted(df["frequency"].unique().to_list()))47 lines.append(48 f" ✓ {label}: {df.height} rows, {n_c} countries, {pmin}…{pmax} [{freq}]"49 + (f", {n_fc} forecast" if n_fc else "")50 + (f", {n_q} quarantined" if n_q else "")51 + (f", {n_w} warnings" if n_w else "")52 )53 return lines545556def main() -> int:57 ap = argparse.ArgumentParser(description=__doc__, formatter_class=argparse.RawDescriptionHelpFormatter)58 ap.add_argument("connectors", nargs="+", help="connector ids (who fred bis ilo …)")59 ap.add_argument("--indicator", default=None)60 ap.add_argument("--mode", choices=["fetch", "normalize"], default="fetch")61 ap.add_argument("--concurrency", type=int, default=None)62 ap.add_argument("-v", "--verbose", action="store_true")63 args = ap.parse_args()64 logging.basicConfig(level=logging.DEBUG if args.verbose else logging.INFO,65 format="%(asctime)s %(levelname)s %(name)s: %(message)s")66 logging.getLogger("httpx").setLevel(logging.WARNING)67 print(f"data dir: {settings.data_dir}")68 summary = run_fetch(connectors=args.connectors, indicator=args.indicator, mode=args.mode, concurrency=args.concurrency)69 print(f"\nrun {summary.run_id}: {len(summary.ok)} ok / {len(summary.runs)} specs in {summary.duration_s:.0f}s")70 for cid in args.connectors:71 print(f"\n[{cid}]")72 for line in summarize(cid, args.indicator):73 print(line)74 for r in summary.runs:75 if r.connector == cid and r.status not in ("ok",):76 print(f" ! {r.status.upper()} {r.dataset}: {r.message}")77 return 0 if not summary.failed else 1787980if __name__ == "__main__":81 raise SystemExit(main())82