gendesign/backend/app/services/site_finder/rosseti_wfs_loader.py
bot-backend 6a2d873c70
All checks were successful
CI Trade-In / changes (pull_request) Successful in 12s
CI / changes (pull_request) Successful in 14s
CI Trade-In / backend-tests (pull_request) Has been skipped
CI Trade-In / browser-tests (pull_request) Has been skipped
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 2m24s
CI / backend-tests (pull_request) Successful in 17m41s
fix(rosseti): считать координату ключа в SQL тем же double, что в питоне
Ревью PR #3329. round(ST_X(geom)::numeric * 100000) округляет по кратчайшему
десятичному представлению float8, а питон — по двоичному double: расхождение на
0.19% реальных координат (761 из 400000), напр. 64.423605 → питон 6442360
(6442360.499999999), numeric-путь 6442361. Каждое расхождение = вечный дубль ЦП,
который сам не зарастёт — миграция применяется один раз (_schema_migrations).
Теперь в SQL sign/floor/abs над float8 без каста в numeric: IEEE754 бит в бит
как math.floor в питоне.

test_coord_e5 брал 60.123455, где двоичное и десятичное округление совпадают —
защита, которая не защищает. Добавлено расходящееся значение 64.423605.

RAISE WARNING при rows_after > 700 заменён на RAISE EXCEPTION: warning не
останавливает прогон, файл помечался бы applied навсегда вместе с дублями.

Refs #3322
2026-09-02 14:51:47 +05:00

297 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 разное.
Округление ДВОИЧНОЕ (по значению double, не по десятичному представлению):
64.423605*1e5 == 6442360.499999999 → 6442360, хотя «по десятичному» было бы
6442361. В 99c та же семантика: floor/abs/sign над float8, БЕЗ каста в numeric
(каст округляет по кратчайшему десятичному repr и расходится в 0.19% координат).
"""
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