589 lines
25 KiB
Python
589 lines
25 KiB
Python
"""Provision official Waterinfo station time series for a persisted Area.
|
|
|
|
This operator discovers annual station series through the public Waterinfo
|
|
KiWIS service, filters station points against the exact persisted Area and
|
|
imports one immutable GeoJSON point dataset per station and observation year.
|
|
It never averages different stations and never interprets a point measurement
|
|
as area-wide water volume.
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
import argparse
|
|
import hashlib
|
|
import json
|
|
import math
|
|
import os
|
|
import re
|
|
import sys
|
|
from dataclasses import dataclass
|
|
from datetime import datetime, timezone
|
|
from pathlib import Path
|
|
from typing import Any
|
|
|
|
import requests
|
|
from requests.adapters import HTTPAdapter
|
|
from shapely.geometry import Point, mapping, shape
|
|
from urllib3.util.retry import Retry
|
|
|
|
|
|
DEFAULT_API_URL = "http://127.0.0.1:8000"
|
|
DEFAULT_PROJECT_NAME = "Kempen Regional Workbench"
|
|
DEFAULT_AREA_NAME = "Gemeente Mol"
|
|
DEFAULT_OUTPUT_DIR = Path("/app/storage/operator-data/waterinfo/mol")
|
|
KIWIS_URL = "https://download.waterinfo.be/tsmdownload/KiWIS/KiWIS"
|
|
CATALOG_URL = "https://waterinfo.vlaanderen.be/"
|
|
ATTRIBUTION = "Bron: Waterinfo Vlaanderen / Vlaamse Milieumaatschappij"
|
|
|
|
|
|
@dataclass(frozen=True)
|
|
class ParameterDefinition:
|
|
key: str
|
|
group_id: str
|
|
property_name: str
|
|
label: str
|
|
unit: str
|
|
reference_layer_name: str
|
|
limitation: str
|
|
|
|
|
|
PARAMETERS = {
|
|
"water_level": ParameterDefinition(
|
|
key="water_level",
|
|
group_id="192784",
|
|
property_name="annual_mean_water_level_m",
|
|
label="Jaargemiddelde waterstand",
|
|
unit="m",
|
|
reference_layer_name="water_level_station",
|
|
limitation=(
|
|
"Dit is een jaargemiddelde op één meetstation. De waarde geldt niet voor het volledige geselecteerde gebied "
|
|
"en levert zonder profiel- of bathymetriegegevens geen watervolume op."
|
|
),
|
|
),
|
|
"discharge": ParameterDefinition(
|
|
key="discharge",
|
|
group_id="192895",
|
|
property_name="annual_mean_discharge_m3_s",
|
|
label="Jaargemiddeld debiet",
|
|
unit="m³/s",
|
|
reference_layer_name="discharge_station",
|
|
limitation=(
|
|
"Dit is een jaargemiddeld debiet op één meetstation. De waarde geldt niet voor alle waterlopen in het "
|
|
"geselecteerde gebied en is geen gebiedsdekkend watervolume."
|
|
),
|
|
),
|
|
}
|
|
|
|
|
|
def parse_args() -> argparse.Namespace:
|
|
parser = argparse.ArgumentParser(description="Provision official annual Waterinfo station histories.")
|
|
parser.add_argument("--base-url", default=os.environ.get("GEOINTEL_INTERNAL_API_URL", DEFAULT_API_URL))
|
|
parser.add_argument("--project-name", default=DEFAULT_PROJECT_NAME)
|
|
parser.add_argument("--area-name", default=DEFAULT_AREA_NAME)
|
|
parser.add_argument("--parameters", default="water_level,discharge")
|
|
parser.add_argument("--from-year", type=int, default=2013)
|
|
parser.add_argument("--to-year", type=int, default=datetime.now(timezone.utc).year)
|
|
parser.add_argument("--min-observations", type=int, default=2)
|
|
parser.add_argument("--max-stations", type=int, default=50)
|
|
parser.add_argument("--output-dir", type=Path, default=Path(os.environ.get("WATERINFO_OUTPUT_DIR", DEFAULT_OUTPUT_DIR)))
|
|
parser.add_argument("--request-timeout", type=int, default=120)
|
|
parser.add_argument("--import-timeout", type=int, default=300)
|
|
parser.add_argument("--force", action="store_true", help="Refetch source artifacts; persisted datasets remain immutable.")
|
|
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-Waterinfo-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 list_paginated_items(session: requests.Session, url: str, *, timeout: int) -> list[dict[str, Any]]:
|
|
items: list[dict[str, Any]] = []
|
|
offset = 0
|
|
total: int | None = None
|
|
while total is None or offset < total:
|
|
page = response_data(session.get(url, params={"limit": 200, "offset": offset}, timeout=timeout))
|
|
page_items = page.get("items") if isinstance(page, dict) else None
|
|
if not isinstance(page_items, list):
|
|
raise RuntimeError(f"GeoIntel list response for {url} has no items array")
|
|
if total is None:
|
|
total = int(page.get("total", len(page_items)))
|
|
items.extend(page_items)
|
|
if not page_items:
|
|
break
|
|
offset += len(page_items)
|
|
if total is not None and len(items) != total:
|
|
raise RuntimeError(f"GeoIntel list response for {url} returned {len(items)} of {total} items")
|
|
return items
|
|
|
|
|
|
def locate_workspace(session: requests.Session, base_url: str, args: argparse.Namespace):
|
|
projects = list_paginated_items(session, f"{base_url}/api/v1/projects", timeout=args.import_timeout)
|
|
project = next((item for item in projects if item.get("name") == args.project_name), None)
|
|
if not project:
|
|
raise RuntimeError(f"Project {args.project_name!r} is missing")
|
|
project_id = str(project["id"])
|
|
areas = list_paginated_items(
|
|
session,
|
|
f"{base_url}/api/v1/projects/{project_id}/areas",
|
|
timeout=args.import_timeout,
|
|
)
|
|
fragment = args.area_name.strip().casefold()
|
|
matches = [item for item in areas if fragment in str(item.get("name") or "").casefold()]
|
|
if len(matches) != 1 or not isinstance(matches[0].get("geometry"), dict):
|
|
raise RuntimeError(f"Expected one persisted Area with geometry matching {args.area_name!r}, received {len(matches)}")
|
|
area = matches[0]
|
|
datasets = list_paginated_items(
|
|
session,
|
|
f"{base_url}/api/v1/projects/{project_id}/datasets",
|
|
timeout=args.import_timeout,
|
|
)
|
|
return project_id, str(area["id"]), shape(area["geometry"]), datasets
|
|
|
|
|
|
def fetch_json(session: requests.Session, params: dict[str, Any], *, timeout: int) -> Any:
|
|
response = session.get(KIWIS_URL, params=params, timeout=timeout)
|
|
response.raise_for_status()
|
|
try:
|
|
return response.json()
|
|
except ValueError as exc:
|
|
raise RuntimeError(f"Waterinfo returned non-JSON: {response.text[:300]}") from exc
|
|
|
|
|
|
def discover_station_series(
|
|
session: requests.Session,
|
|
parameter: ParameterDefinition,
|
|
area_geometry,
|
|
*,
|
|
timeout: int,
|
|
) -> tuple[dict[str, Any], list[dict[str, Any]]]:
|
|
payload = fetch_json(
|
|
session,
|
|
{
|
|
"service": "kisters",
|
|
"type": "queryServices",
|
|
"request": "getTimeseriesValueLayer",
|
|
"datasource": 1,
|
|
"format": "geojson",
|
|
"timeseriesgroup_id": parameter.group_id,
|
|
"metadata": "true",
|
|
"md_returnfields": (
|
|
"custom_attributes,station_id,station_no,station_name,ts_id,ts_name,"
|
|
"stationparameter_name,ts_unitsymbol,parametertype_name"
|
|
),
|
|
"custattr_returnfields": "dataprovider,dataowner",
|
|
"invalidValue": -9999,
|
|
"invalidPeriod": "P2Y",
|
|
},
|
|
timeout=timeout,
|
|
)
|
|
if not isinstance(payload, dict) or payload.get("type") != "FeatureCollection":
|
|
raise RuntimeError(f"Waterinfo group {parameter.group_id} did not return a GeoJSON FeatureCollection")
|
|
selected: list[dict[str, Any]] = []
|
|
for feature in payload.get("features") or []:
|
|
geometry = feature.get("geometry") if isinstance(feature, dict) else None
|
|
properties = feature.get("properties") if isinstance(feature, dict) else None
|
|
if not isinstance(geometry, dict) or geometry.get("type") != "Point" or not isinstance(properties, dict):
|
|
continue
|
|
point = shape(geometry)
|
|
ts_id = properties.get("ts_id")
|
|
if point.is_empty or ts_id in (None, "") or not area_geometry.covers(point):
|
|
continue
|
|
selected.append({"geometry": geometry, "properties": properties, "ts_id": str(ts_id)})
|
|
return payload, selected
|
|
|
|
|
|
def fetch_annual_values(
|
|
session: requests.Session,
|
|
ts_id: str,
|
|
*,
|
|
from_year: int,
|
|
to_year: int,
|
|
timeout: int,
|
|
) -> tuple[Any, dict[int, float]]:
|
|
payload = fetch_json(
|
|
session,
|
|
{
|
|
"service": "kisters",
|
|
"type": "queryServices",
|
|
"request": "getTimeseriesValues",
|
|
"datasource": 1,
|
|
"format": "json",
|
|
"ts_id": ts_id,
|
|
"metadata": "true",
|
|
"md_returnfields": "station_name,station_no,stationparameter_name,ts_id,ts_unitsymbol",
|
|
"from": f"{from_year}-01-01",
|
|
"to": f"{to_year}-12-31",
|
|
},
|
|
timeout=timeout,
|
|
)
|
|
entries = payload if isinstance(payload, list) else [payload]
|
|
values: dict[int, float] = {}
|
|
for entry in entries:
|
|
if not isinstance(entry, dict):
|
|
continue
|
|
for row in entry.get("data") or []:
|
|
if not isinstance(row, list) or len(row) < 2:
|
|
continue
|
|
try:
|
|
year = int(str(row[0])[:4])
|
|
value = float(row[1])
|
|
except (TypeError, ValueError):
|
|
continue
|
|
if year < from_year or year > to_year or not math.isfinite(value) or value <= -9999:
|
|
continue
|
|
if year in values and not math.isclose(values[year], value, rel_tol=0.0, abs_tol=1e-12):
|
|
raise RuntimeError(f"Waterinfo series {ts_id} contains multiple annual values for {year}")
|
|
values[year] = value
|
|
return payload, values
|
|
|
|
|
|
def safe_slug(value: str) -> str:
|
|
normalized = re.sub(r"[^a-z0-9]+", "-", value.casefold()).strip("-")
|
|
return normalized[:80] or "station"
|
|
|
|
|
|
def series_key(parameter: ParameterDefinition, station: dict[str, Any]) -> str:
|
|
properties = station["properties"]
|
|
station_identity = str(properties.get("station_no") or properties.get("station_id") or station["ts_id"])
|
|
return f"waterinfo:{parameter.key}:annual:{safe_slug(station_identity)}"
|
|
|
|
|
|
def build_snapshot(
|
|
parameter: ParameterDefinition,
|
|
station: dict[str, Any],
|
|
year: int,
|
|
value: float,
|
|
) -> dict[str, Any]:
|
|
properties = station["properties"]
|
|
station_name = str(properties.get("station_name") or "Waterinfo meetstation")
|
|
station_no = str(properties.get("station_no") or properties.get("station_id") or station["ts_id"])
|
|
feature_id = f"waterinfo-{parameter.key}-{safe_slug(station_no)}-{year}"
|
|
return {
|
|
"type": "FeatureCollection",
|
|
"name": f"{parameter.label} - {station_name} - {year}",
|
|
"crs": {"type": "name", "properties": {"name": "EPSG:4326"}},
|
|
"features": [
|
|
{
|
|
"type": "Feature",
|
|
"id": feature_id,
|
|
"geometry": station["geometry"],
|
|
"properties": {
|
|
"source_feature_id": feature_id,
|
|
"source_name": "waterinfo",
|
|
"theme": "water",
|
|
"measurement_type": parameter.key,
|
|
"observation_year": year,
|
|
parameter.property_name: value,
|
|
"station_id": properties.get("station_id"),
|
|
"station_no": properties.get("station_no"),
|
|
"station_name": station_name,
|
|
"timeseries_id": station["ts_id"],
|
|
"timeseries_name": properties.get("ts_name"),
|
|
"reported_unit": properties.get("ts_unitsymbol"),
|
|
"data_provider": properties.get("dataprovider"),
|
|
"data_owner": properties.get("dataowner"),
|
|
"authority_level": "authoritative",
|
|
"attribution": ATTRIBUTION,
|
|
},
|
|
}
|
|
],
|
|
}
|
|
|
|
|
|
def sha256_bytes(content: bytes) -> str:
|
|
return hashlib.sha256(content).hexdigest()
|
|
|
|
|
|
def write_json(path: Path, payload: Any, *, pretty: bool = False) -> str:
|
|
content = json.dumps(
|
|
payload,
|
|
ensure_ascii=False,
|
|
indent=2 if pretty else None,
|
|
separators=None if pretty else (",", ":"),
|
|
sort_keys=pretty,
|
|
).encode("utf-8")
|
|
temporary = path.with_suffix(f"{path.suffix}.partial")
|
|
temporary.write_bytes(content)
|
|
temporary.replace(path)
|
|
return sha256_bytes(content)
|
|
|
|
|
|
def upload_snapshot(
|
|
session: requests.Session,
|
|
*,
|
|
base_url: str,
|
|
project_id: str,
|
|
area_id: str,
|
|
parameter: ParameterDefinition,
|
|
station: dict[str, Any],
|
|
year: int,
|
|
path: Path,
|
|
raw_sha256: str,
|
|
coverage: tuple[int, int, int],
|
|
timeout: int,
|
|
) -> dict[str, Any]:
|
|
properties = station["properties"]
|
|
station_name = str(properties.get("station_name") or "Waterinfo meetstation")
|
|
key = series_key(parameter, station)
|
|
observed_at = f"{year}-01-01T00:00:00Z"
|
|
source_metadata = {
|
|
"provider": "Waterinfo Vlaanderen / Vlaamse Milieumaatschappij",
|
|
"authority_level": "authoritative",
|
|
"theme": "water",
|
|
"semantic_metrics": False,
|
|
"geometry_clipped_to_area": True,
|
|
"measurement_type": parameter.key,
|
|
"station_name": station_name,
|
|
"station_no": properties.get("station_no"),
|
|
"timeseries_id": station["ts_id"],
|
|
"attribution": ATTRIBUTION,
|
|
"catalog_url": CATALOG_URL,
|
|
"temporal_series_label": f"{station_name} - {parameter.label.lower()}",
|
|
"observation_date_precision": "year",
|
|
"selection_aggregation": {
|
|
"metric_key": parameter.key,
|
|
"method": "mean",
|
|
"property": parameter.property_name,
|
|
"label": parameter.label,
|
|
"unit": parameter.unit,
|
|
"warning": parameter.limitation,
|
|
},
|
|
}
|
|
provenance_metadata = {
|
|
"operator_tool": "provision_waterinfo_station_history.py",
|
|
"operator_explicit_fetch": True,
|
|
"geometry_clipped_to_area": True,
|
|
"kiwis_url": KIWIS_URL,
|
|
"timeseries_group_id": parameter.group_id,
|
|
"timeseries_id": station["ts_id"],
|
|
"station_id": properties.get("station_id"),
|
|
"station_no": properties.get("station_no"),
|
|
"raw_timeseries_sha256": raw_sha256,
|
|
"coverage_first_year": coverage[0],
|
|
"coverage_last_year": coverage[1],
|
|
"coverage_observation_count": coverage[2],
|
|
"generated_at": datetime.now(timezone.utc).isoformat(),
|
|
"limitation_message": parameter.limitation,
|
|
}
|
|
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": "waterinfo",
|
|
"reference_layer_name": parameter.reference_layer_name,
|
|
"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": key,
|
|
"observed_at": observed_at,
|
|
"valid_from": observed_at,
|
|
"valid_to": f"{year}-12-31T23:59:59Z",
|
|
"temporal_granularity": "year",
|
|
"source_version": f"{station['ts_id']}:{year}",
|
|
},
|
|
files={"file": (path.name, handle, "application/geo+json")},
|
|
timeout=timeout,
|
|
)
|
|
return response_data(response)
|
|
|
|
|
|
def main() -> int:
|
|
args = parse_args()
|
|
if args.from_year > args.to_year or args.min_observations < 2 or args.max_stations < 1:
|
|
print(json.dumps({"status": "error", "message": "Invalid year, observation or station limits"}), file=sys.stderr)
|
|
return 2
|
|
requested_keys = [value.strip() for value in args.parameters.split(",") if value.strip()]
|
|
unsupported = [key for key in requested_keys if key not in PARAMETERS]
|
|
if unsupported or not requested_keys:
|
|
print(json.dumps({"status": "error", "message": f"Unsupported parameters: {unsupported}"}), file=sys.stderr)
|
|
return 2
|
|
|
|
args.output_dir.mkdir(parents=True, exist_ok=True)
|
|
base_url = args.base_url.rstrip("/")
|
|
results: list[dict[str, Any]] = []
|
|
try:
|
|
with requests.Session() as api_session:
|
|
project_id, area_id, area_geometry, existing = locate_workspace(api_session, base_url, args)
|
|
with build_session() as source_session:
|
|
for parameter_key in requested_keys:
|
|
parameter = PARAMETERS[parameter_key]
|
|
station_layer, stations = discover_station_series(
|
|
source_session,
|
|
parameter,
|
|
area_geometry,
|
|
timeout=args.request_timeout,
|
|
)
|
|
layer_path = args.output_dir / f"waterinfo_{parameter.key}_station_layer.json"
|
|
layer_sha256 = write_json(layer_path, station_layer)
|
|
if not stations:
|
|
results.append(
|
|
{
|
|
"parameter": parameter.key,
|
|
"timeseries_group_id": parameter.group_id,
|
|
"status": "no_stations_in_area",
|
|
"observation_count": 0,
|
|
}
|
|
)
|
|
if len(stations) > args.max_stations:
|
|
raise RuntimeError(
|
|
f"Waterinfo returned {len(stations)} in-area {parameter.key} stations; safety limit is {args.max_stations}"
|
|
)
|
|
for station in stations:
|
|
raw_path = args.output_dir / f"waterinfo_{parameter.key}_{station['ts_id']}_annual.json"
|
|
if args.force or not raw_path.exists():
|
|
raw_payload, values = fetch_annual_values(
|
|
source_session,
|
|
station["ts_id"],
|
|
from_year=args.from_year,
|
|
to_year=args.to_year,
|
|
timeout=args.request_timeout,
|
|
)
|
|
raw_sha256 = write_json(raw_path, raw_payload)
|
|
else:
|
|
raw_payload = json.loads(raw_path.read_text(encoding="utf-8"))
|
|
raw_sha256 = sha256_bytes(raw_path.read_bytes())
|
|
entries = raw_payload if isinstance(raw_payload, list) else [raw_payload]
|
|
values = {}
|
|
for entry in entries:
|
|
for row in entry.get("data", []) if isinstance(entry, dict) else []:
|
|
try:
|
|
year = int(str(row[0])[:4])
|
|
value = float(row[1])
|
|
except (IndexError, TypeError, ValueError):
|
|
continue
|
|
if args.from_year <= year <= args.to_year and math.isfinite(value) and value > -9999:
|
|
values[year] = value
|
|
if len(values) < args.min_observations:
|
|
results.append(
|
|
{
|
|
"parameter": parameter.key,
|
|
"timeseries_id": station["ts_id"],
|
|
"status": "insufficient_observations",
|
|
"observation_count": len(values),
|
|
}
|
|
)
|
|
continue
|
|
coverage = (min(values), max(values), len(values))
|
|
key = series_key(parameter, station)
|
|
for year, value in sorted(values.items()):
|
|
path = args.output_dir / f"{safe_slug(key)}_{year}.geojson"
|
|
snapshot = build_snapshot(parameter, station, year, value)
|
|
write_json(path, snapshot)
|
|
existing_dataset = next(
|
|
(
|
|
item
|
|
for item in existing
|
|
if item.get("temporal_series_key") == key
|
|
and str(item.get("observed_at") or "").startswith(str(year))
|
|
),
|
|
None,
|
|
)
|
|
if existing_dataset:
|
|
results.append(
|
|
{
|
|
"parameter": parameter.key,
|
|
"timeseries_id": station["ts_id"],
|
|
"year": year,
|
|
"dataset_id": existing_dataset["id"],
|
|
"status": "existing",
|
|
}
|
|
)
|
|
elif args.fetch_only:
|
|
results.append(
|
|
{
|
|
"parameter": parameter.key,
|
|
"timeseries_id": station["ts_id"],
|
|
"year": year,
|
|
"path": str(path),
|
|
"status": "prepared",
|
|
}
|
|
)
|
|
else:
|
|
dataset = upload_snapshot(
|
|
api_session,
|
|
base_url=base_url,
|
|
project_id=project_id,
|
|
area_id=area_id,
|
|
parameter=parameter,
|
|
station=station,
|
|
year=year,
|
|
path=path,
|
|
raw_sha256=raw_sha256,
|
|
coverage=coverage,
|
|
timeout=args.import_timeout,
|
|
)
|
|
results.append(
|
|
{
|
|
"parameter": parameter.key,
|
|
"timeseries_id": station["ts_id"],
|
|
"year": year,
|
|
"dataset_id": dataset["id"],
|
|
"status": "imported",
|
|
}
|
|
)
|
|
manifest = {
|
|
"schema_version": 1,
|
|
"parameter": parameter.key,
|
|
"timeseries_group_id": parameter.group_id,
|
|
"area_name": args.area_name,
|
|
"station_count": len(stations),
|
|
"station_layer_path": str(layer_path),
|
|
"station_layer_sha256": layer_sha256,
|
|
"from_year": args.from_year,
|
|
"to_year": args.to_year,
|
|
"generated_at": datetime.now(timezone.utc).isoformat(),
|
|
}
|
|
write_json(args.output_dir / f"waterinfo_{parameter.key}_manifest.json", manifest, pretty=True)
|
|
except (OSError, RuntimeError, requests.RequestException, ValueError, KeyError) as exc:
|
|
print(json.dumps({"status": "error", "message": str(exc)}, ensure_ascii=False), file=sys.stderr)
|
|
return 1
|
|
|
|
print(
|
|
json.dumps(
|
|
{
|
|
"status": "ok",
|
|
"project": args.project_name,
|
|
"area": args.area_name,
|
|
"results": results,
|
|
},
|
|
ensure_ascii=False,
|
|
indent=2,
|
|
)
|
|
)
|
|
return 0
|
|
|
|
|
|
if __name__ == "__main__":
|
|
sys.exit(main())
|