SPB Git forge

spb/countryatlas

Public
20commits 1branches 0releases
268.3 MBsize
maindefault branch
12 days agolast push
TypeScript 57% Python 38.6% JavaScript 3.6% CSS 0.6%
3.6 KB · 82 lines python
Raw Blame History
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