feat: synchronize regional official time series
GeoIntel CI / docs-smoke (push) Canceled after 0s
GeoIntel CI / contract-smoke (push) Canceled after 0s

This commit is contained in:
Codex
2026-07-14 22:19:30 +02:00
parent cfd13307e4
commit a5b4bbdf8f
14 changed files with 609 additions and 51 deletions
+153 -38
View File
@@ -1,9 +1,14 @@
"""Provision official annual Statbel population snapshots for Mol.
"""Provision official annual Statbel population snapshots for an approved scope.
The command joins annual population totals to the matching official
statistical-sector geometries, clips the result to Mol and imports each year
statistical-sector geometries, clips the result to an approved geographic
scope and imports each year
through the existing GeoIntel upload API. It never runs during application
startup and it never synthesizes missing population values.
Mol remains the backwards-compatible default. The canonical regional operator
uses ``--scope kempen-transport-region`` and the persisted official scope
boundary produced by ``provision_geographic_scope.py``.
"""
from __future__ import annotations
@@ -27,6 +32,8 @@ from shapely.ops import transform
from shapely.validation import make_valid
from urllib3.util.retry import Retry
from geographic_scopes import GEOGRAPHIC_SCOPES, GeographicScope
MUNICIPALITY_NAME = "Mol"
MUNICIPALITY_NIS_CODE = "13025"
@@ -36,6 +43,9 @@ ATTRIBUTION = "Bron: Statbel, bevolking per statistische sector, CC BY 4.0"
DEFAULT_API_URL = "http://127.0.0.1:8000"
DEFAULT_OUTPUT_DIR = Path("/app/storage/operator-data/mol-population-history")
DEFAULT_BOUNDARY_PATH = Path("/app/storage/operator-data/mol-municipality/mol_municipality_boundary.geojson")
DEFAULT_SCOPE_KEY = "mol"
DEFAULT_SCOPE_OUTPUT_ROOT = Path("/app/storage/operator-data/geographic-scopes")
DEFAULT_REGIONAL_OUTPUT_ROOT = Path("/app/storage/operator-data/official-population")
SECTOR_URL = (
"https://statbel.fgov.be/sites/default/files/files/opendata/Statistische%20sectoren/"
"sh_statbel_statistical_sectors_31370_{year}0101.geojson.zip"
@@ -50,12 +60,19 @@ POPULATION_URLS = {
def parse_args() -> argparse.Namespace:
parser = argparse.ArgumentParser(description="Provision official annual Statbel population snapshots for Mol.")
parser = argparse.ArgumentParser(description="Provision official annual Statbel population snapshots.")
parser.add_argument("--scope", choices=sorted(GEOGRAPHIC_SCOPES), default=DEFAULT_SCOPE_KEY)
parser.add_argument("--base-url", default=os.environ.get("GEOINTEL_INTERNAL_API_URL", DEFAULT_API_URL))
parser.add_argument("--project-name", default=PROJECT_NAME)
parser.add_argument("--project-name", default=None)
parser.add_argument("--area-name", default=None, help="Case-insensitive fragment identifying the persisted Area.")
parser.add_argument("--years", default="2021,2022,2023,2024,2025")
parser.add_argument("--output-dir", type=Path, default=Path(os.environ.get("MOL_POPULATION_OUTPUT_DIR", DEFAULT_OUTPUT_DIR)))
parser.add_argument("--boundary-path", type=Path, default=Path(os.environ.get("MOL_BOUNDARY_PATH", DEFAULT_BOUNDARY_PATH)))
parser.add_argument("--output-dir", type=Path, default=None)
parser.add_argument("--boundary-path", type=Path, default=None)
parser.add_argument(
"--scope-output-root",
type=Path,
default=Path(os.environ.get("GEOINTEL_SCOPE_OUTPUT_ROOT", DEFAULT_SCOPE_OUTPUT_ROOT)),
)
parser.add_argument("--request-timeout", type=int, default=180)
parser.add_argument("--import-timeout", type=int, default=900)
parser.add_argument("--force", action="store_true")
@@ -75,7 +92,7 @@ def build_session() -> requests.Session:
raise_on_status=True,
)
session = requests.Session()
session.headers.update({"User-Agent": "GeoIntel-Mol-Population-Operator/1.0"})
session.headers.update({"User-Agent": "GeoIntel-Official-Population-Operator/1.0"})
adapter = HTTPAdapter(max_retries=retry)
session.mount("https://", adapter)
session.mount("http://", adapter)
@@ -94,18 +111,56 @@ def response_data(response: requests.Response) -> Any:
return payload["data"]
def load_boundary(path: Path):
def series_key(scope: GeographicScope) -> str:
return f"statbel:population-statistical-sector:{scope.key}"
def resolve_output_dir(args: argparse.Namespace, scope: GeographicScope) -> Path:
if args.output_dir is not None:
return args.output_dir
legacy = os.environ.get("MOL_POPULATION_OUTPUT_DIR")
if scope.key == "mol" and legacy:
return Path(legacy)
if scope.key == "mol":
return DEFAULT_OUTPUT_DIR
return DEFAULT_REGIONAL_OUTPUT_ROOT / scope.key
def resolve_boundary_path(args: argparse.Namespace, scope: GeographicScope) -> Path:
if args.boundary_path is not None:
return args.boundary_path
legacy = os.environ.get("MOL_BOUNDARY_PATH")
if scope.key == "mol" and legacy:
return Path(legacy)
if scope.key == "mol" and DEFAULT_BOUNDARY_PATH.exists():
return DEFAULT_BOUNDARY_PATH
scope_dir = args.scope_output_root / scope.key
manifest_path = scope_dir / f"{scope.key.replace('-', '_')}_scope_manifest.json"
if not manifest_path.exists():
raise RuntimeError(
f"Official scope manifest is missing at {manifest_path}; run provision_geographic_scope.py --scope {scope.key} first"
)
manifest = json.loads(manifest_path.read_text(encoding="utf-8"))
if manifest.get("scope_key") != scope.key or manifest.get("status") != "complete":
raise RuntimeError(f"Official scope manifest at {manifest_path} is incomplete or belongs to another scope")
boundary_path = scope_dir / str(manifest.get("boundary_filename") or "")
if not boundary_path.is_file():
raise RuntimeError(f"Official scope boundary referenced by {manifest_path} is missing")
return boundary_path
def load_boundary(path: Path, scope: GeographicScope):
if not path.exists():
raise RuntimeError(f"Mol boundary is missing at {path}; run provision_mol_municipality_workspace.py first")
raise RuntimeError(f"Boundary for {scope.display_name} is missing at {path}")
payload = json.loads(path.read_text(encoding="utf-8"))
features = payload.get("features") or []
if len(features) != 1:
raise RuntimeError("Mol boundary artifact must contain exactly one feature")
raise RuntimeError(f"Boundary artifact for {scope.display_name} must contain exactly one feature")
boundary = shape(features[0]["geometry"])
if not boundary.is_valid:
boundary = make_valid(boundary)
if boundary.is_empty or not boundary.is_valid:
raise RuntimeError("Mol boundary artifact is invalid")
raise RuntimeError(f"Boundary artifact for {scope.display_name} is invalid")
return boundary
@@ -117,7 +172,7 @@ def zip_member_json(content: bytes) -> dict[str, Any]:
return json.loads(archive.read(member).decode("utf-8"))
def population_rows(content: bytes) -> dict[str, dict[str, Any]]:
def population_rows(content: bytes, scope: GeographicScope) -> dict[str, dict[str, Any]]:
with zipfile.ZipFile(io.BytesIO(content)) as archive:
member = next((name for name in archive.namelist() if name.lower().endswith((".txt", ".csv"))), None)
if not member:
@@ -127,9 +182,11 @@ def population_rows(content: bytes) -> dict[str, dict[str, Any]]:
text = raw.decode("utf-8-sig")
except UnicodeDecodeError:
text = raw.decode("cp1252")
members = {member.nis_code: member.name for member in scope.members}
rows: dict[str, dict[str, Any]] = {}
for row in csv.DictReader(io.StringIO(text), delimiter="|"):
if str(row.get("CD_REFNIS") or "").strip() != MUNICIPALITY_NIS_CODE:
nis_code = str(row.get("CD_REFNIS") or "").strip()
if nis_code not in members:
continue
sector_code = str(row.get("CD_SECTOR") or "").strip()
total_raw = str(row.get("TOTAL") or "").strip()
@@ -139,19 +196,28 @@ def population_rows(content: bytes) -> dict[str, dict[str, Any]]:
"population_total": int(total_raw),
"sector_name_nl": row.get("TX_DESCR_SECTOR_NL"),
"municipality_name_nl": row.get("TX_DESCR_NL"),
"municipality": members[nis_code],
"nis_code": nis_code,
}
if not rows:
raise RuntimeError("Statbel population table contains no usable Mol sectors")
raise RuntimeError(f"Statbel population table contains no usable sectors for {scope.display_name}")
return rows
def build_snapshot(year: int, sector_payload: dict[str, Any], population: dict[str, dict[str, Any]], boundary) -> dict[str, Any]:
def build_snapshot(
year: int,
sector_payload: dict[str, Any],
population: dict[str, dict[str, Any]],
boundary,
scope: GeographicScope,
) -> dict[str, Any]:
transformer = Transformer.from_crs("EPSG:31370", "EPSG:4326", always_xy=True)
member_codes = set(scope.nis_codes)
features: list[dict[str, Any]] = []
missing_population = 0
for source_feature in sector_payload.get("features") or []:
properties = source_feature.get("properties") or {}
if str(properties.get("cd_munty_refnis") or "") != MUNICIPALITY_NIS_CODE:
if str(properties.get("cd_munty_refnis") or "") not in member_codes:
continue
sector_code = str(properties.get("cd_sector") or "").strip()
population_values = population.get(sector_code)
@@ -173,8 +239,8 @@ def build_snapshot(year: int, sector_payload: dict[str, Any], population: dict[s
"source_feature_id": sector_code,
"reference_layer_name": "population",
"authority_level": "authoritative",
"municipality": MUNICIPALITY_NAME,
"nis_code": MUNICIPALITY_NIS_CODE,
"municipality": population_values["municipality"],
"nis_code": population_values["nis_code"],
"observation_year": year,
"attribution": ATTRIBUTION,
}
@@ -183,28 +249,37 @@ def build_snapshot(year: int, sector_payload: dict[str, Any], population: dict[s
raise RuntimeError(f"No joined population sectors were produced for {year}")
return {
"type": "FeatureCollection",
"name": f"Statbel population by statistical sector - Mol {year}",
"name": f"Statbel population by statistical sector - {scope.display_name} {year}",
"features": features,
"municipality": MUNICIPALITY_NAME,
"nis_code": MUNICIPALITY_NIS_CODE,
"coverage_scope": scope.key,
"scope_type": scope.scope_type,
"member_count": len(scope.members),
"member_nis_codes": list(scope.nis_codes),
"observation_year": year,
"missing_population_sector_count": missing_population,
"attribution": ATTRIBUTION,
}
def locate_workspace(session: requests.Session, base_url: str, project_name: str, timeout: int):
def locate_workspace(
session: requests.Session,
base_url: str,
project_name: str,
area_name: str,
timeout: int,
):
projects = response_data(session.get(f"{base_url}/api/v1/projects", params={"limit": 200}, timeout=timeout))
project = next((item for item in projects.get("items") or [] if item.get("name") == project_name), None)
if not project:
raise RuntimeError(f"Project {project_name!r} is missing")
project_id = str(project["id"])
areas = response_data(session.get(f"{base_url}/api/v1/projects/{project_id}/areas", params={"limit": 200}, timeout=timeout))
area = next((item for item in areas.get("items") or [] if "gemeente mol" in str(item.get("name", "")).lower()), None)
if not area:
raise RuntimeError("Official Mol area is missing")
area_fragment = area_name.strip().casefold()
matches = [item for item in areas.get("items") or [] if area_fragment in str(item.get("name") or "").casefold()]
if len(matches) != 1:
raise RuntimeError(f"Expected one official Area matching {area_name!r}, received {len(matches)}")
datasets = response_data(session.get(f"{base_url}/api/v1/projects/{project_id}/datasets", params={"limit": 200}, timeout=timeout))
return project_id, str(area["id"]), list(datasets.get("items") or [])
return project_id, str(matches[0]["id"]), list(datasets.get("items") or [])
def upload_snapshot(
@@ -215,16 +290,21 @@ def upload_snapshot(
year: int,
path: Path,
timeout: int,
scope: GeographicScope,
) -> dict[str, Any]:
observed_at = f"{year}-01-01T00:00:00Z"
source_metadata = {
"provider": "Statbel",
"authority_level": "authoritative",
"coverage_scope": "municipality",
"municipality": MUNICIPALITY_NAME,
"nis_code": MUNICIPALITY_NIS_CODE,
"coverage_scope": scope.key,
"scope_type": scope.scope_type,
"scope_display_name": scope.display_name,
"member_count": len(scope.members),
"member_nis_codes": list(scope.nis_codes),
"attribution": ATTRIBUTION,
"license": "CC BY 4.0",
"temporal_series_label": "Officiële bevolkingscijfers per statistische sector",
"observation_date_precision": "year",
"identity_stable": False,
"identity_limitation": "Statistical-sector codes and boundaries can change between annual editions.",
"selection_aggregation": {
@@ -239,6 +319,7 @@ def upload_snapshot(
provenance_metadata = {
"operator_tool": "provision_mol_population_history.py",
"operator_explicit_fetch": True,
"scope_key": scope.key,
"sector_geometry_url": SECTOR_URL.format(year=year),
"population_url": POPULATION_URLS[year],
"generated_at": datetime.now(timezone.utc).isoformat(),
@@ -255,7 +336,7 @@ def upload_snapshot(
"source_metadata_json": json.dumps(source_metadata, ensure_ascii=False),
"provenance_metadata_json": json.dumps(provenance_metadata, ensure_ascii=False),
"area_id": area_id,
"temporal_series_key": SERIES_KEY,
"temporal_series_key": series_key(scope),
"observed_at": observed_at,
"valid_from": observed_at,
"valid_to": f"{year}-12-31T23:59:59Z",
@@ -270,6 +351,9 @@ def upload_snapshot(
def main() -> int:
args = parse_args()
scope = GEOGRAPHIC_SCOPES[args.scope]
project_name = args.project_name or scope.project_name
area_name = args.area_name or scope.area_name
try:
years = sorted({int(value.strip()) for value in args.years.split(",") if value.strip()})
except ValueError:
@@ -280,14 +364,16 @@ def main() -> int:
print(json.dumps({"status": "error", "message": f"Unsupported years: {unsupported}"}), file=sys.stderr)
return 2
args.output_dir.mkdir(parents=True, exist_ok=True)
output_dir = resolve_output_dir(args, scope)
output_dir.mkdir(parents=True, exist_ok=True)
results: list[dict[str, Any]] = []
try:
boundary = load_boundary(args.boundary_path)
boundary_path = resolve_boundary_path(args, scope)
boundary = load_boundary(boundary_path, scope)
prepared: list[tuple[int, Path, int]] = []
with build_session() as source_session:
for year in years:
path = args.output_dir / f"mol_statbel_population_{year}.geojson"
path = output_dir / f"{scope.key.replace('-', '_')}_statbel_population_{year}.geojson"
if args.force or not path.exists():
sectors_response = source_session.get(SECTOR_URL.format(year=year), timeout=args.request_timeout)
sectors_response.raise_for_status()
@@ -296,8 +382,9 @@ def main() -> int:
snapshot = build_snapshot(
year,
zip_member_json(sectors_response.content),
population_rows(population_response.content),
population_rows(population_response.content, scope),
boundary,
scope,
)
path.write_text(json.dumps(snapshot, ensure_ascii=False, separators=(",", ":")), encoding="utf-8")
payload = json.loads(path.read_text(encoding="utf-8"))
@@ -308,14 +395,20 @@ def main() -> int:
else:
base_url = args.base_url.rstrip("/")
with requests.Session() as api_session:
project_id, area_id, existing = locate_workspace(api_session, base_url, args.project_name, args.import_timeout)
project_id, area_id, existing = locate_workspace(
api_session,
base_url,
project_name,
area_name,
args.import_timeout,
)
for year, path, count in prepared:
observed_at = f"{year}-01-01T00:00:00+00:00"
dataset = next(
(
item
for item in existing
if item.get("temporal_series_key") == SERIES_KEY
if item.get("temporal_series_key") == series_key(scope)
and str(item.get("observed_at") or "").startswith(observed_at[:10])
),
None,
@@ -323,13 +416,35 @@ def main() -> int:
if dataset:
results.append({"year": year, "dataset_id": dataset["id"], "feature_count": dataset.get("feature_count"), "status": "existing"})
continue
dataset = upload_snapshot(api_session, base_url, project_id, area_id, year, path, args.import_timeout)
dataset = upload_snapshot(
api_session,
base_url,
project_id,
area_id,
year,
path,
args.import_timeout,
scope,
)
results.append({"year": year, "dataset_id": dataset["id"], "feature_count": dataset.get("feature_count"), "status": "imported"})
except (OSError, RuntimeError, requests.RequestException, ValueError, KeyError, zipfile.BadZipFile) as exc:
print(json.dumps({"status": "error", "message": str(exc)}, ensure_ascii=False), file=sys.stderr)
return 1
print(json.dumps({"status": "ok", "municipality": MUNICIPALITY_NAME, "series": SERIES_KEY, "snapshots": results}, ensure_ascii=False, indent=2))
print(
json.dumps(
{
"status": "ok",
"scope": scope.key,
"display_name": scope.display_name,
"member_count": len(scope.members),
"series": series_key(scope),
"snapshots": results,
},
ensure_ascii=False,
indent=2,
)
)
return 0