Harden Statbel population import compatibility
GeoIntel CI / docs-smoke (push) Canceled after 0s
GeoIntel CI / contract-smoke (push) Canceled after 0s

This commit is contained in:
Codex
2026-07-16 22:54:21 +02:00
parent 5942712c7b
commit 4f4d467d11
11 changed files with 1533 additions and 18 deletions
+241 -18
View File
@@ -15,6 +15,7 @@ from __future__ import annotations
import argparse
import csv
from hashlib import sha256
import io
import json
import os
@@ -30,9 +31,15 @@ from requests.adapters import HTTPAdapter
from shapely.geometry import mapping, shape
from shapely.ops import transform
from shapely.validation import make_valid
from urllib.parse import unquote, urlsplit
from urllib3.util.retry import Retry
from geographic_scopes import GEOGRAPHIC_SCOPES, GeographicScope
from statbel_population_preflight import (
StatbelPreflightError,
validate_statbel_release,
write_manifest,
)
MUNICIPALITY_NAME = "Mol"
@@ -57,6 +64,7 @@ POPULATION_URLS = {
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",
}
POPULATION_LAYOUTS = {year: ("new" if year == 2025 else "standard") for year in POPULATION_URLS}
def parse_args() -> argparse.Namespace:
@@ -182,16 +190,35 @@ def population_rows(content: bytes, scope: GeographicScope) -> dict[str, dict[st
text = raw.decode("utf-8-sig")
except UnicodeDecodeError:
text = raw.decode("cp1252")
required = {"CD_REFNIS", "CD_SECTOR", "TOTAL", "TX_DESCR_SECTOR_NL", "TX_DESCR_NL"}
reader = csv.DictReader(io.StringIO(text), delimiter="|")
missing_columns = sorted(required - set(reader.fieldnames or ()))
if missing_columns:
raise RuntimeError(f"Statbel population table is missing required columns: {', '.join(missing_columns)}")
return scoped_population_rows(list(reader), scope)
def scoped_population_rows(source_rows: list[dict[str, Any]], scope: GeographicScope) -> dict[str, dict[str, Any]]:
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="|"):
seen: set[str] = set()
for row_number, row in enumerate(source_rows, start=2):
nis_code = str(row.get("CD_REFNIS") or "").strip()
sector_code = str(row.get("CD_SECTOR") or "").strip().upper()
total_raw = str(row.get("TOTAL") if row.get("TOTAL") is not None else "").strip()
if (
len(nis_code) != 5
or not nis_code.isdigit()
or len(sector_code) != 9
or not sector_code[:5].isdigit()
or not total_raw.isdigit()
):
raise RuntimeError(f"Statbel population row {row_number} has invalid code or TOTAL values")
if sector_code in seen:
raise RuntimeError(f"Statbel population table contains duplicate sector {sector_code}")
seen.add(sector_code)
if nis_code not in members:
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"),
@@ -210,6 +237,7 @@ def build_snapshot(
population: dict[str, dict[str, Any]],
boundary,
scope: GeographicScope,
preflight_manifest: dict[str, Any] | None = None,
) -> dict[str, Any]:
transformer = Transformer.from_crs("EPSG:31370", "EPSG:4326", always_xy=True)
member_codes = set(scope.nis_codes)
@@ -245,8 +273,20 @@ def build_snapshot(
"attribution": ATTRIBUTION,
}
features.append({"type": "Feature", "id": sector_code, "geometry": mapping(geometry), "properties": combined})
if missing_population:
raise RuntimeError(f"Statbel geometry has {missing_population} sectors without population rows for {year}")
if not features:
raise RuntimeError(f"No joined population sectors were produced for {year}")
spatial_population_total = sum(int(feature["properties"]["population_total"]) for feature in features)
accounting = (preflight_manifest or {}).get("scope_accounting") or {}
if accounting:
expected_count = int(accounting.get("spatial_sector_count") or 0)
expected_total = int(accounting.get("spatial_population_total") or -1)
if len(features) != expected_count or spatial_population_total != expected_total:
raise RuntimeError(
f"Derived snapshot accounting differs from the passed Statbel preflight for {year}: "
f"features {len(features)}/{expected_count}, population {spatial_population_total}/{expected_total}"
)
return {
"type": "FeatureCollection",
"name": f"Statbel population by statistical sector - {scope.display_name} {year}",
@@ -258,10 +298,145 @@ def build_snapshot(
"geometry_clipped_to_area": True,
"observation_year": year,
"missing_population_sector_count": missing_population,
"spatial_population_total": spatial_population_total,
"unlocated_population_row_count": int(accounting.get("unlocated_row_count") or 0),
"unlocated_population_total": int(accounting.get("unlocated_population_total") or 0),
"accounted_population_total": int(accounting.get("accounted_population_total") or spatial_population_total),
"population_accounting_limitation": (
"Statbel ZZZZ rows cannot be mapped and are excluded from spatial selection metrics."
),
"attribution": ATTRIBUTION,
}
def sha256_path(path: Path) -> str:
digest = sha256()
with path.open("rb") as handle:
for chunk in iter(lambda: handle.read(1024 * 1024), b""):
digest.update(chunk)
return digest.hexdigest()
def write_bytes_atomic(path: Path, content: bytes) -> None:
path.parent.mkdir(parents=True, exist_ok=True)
temporary = path.with_suffix(path.suffix + ".tmp")
temporary.write_bytes(content)
temporary.replace(path)
def write_text_atomic(path: Path, content: str) -> None:
path.parent.mkdir(parents=True, exist_ok=True)
temporary = path.with_suffix(path.suffix + ".tmp")
temporary.write_text(content, encoding="utf-8")
temporary.replace(path)
def snapshot_path(output_dir: Path, scope: GeographicScope, year: int) -> Path:
return output_dir / f"{scope.key.replace('-', '_')}_statbel_population_{year}.geojson"
def preflight_manifest_path(output_dir: Path, scope: GeographicScope, year: int) -> Path:
return output_dir / f"{scope.key.replace('-', '_')}_statbel_population_{year}.preflight.json"
def previous_snapshot_path(output_dir: Path, scope: GeographicScope, year: int) -> Path | None:
candidates = [snapshot_path(output_dir, scope, candidate) for candidate in POPULATION_URLS if candidate < year]
available = [path for path in candidates if path.is_file()]
return max(available, key=lambda path: int(path.stem.rsplit("_", 1)[-1])) if available else None
def load_preflight_manifest(path: Path, snapshot: Path, year: int, scope: GeographicScope) -> dict[str, Any]:
try:
manifest = json.loads(path.read_text(encoding="utf-8"))
except (OSError, UnicodeDecodeError, json.JSONDecodeError) as exc:
raise RuntimeError(f"Statbel preflight manifest is unreadable at {path}") from exc
release = manifest.get("release") or {}
accounting = manifest.get("scope_accounting") or {}
derived = (manifest.get("artifacts") or {}).get("derived_snapshot") or {}
if (
manifest.get("status") != "passed"
or manifest.get("import_eligible") is not True
or int(release.get("year") or 0) != year
or accounting.get("scope_key") != scope.key
or derived.get("sha256") != sha256_path(snapshot)
):
raise RuntimeError(f"Statbel preflight manifest at {path} does not authorize the retained snapshot")
for artifact_name in ("population", "geometry"):
artifact = (manifest.get("artifacts") or {}).get(artifact_name) or {}
retained_path = Path(str(artifact.get("retained_path") or ""))
if (
not retained_path.is_file()
or artifact.get("archive_sha256") != sha256_path(retained_path)
or int(artifact.get("archive_size_bytes") or -1) != retained_path.stat().st_size
):
raise RuntimeError(
f"Statbel preflight manifest at {path} does not authorize the retained {artifact_name} archive"
)
return manifest
def stage_release(
*,
year: int,
population_content: bytes,
geometry_content: bytes,
output_dir: Path,
boundary,
scope: GeographicScope,
) -> tuple[Path, Path, dict[str, Any]]:
layout = POPULATION_LAYOUTS[year]
population_url = POPULATION_URLS[year]
geometry_url = SECTOR_URL.format(year=year)
result = validate_statbel_release(
year=year,
layout=layout,
population_content=population_content,
population_url=population_url,
geometry_content=geometry_content,
geometry_url=geometry_url,
scope=scope,
baseline_snapshot=previous_snapshot_path(output_dir, scope, year),
)
raw_dir = output_dir / "raw" / str(year)
population_archive_path = raw_dir / Path(unquote(urlsplit(population_url).path)).name
geometry_archive_path = raw_dir / Path(unquote(urlsplit(geometry_url).path)).name
write_bytes_atomic(population_archive_path, population_content)
write_bytes_atomic(geometry_archive_path, geometry_content)
manifest = dict(result.manifest)
manifest["artifacts"] = {
**manifest["artifacts"],
"population": {
**manifest["artifacts"]["population"],
"retained_path": str(population_archive_path),
},
"geometry": {
**manifest["artifacts"]["geometry"],
"retained_path": str(geometry_archive_path),
},
}
population = scoped_population_rows(list(result.population.rows.values()), scope)
path = snapshot_path(output_dir, scope, year)
snapshot = build_snapshot(
year,
result.geometry.payload,
population,
boundary,
scope,
preflight_manifest=manifest,
)
write_text_atomic(path, json.dumps(snapshot, ensure_ascii=False, separators=(",", ":")))
manifest["artifacts"]["derived_snapshot"] = {
"retained_path": str(path),
"size_bytes": path.stat().st_size,
"sha256": sha256_path(path),
"feature_count": len(snapshot["features"]),
}
manifest_path = preflight_manifest_path(output_dir, scope, year)
write_manifest(manifest_path, manifest)
return path, manifest_path, manifest
def locate_workspace(
session: requests.Session,
base_url: str,
@@ -292,8 +467,11 @@ def upload_snapshot(
path: Path,
timeout: int,
scope: GeographicScope,
preflight_path: Path,
) -> dict[str, Any]:
observed_at = f"{year}-01-01T00:00:00Z"
preflight = load_preflight_manifest(preflight_path, path, year, scope)
accounting = preflight["scope_accounting"]
source_metadata = {
"provider": "Statbel",
"authority_level": "authoritative",
@@ -309,6 +487,11 @@ def upload_snapshot(
"observation_date_precision": "year",
"identity_stable": False,
"identity_limitation": "Statistical-sector codes and boundaries can change between annual editions.",
"population_layout": preflight["release"]["population_layout"],
"population_accounting": accounting,
"spatial_population_limitation": (
"ZZZZ population rows have no geometry and are excluded from spatial selection metrics."
),
"selection_aggregation": {
"method": "area_weighted_sum",
"property": "population_total",
@@ -325,6 +508,12 @@ def upload_snapshot(
"geometry_clipped_to_area": True,
"sector_geometry_url": SECTOR_URL.format(year=year),
"population_url": POPULATION_URLS[year],
"population_layout": preflight["release"]["population_layout"],
"preflight_manifest_path": str(preflight_path),
"preflight_manifest_sha256": sha256_path(preflight_path),
"population_archive_sha256": preflight["artifacts"]["population"]["archive_sha256"],
"sector_archive_sha256": preflight["artifacts"]["geometry"]["archive_sha256"],
"derived_snapshot_sha256": preflight["artifacts"]["derived_snapshot"]["sha256"],
"generated_at": datetime.now(timezone.utc).isoformat(),
}
with path.open("rb") as handle:
@@ -373,28 +562,52 @@ def main() -> int:
try:
boundary_path = resolve_boundary_path(args, scope)
boundary = load_boundary(boundary_path, scope)
prepared: list[tuple[int, Path, int]] = []
prepared: list[dict[str, Any]] = []
with build_session() as source_session:
for year in years:
path = output_dir / f"{scope.key.replace('-', '_')}_statbel_population_{year}.geojson"
path = snapshot_path(output_dir, scope, year)
manifest_path = preflight_manifest_path(output_dir, scope, year)
preflight_status = "passed"
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, scope),
boundary,
scope,
path, manifest_path, _manifest = stage_release(
year=year,
population_content=population_response.content,
geometry_content=sectors_response.content,
output_dir=output_dir,
boundary=boundary,
scope=scope,
)
path.write_text(json.dumps(snapshot, ensure_ascii=False, separators=(",", ":")), encoding="utf-8")
elif manifest_path.is_file():
load_preflight_manifest(manifest_path, path, year, scope)
else:
preflight_status = "legacy_existing_only"
payload = json.loads(path.read_text(encoding="utf-8"))
prepared.append((year, path, len(payload.get("features") or [])))
prepared.append(
{
"year": year,
"path": path,
"manifest_path": manifest_path if manifest_path.is_file() else None,
"feature_count": len(payload.get("features") or []),
"preflight_status": preflight_status,
}
)
if args.fetch_only:
results = [{"year": year, "path": str(path), "feature_count": count, "status": "prepared"} for year, path, count in prepared]
results = [
{
"year": item["year"],
"path": str(item["path"]),
"preflight_manifest_path": str(item["manifest_path"]) if item["manifest_path"] else None,
"preflight_status": item["preflight_status"],
"feature_count": item["feature_count"],
"status": "prepared" if item["preflight_status"] == "passed" else "legacy_cached",
}
for item in prepared
]
else:
base_url = args.base_url.rstrip("/")
with requests.Session() as api_session:
@@ -405,7 +618,9 @@ def main() -> int:
area_name,
args.import_timeout,
)
for year, path, count in prepared:
for item in prepared:
year = int(item["year"])
path = item["path"]
observed_at = f"{year}-01-01T00:00:00+00:00"
dataset = next(
(
@@ -419,6 +634,10 @@ def main() -> int:
if dataset:
results.append({"year": year, "dataset_id": dataset["id"], "feature_count": dataset.get("feature_count"), "status": "existing"})
continue
if item["preflight_status"] != "passed" or item["manifest_path"] is None:
raise RuntimeError(
f"Statbel {year} has only a legacy cached snapshot; rerun with --force to create preflight evidence before import"
)
dataset = upload_snapshot(
api_session,
base_url,
@@ -428,10 +647,14 @@ def main() -> int:
path,
args.import_timeout,
scope,
item["manifest_path"],
)
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)
payload = {"status": "error", "message": str(exc)}
if isinstance(exc, StatbelPreflightError):
payload.update({"error_code": exc.code, "details": exc.details})
print(json.dumps(payload, ensure_ascii=False), file=sys.stderr)
return 1
print(