338 lines
15 KiB
Python
338 lines
15 KiB
Python
"""Provision official annual Statbel population snapshots for Mol.
|
|
|
|
The command joins annual population totals to the matching official
|
|
statistical-sector geometries, clips the result to Mol and imports each year
|
|
through the existing GeoIntel upload API. It never runs during application
|
|
startup and it never synthesizes missing population values.
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
import argparse
|
|
import csv
|
|
import io
|
|
import json
|
|
import os
|
|
import sys
|
|
import zipfile
|
|
from datetime import datetime, timezone
|
|
from pathlib import Path
|
|
from typing import Any
|
|
|
|
import requests
|
|
from pyproj import Transformer
|
|
from requests.adapters import HTTPAdapter
|
|
from shapely.geometry import mapping, shape
|
|
from shapely.ops import transform
|
|
from shapely.validation import make_valid
|
|
from urllib3.util.retry import Retry
|
|
|
|
|
|
MUNICIPALITY_NAME = "Mol"
|
|
MUNICIPALITY_NIS_CODE = "13025"
|
|
PROJECT_NAME = "Mol Municipality Workbench"
|
|
SERIES_KEY = "statbel:population-statistical-sector:mol"
|
|
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")
|
|
SECTOR_URL = (
|
|
"https://statbel.fgov.be/sites/default/files/files/opendata/Statistische%20sectoren/"
|
|
"sh_statbel_statistical_sectors_31370_{year}0101.geojson.zip"
|
|
)
|
|
POPULATION_URLS = {
|
|
2021: "https://statbel.fgov.be/sites/default/files/files/opendata/bevolking/sectoren/OPENDATA_SECTOREN_2021.zip",
|
|
2022: "https://statbel.fgov.be/sites/default/files/files/opendata/bevolking/sectoren/OPENDATA_SECTOREN_2022.zip",
|
|
2023: "https://statbel.fgov.be/sites/default/files/files/opendata/bevolking/sectoren/OPENDATA_SECTOREN_2023.zip",
|
|
2024: "https://statbel.fgov.be/sites/default/files/files/opendata/bevolking/sectoren/OPENDATA_SECTOREN_2024.zip",
|
|
2025: "https://statbel.fgov.be/sites/default/files/files/opendata/bevolking/sectoren/OPENDATA_SECTOREN_2025_NEW.zip",
|
|
}
|
|
|
|
|
|
def parse_args() -> argparse.Namespace:
|
|
parser = argparse.ArgumentParser(description="Provision official annual Statbel population snapshots for Mol.")
|
|
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("--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("--request-timeout", type=int, default=180)
|
|
parser.add_argument("--import-timeout", type=int, default=900)
|
|
parser.add_argument("--force", action="store_true")
|
|
parser.add_argument("--fetch-only", action="store_true")
|
|
return parser.parse_args()
|
|
|
|
|
|
def build_session() -> requests.Session:
|
|
retry = Retry(
|
|
total=5,
|
|
connect=5,
|
|
read=5,
|
|
status=5,
|
|
backoff_factor=1.0,
|
|
status_forcelist=(429, 500, 502, 503, 504),
|
|
allowed_methods=frozenset({"GET"}),
|
|
raise_on_status=True,
|
|
)
|
|
session = requests.Session()
|
|
session.headers.update({"User-Agent": "GeoIntel-Mol-Population-Operator/1.0"})
|
|
adapter = HTTPAdapter(max_retries=retry)
|
|
session.mount("https://", adapter)
|
|
session.mount("http://", adapter)
|
|
return session
|
|
|
|
|
|
def response_data(response: requests.Response) -> Any:
|
|
try:
|
|
payload = response.json()
|
|
except ValueError as exc:
|
|
raise RuntimeError(f"GeoIntel API returned non-JSON ({response.status_code}): {response.text[:300]}") from exc
|
|
if not response.ok:
|
|
raise RuntimeError(f"GeoIntel API failed ({response.status_code}): {json.dumps(payload, ensure_ascii=False)[:800]}")
|
|
if not isinstance(payload, dict) or "data" not in payload:
|
|
raise RuntimeError("GeoIntel API response does not use the canonical data envelope")
|
|
return payload["data"]
|
|
|
|
|
|
def load_boundary(path: Path):
|
|
if not path.exists():
|
|
raise RuntimeError(f"Mol boundary is missing at {path}; run provision_mol_municipality_workspace.py first")
|
|
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")
|
|
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")
|
|
return boundary
|
|
|
|
|
|
def zip_member_json(content: bytes) -> dict[str, Any]:
|
|
with zipfile.ZipFile(io.BytesIO(content)) as archive:
|
|
member = next((name for name in archive.namelist() if name.lower().endswith(".geojson")), None)
|
|
if not member:
|
|
raise RuntimeError("Statbel sector archive contains no GeoJSON file")
|
|
return json.loads(archive.read(member).decode("utf-8"))
|
|
|
|
|
|
def population_rows(content: bytes) -> 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:
|
|
raise RuntimeError("Statbel population archive contains no text table")
|
|
raw = archive.read(member)
|
|
try:
|
|
text = raw.decode("utf-8-sig")
|
|
except UnicodeDecodeError:
|
|
text = raw.decode("cp1252")
|
|
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:
|
|
continue
|
|
sector_code = str(row.get("CD_SECTOR") or "").strip()
|
|
total_raw = str(row.get("TOTAL") or "").strip()
|
|
if not sector_code or not total_raw or not total_raw.isdigit():
|
|
continue
|
|
rows[sector_code] = {
|
|
"population_total": int(total_raw),
|
|
"sector_name_nl": row.get("TX_DESCR_SECTOR_NL"),
|
|
"municipality_name_nl": row.get("TX_DESCR_NL"),
|
|
}
|
|
if not rows:
|
|
raise RuntimeError("Statbel population table contains no usable Mol sectors")
|
|
return rows
|
|
|
|
|
|
def build_snapshot(year: int, sector_payload: dict[str, Any], population: dict[str, dict[str, Any]], boundary) -> dict[str, Any]:
|
|
transformer = Transformer.from_crs("EPSG:31370", "EPSG:4326", always_xy=True)
|
|
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:
|
|
continue
|
|
sector_code = str(properties.get("cd_sector") or "").strip()
|
|
population_values = population.get(sector_code)
|
|
if not population_values:
|
|
missing_population += 1
|
|
continue
|
|
geometry = transform(transformer.transform, shape(source_feature["geometry"]))
|
|
if not geometry.is_valid:
|
|
geometry = make_valid(geometry)
|
|
geometry = geometry.intersection(boundary)
|
|
if geometry.is_empty:
|
|
continue
|
|
if not geometry.is_valid:
|
|
geometry = make_valid(geometry)
|
|
combined = {
|
|
**properties,
|
|
**population_values,
|
|
"source_name": "statbel",
|
|
"source_feature_id": sector_code,
|
|
"reference_layer_name": "population",
|
|
"authority_level": "authoritative",
|
|
"municipality": MUNICIPALITY_NAME,
|
|
"nis_code": MUNICIPALITY_NIS_CODE,
|
|
"observation_year": year,
|
|
"attribution": ATTRIBUTION,
|
|
}
|
|
features.append({"type": "Feature", "id": sector_code, "geometry": mapping(geometry), "properties": combined})
|
|
if not features:
|
|
raise RuntimeError(f"No joined population sectors were produced for {year}")
|
|
return {
|
|
"type": "FeatureCollection",
|
|
"name": f"Statbel population by statistical sector - Mol {year}",
|
|
"features": features,
|
|
"municipality": MUNICIPALITY_NAME,
|
|
"nis_code": MUNICIPALITY_NIS_CODE,
|
|
"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):
|
|
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")
|
|
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 [])
|
|
|
|
|
|
def upload_snapshot(
|
|
session: requests.Session,
|
|
base_url: str,
|
|
project_id: str,
|
|
area_id: str,
|
|
year: int,
|
|
path: Path,
|
|
timeout: int,
|
|
) -> 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,
|
|
"attribution": ATTRIBUTION,
|
|
"license": "CC BY 4.0",
|
|
"identity_stable": False,
|
|
"identity_limitation": "Statistical-sector codes and boundaries can change between annual editions.",
|
|
"selection_aggregation": {
|
|
"method": "area_weighted_sum",
|
|
"property": "population_total",
|
|
"label": "Inwoners",
|
|
"unit": "inwoners",
|
|
"warning_only_when_estimate": True,
|
|
"warning": "Bevolking binnen een gedeeltelijke statistische sector is oppervlaktegewogen en blijft een schatting.",
|
|
},
|
|
}
|
|
provenance_metadata = {
|
|
"operator_tool": "provision_mol_population_history.py",
|
|
"operator_explicit_fetch": True,
|
|
"sector_geometry_url": SECTOR_URL.format(year=year),
|
|
"population_url": POPULATION_URLS[year],
|
|
"generated_at": datetime.now(timezone.utc).isoformat(),
|
|
}
|
|
with path.open("rb") as handle:
|
|
response = session.post(
|
|
f"{base_url}/api/v1/projects/{project_id}/datasets/upload",
|
|
data={
|
|
"dataset_type": "vector",
|
|
"source": "operator_official_import",
|
|
"dataset_role": "reference",
|
|
"source_name": "statbel",
|
|
"reference_layer_name": "population",
|
|
"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,
|
|
"observed_at": observed_at,
|
|
"valid_from": observed_at,
|
|
"valid_to": f"{year}-12-31T23:59:59Z",
|
|
"temporal_granularity": "year",
|
|
"source_version": str(year),
|
|
},
|
|
files={"file": (path.name, handle, "application/geo+json")},
|
|
timeout=timeout,
|
|
)
|
|
return response_data(response)
|
|
|
|
|
|
def main() -> int:
|
|
args = parse_args()
|
|
try:
|
|
years = sorted({int(value.strip()) for value in args.years.split(",") if value.strip()})
|
|
except ValueError:
|
|
print(json.dumps({"status": "error", "message": "Years must be comma-separated integers"}), file=sys.stderr)
|
|
return 2
|
|
unsupported = [year for year in years if year not in POPULATION_URLS]
|
|
if unsupported or not years:
|
|
print(json.dumps({"status": "error", "message": f"Unsupported years: {unsupported}"}), file=sys.stderr)
|
|
return 2
|
|
|
|
args.output_dir.mkdir(parents=True, exist_ok=True)
|
|
results: list[dict[str, Any]] = []
|
|
try:
|
|
boundary = load_boundary(args.boundary_path)
|
|
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"
|
|
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()
|
|
population_response = source_session.get(POPULATION_URLS[year], timeout=args.request_timeout)
|
|
population_response.raise_for_status()
|
|
snapshot = build_snapshot(
|
|
year,
|
|
zip_member_json(sectors_response.content),
|
|
population_rows(population_response.content),
|
|
boundary,
|
|
)
|
|
path.write_text(json.dumps(snapshot, ensure_ascii=False, separators=(",", ":")), encoding="utf-8")
|
|
payload = json.loads(path.read_text(encoding="utf-8"))
|
|
prepared.append((year, path, len(payload.get("features") or [])))
|
|
|
|
if args.fetch_only:
|
|
results = [{"year": year, "path": str(path), "feature_count": count, "status": "prepared"} for year, path, count in prepared]
|
|
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)
|
|
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
|
|
and str(item.get("observed_at") or "").startswith(observed_at[:10])
|
|
),
|
|
None,
|
|
)
|
|
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)
|
|
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))
|
|
return 0
|
|
|
|
|
|
if __name__ == "__main__":
|
|
sys.exit(main())
|