gendesign/backend/app/services/site_finder/rosseti_wfs_loader.py
bot-backend b35f444d50
All checks were successful
CI Trade-In / changes (pull_request) Successful in 9s
CI Trade-In / backend-tests (pull_request) Has been skipped
CI Trade-In / browser-tests (pull_request) Has been skipped
CI / changes (pull_request) Successful in 11s
CI Trade-In / frontend-checks (pull_request) Has been skipped
CI / frontend-tests (pull_request) Has been skipped
CI / openapi-codegen-check (pull_request) Successful in 2m22s
CI / backend-tests (pull_request) Successful in 17m46s
fix(rosseti): стабильный external_id ЦП вместо сессионного fid GeoServer
feature['id'] в WFS-ответе — сессионный fid, новый на каждый GetFeature.
ON CONFLICT (source, external_id) не срабатывал ни разу, каждый weekly-прогон
дописывал полный комплект ~488 фич: 4880 строк на 481 ЦП. Задуманный sha1-фолбэк
по атрибутам был мёртв — fid присутствует всегда.

Ключ теперь считается ТОЛЬКО по стабильным атрибутам (нормализованное имя, класс
напряжения, координаты в 1e-5 градуса), fid игнорируется. Координата квантуется
в целое, а не форматируется как float: ключ обязан совпадать байт-в-байт с
бэкфиллом в SQL. sha256 вместо sha1 — встроен в PG16, pgcrypto не нужен.

Починка разбора старые строки не убирает (ключи не совпадут, ON CONFLICT ничего
не перезапишет) → data/sql/99c_power_supply_centers_dedup.sql: пересчёт ключа
существующих строк + схлопывание копий (победитель — свежайший snapshot,
NULLS LAST явно), с печатью чисел до/после и идемпотентностью.

Refs #3322
2026-09-02 14:43:58 +05:00

294 lines
14 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

"""Загрузчик центров питания (ЦП) Россети из WFS — свободная мощность §3 (#2119).
No-B2B верифицированный источник: WFS геосервера портала-тп.рф отдаёт ~488
центров питания (ПС 35/110 кВ) по Свердловской области с координатами (WGS84,
округлены ~100 м), классом напряжения и индексом загрузки (открыт/ограничен/
закрыт для техприсоединения). UPSERT-ит в ``power_supply_centers``.
Источник ГЕО-БЛОКИРУЕТ non-RU IP → загрузчик РАБОТАЕТ НА ПРОДЕ (Celery weekly /
manual docker exec). Punycode-host ОБЯЗАТЕЛЕН — кириллический алиас портал-тп.рф
не резолвится корректно из curl/httpx.
Структура зеркалит utility_infrastructure_loader: httpx с явным таймаутом,
per-row SAVEPOINT при UPSERT (битая фича не валит weekly-sync).
Также экспортирует ``normalize_sc_name`` — общий нормализатор имён ЦП, которым
пользуется и ``rosseti_reserve_loader`` (матч xlsx-резерва по имени), и endpoint.
"""
import hashlib
import json
import logging
import math
import re
import httpx
from sqlalchemy import text
from sqlalchemy.orm import Session
from app.core.db import SessionLocal
logger = logging.getLogger(__name__)
# Punycode-host ОБЯЗАТЕЛЕН (кириллический алиас портал-тп.рф не для curl/httpx).
# xn----7sb7akeedqd.xn--p1ai == портал-тп.рф.
WFS_BASE_URL = "https://xn----7sb7akeedqd.xn--p1ai/geoserver/wfs"
# Слой центров питания (полный вид с резервами/напряжением).
_WFS_TYPE_NAME = "gisosslabs:sc_points_fullview"
# CQL-фильтр по региону — только Свердловская область.
_WFS_CQL_FILTER = "rg_code='SverdlovskOblast'"
# Таймаут WFS-запроса (сек). ~488 фич — с запасом.
_WFS_TIMEOUT = 60
# sc_indexload_id (справочник Россети) → индекс загрузки ЦП. Неизвестный id →
# None (сырое значение остаётся в raw). Верифицированные id из справочника:
_INDEXLOAD_MAP: dict[int, str] = {
256184: "open", # открыт для техприсоединения (есть резерв)
256186: "limited", # ограниченно
256188: "closed", # закрыт (резерва нет)
}
# «ПС»/«подстанция» + класс напряжения-префикс — срезаем при нормализации имени.
# Оба компонента опциональны, но хотя бы один должен присутствовать (иначе матч
# пустой и sub ничего не делает). Примеры:
# «ПС 110/35/10 Нижне-Исетская» → «нижне-исетская»
# «ПС Уктус» → «уктус» (только «ПС», без напряжения)
# «110/10 Уктус» → «уктус» (только напряжение)
_VOLTAGE_PREFIX_RE = re.compile(
r"^\s*(?:(?:пс|подстанци\w*)\s*)?(?:[\d/.,]+\s*(?:кв)?\s*)?",
re.IGNORECASE,
)
# Класс напряжения внутри имени ЦП, напр. «110/35/10» или «110/10».
_VOLTAGE_CLASS_RE = re.compile(r"\b(\d{1,3}(?:[/.]\d{1,3}){1,3})\b")
_MULTISPACE_RE = re.compile(r"\s+")
def normalize_sc_name(name: str | None) -> str:
"""Нормализует имя ЦП для матча xlsx-резерва с WFS-фичей.
Правила: lower, ё→е, срез «ПС»/подстанция + класс напряжения-префикса,
удаление кавычек, схлопывание тире/пробелов. Пусто → ''.
Examples:
«ПС 110/35/10 Нижне-Исетская» → «нижне-исетская»
«Нижне — Исетская» → «нижне-исетская»
«ПС "Уктус"» → «уктус»
"""
if not name:
return ""
s = name.strip().lower().replace("ё", "е")
s = s.replace("«", "").replace("»", "").replace('"', "").replace("'", "")
# Срез voltage/«ПС»-префикса в начале имени (не трогаем цифры внутри имени).
s = _VOLTAGE_PREFIX_RE.sub("", s, count=1)
# Схлопываем все виды тире (—, , -) в одиночный дефис, пробелы вокруг него убираем.
s = re.sub(r"\s*[—–-]\s*", "-", s)
s = _MULTISPACE_RE.sub(" ", s).strip()
return s
def parse_voltage_class(name: str | None) -> str | None:
"""Извлекает класс напряжения из текста имени ЦП («ПС 110/35/10 …» → '110/35/10').
Возвращает первую последовательность вида «NNN/NNN[/NNN…]». None — если нет.
"""
if not name:
return None
m = _VOLTAGE_CLASS_RE.search(name)
if not m:
return None
return m.group(1).replace(".", "/")
def _coord_e5(value: float | None) -> str:
"""Координата → целое в единицах 1e-5 градуса (~1 м), полукруглением от нуля.
Целое, а не форматированный float: ключ обязан совпадать байт-в-байт с SQL-
бэкфиллом (99c), а текстовое представление double в питоне и в PG разное.
``floor(|v|*1e5 + 0.5)`` со знаком = ``round(numeric)`` в PG (half-away-from-zero).
"""
if value is None:
return ""
n = math.floor(abs(value) * 100000 + 0.5)
return str(-n if value < 0 else n)
def _stable_external_id(feature: dict, props: dict) -> str:
"""Стабильный external_id фичи — хэш атрибутов. ``feature['id']`` ИГНОРИРУЕТСЯ.
GeoServer отдаёт СЕССИОННЫЙ fid (``sc_points_fullview.fid--<random>``), новый на
каждый GetFeature → ON CONFLICT (source, external_id) не срабатывал ни разу и
таблица росла ×10 (4880 строк на 481 ЦП, #3322). Ключ считаем только по стабильным
атрибутам: нормализованное имя | класс напряжения | координаты в 1e-5 градуса.
ФОРМУЛА ПРОДУБЛИРОВАНА в ``data/sql/99c_power_supply_centers_dedup.sql`` (бэкфилл
существующих строк) — менять только синхронно. sha256, а не sha1: sha256 встроен
в PG16, sha1 требует pgcrypto. Префикс ``h:`` отличает новый ключ от старого fid.
"""
geom_pair = _point_geom_sql(feature)
coords = geom_pair[1] if geom_pair else {}
sc_name = props.get("sc_name")
seed = "|".join(
(
normalize_sc_name(sc_name),
parse_voltage_class(sc_name) or "",
_coord_e5(coords.get("lon")),
_coord_e5(coords.get("lat")),
)
)
# sha256 здесь — стабильный дедуп-id фичи, не криптография.
return "h:" + hashlib.sha256(seed.encode("utf-8")).hexdigest()[:16]
def _map_load_index(props: dict) -> str | None:
"""sc_indexload_id → 'open'|'limited'|'closed'|None (неизвестный → None)."""
raw = props.get("sc_indexload_id")
if raw is None:
return None
try:
return _INDEXLOAD_MAP.get(int(raw))
except (TypeError, ValueError):
return None
def fetch_power_supply_centers() -> list[dict]:
"""Тянет WFS-фичи центров питания Свердловской области (GeoJSON Feature-list).
RUN-ON-PROD: источник гео-блокирует non-RU IP. httpx с явным таймаутом,
verify по умолчанию (валидный TLS у геосервера). Возвращает list feature-dict'ов
``{"id", "geometry": {...}, "properties": {...}}``.
"""
params = {
"service": "WFS",
"version": "1.1.0",
"request": "GetFeature",
"typeName": _WFS_TYPE_NAME,
"outputFormat": "json",
"CQL_FILTER": _WFS_CQL_FILTER,
}
resp = httpx.get(WFS_BASE_URL, params=params, timeout=_WFS_TIMEOUT)
resp.raise_for_status()
data = resp.json()
features: list[dict] = (data or {}).get("features") or []
logger.info("rosseti_wfs: загружено центров питания=%d", len(features))
return features
def _point_geom_sql(feature: dict) -> tuple[str, dict] | None:
"""SQL-выражение + params для geom из GeoJSON Point. None — если не Point/нет коорд."""
geom = feature.get("geometry") or {}
if geom.get("type") != "Point":
return None
coords = geom.get("coordinates")
if not coords or len(coords) < 2:
return None
lon, lat = coords[0], coords[1]
if lon is None or lat is None:
return None
return "ST_SetSRID(ST_MakePoint(:lon, :lat), 4326)", {"lon": lon, "lat": lat}
def load_power_supply_centers(db: Session | None = None) -> dict[str, int]:
"""Тянет WFS + UPSERT-ит центры питания в power_supply_centers.
UPSERT по UNIQUE(source, external_id). Per-row SAVEPOINT — битая фича не валит
weekly-sync. ``db`` принимается для совместимости сигнатуры; если None — своя
SessionLocal (WFS-fetch + запись в БД sync, как utility_infrastructure_loader).
Returns: счётчики fetched/inserted/updated/skipped.
"""
features = fetch_power_supply_centers()
owns_session = db is None
if db is None:
db = SessionLocal()
inserted = 0
updated = 0
skipped = 0
fetched = len(features)
try:
for feature in features:
props: dict = feature.get("properties") or {}
sc_name = props.get("sc_name")
if not sc_name:
skipped += 1
continue
geom_pair = _point_geom_sql(feature)
geom_sql = geom_pair[0] if geom_pair else "NULL"
geom_params = geom_pair[1] if geom_pair else {}
params: dict = {
"source": "rosseti_wfs",
"external_id": _stable_external_id(feature, props),
"sc_name": sc_name,
"sc_name_norm": normalize_sc_name(sc_name),
"dzo_name": props.get("dzo_name"),
"org_name": props.get("org_name"),
"branch_name": props.get("br_name"),
"voltage_class": parse_voltage_class(sc_name),
"load_index": _map_load_index(props),
"municipality": props.get("municipality"),
"place_city": props.get("sc_placecity"),
"raw": json.dumps(feature, ensure_ascii=False),
**geom_params,
}
try:
with db.begin_nested(): # SAVEPOINT — откат только этой фичи
result = db.execute(
text(f"""
INSERT INTO power_supply_centers
(source, external_id, sc_name, sc_name_norm, dzo_name,
org_name, branch_name, voltage_class, load_index,
municipality, place_city, geom, raw, fetched_at)
VALUES (
:source, :external_id, :sc_name, :sc_name_norm, :dzo_name,
:org_name, :branch_name, :voltage_class, :load_index,
:municipality, :place_city, {geom_sql},
CAST(:raw AS jsonb), NOW()
)
ON CONFLICT (source, external_id) DO UPDATE
SET sc_name = EXCLUDED.sc_name,
sc_name_norm = EXCLUDED.sc_name_norm,
dzo_name = EXCLUDED.dzo_name,
org_name = EXCLUDED.org_name,
branch_name = EXCLUDED.branch_name,
voltage_class = EXCLUDED.voltage_class,
load_index = EXCLUDED.load_index,
municipality = EXCLUDED.municipality,
place_city = EXCLUDED.place_city,
geom = EXCLUDED.geom,
raw = EXCLUDED.raw,
fetched_at = NOW()
RETURNING (xmax = 0) AS is_insert
"""),
params,
).scalar()
if result:
inserted += 1
else:
updated += 1
except Exception as e:
logger.warning("rosseti_wfs upsert failed for %r: %s", sc_name, e)
skipped += 1
db.commit()
except Exception as e:
db.rollback()
logger.exception("rosseti_wfs: unexpected error, outer tx rolled back: %s", e)
raise
finally:
if owns_session:
db.close()
result_dict = {
"fetched": fetched,
"inserted": inserted,
"updated": updated,
"skipped": skipped,
}
logger.info("rosseti_wfs load done: %s", result_dict)
return result_dict