fix(site_finder): Россети — стабильный ключ вместо сессионного fid, дедуп power_supply_centers ×10 → ~481 #3329
3 changed files with 236 additions and 12 deletions
|
|
@ -19,6 +19,7 @@ per-row SAVEPOINT при UPSERT (битая фича не валит weekly-sync
|
|||
import hashlib
|
||||
import json
|
||||
import logging
|
||||
import math
|
||||
import re
|
||||
|
||||
import httpx
|
||||
|
|
@ -101,20 +102,47 @@ def parse_voltage_class(name: str | None) -> str | None:
|
|||
return m.group(1).replace(".", "/")
|
||||
|
||||
|
||||
def _stable_external_id(feature: dict, props: dict) -> str:
|
||||
"""Стабильный external_id фичи: feature['id'] или хэш ключевых полей.
|
||||
def _coord_e5(value: float | None) -> str:
|
||||
"""Координата → целое в единицах 1e-5 градуса (~1 м), полукруглением от нуля.
|
||||
|
||||
WFS обычно отдаёт стабильный ``feature['id']``; если его нет — детерминированный
|
||||
sha1 по (sc_name, координаты) чтобы UPSERT оставался идемпотентным.
|
||||
Целое, а не форматированный float: ключ обязан совпадать байт-в-байт с SQL-
|
||||
бэкфиллом (99c), а текстовое представление double в питоне и в PG разное.
|
||||
Округление ДВОИЧНОЕ (по значению double, не по десятичному представлению):
|
||||
64.423605*1e5 == 6442360.499999999 → 6442360, хотя «по десятичному» было бы
|
||||
6442361. В 99c та же семантика: floor/abs/sign над float8, БЕЗ каста в numeric
|
||||
(каст округляет по кратчайшему десятичному repr и расходится в 0.19% координат).
|
||||
"""
|
||||
fid = feature.get("id")
|
||||
if fid:
|
||||
return str(fid)
|
||||
geom = feature.get("geometry") or {}
|
||||
coords = geom.get("coordinates")
|
||||
seed = f"{props.get('sc_name', '')}|{coords}"
|
||||
# sha1 здесь — стабильный дедуп-id фичи, не криптография.
|
||||
return "h:" + hashlib.sha1(seed.encode("utf-8")).hexdigest()[:16]
|
||||
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:
|
||||
|
|
|
|||
|
|
@ -18,6 +18,7 @@ Pure / mock-based — без реальной сети и БД. Покрывае
|
|||
|
||||
from __future__ import annotations
|
||||
|
||||
import hashlib
|
||||
import io
|
||||
from contextlib import contextmanager
|
||||
from typing import Any
|
||||
|
|
@ -96,6 +97,64 @@ def test_map_load_index_unknown_and_missing() -> None:
|
|||
assert rw._map_load_index({"sc_indexload_id": "мусор"}) is None
|
||||
|
||||
|
||||
# ── _stable_external_id (#3322: сессионный fid раздувал таблицу ×10) ───────────
|
||||
|
||||
|
||||
def _wfs_feature(fid: str, name: str, lon: float, lat: float) -> dict[str, Any]:
|
||||
return {
|
||||
"id": fid,
|
||||
"geometry": {"type": "Point", "coordinates": [lon, lat]},
|
||||
"properties": {"sc_name": name},
|
||||
}
|
||||
|
||||
|
||||
def test_stable_external_id_ignores_session_fid() -> None:
|
||||
"""Разные сессионные fid + одинаковые атрибуты → ОДИН ключ (регрессия #3322)."""
|
||||
a = _wfs_feature("sc_points_fullview.fid--1a2b3c", "ПС 110/10 Уктус", 60.47123, 56.77456)
|
||||
b = _wfs_feature("sc_points_fullview.fid--9f8e7d", "ПС 110/10 Уктус", 60.47123, 56.77456)
|
||||
|
||||
key_a = rw._stable_external_id(a, a["properties"])
|
||||
assert key_a == rw._stable_external_id(b, b["properties"])
|
||||
# Проверка ПО ЗНАЧЕНИЮ: ключ = sha256 по «имя|напряжение|lon_e5|lat_e5»,
|
||||
# не fid. Тот же seed повторён в data/sql/99c_power_supply_centers_dedup.sql.
|
||||
expected = "h:" + hashlib.sha256("уктус|110/10|6047123|5677456".encode()).hexdigest()[:16]
|
||||
assert key_a == expected == "h:844e54f0152d2799"
|
||||
|
||||
|
||||
def test_stable_external_id_differs_on_attributes() -> None:
|
||||
"""Одинаковый fid, разные атрибуты (координата / имя) → РАЗНЫЕ ключи."""
|
||||
base = _wfs_feature("sc_points_fullview.fid--same", "ПС 110/10 Уктус", 60.47123, 56.77456)
|
||||
moved = _wfs_feature("sc_points_fullview.fid--same", "ПС 110/10 Уктус", 60.47124, 56.77456)
|
||||
renamed = _wfs_feature("sc_points_fullview.fid--same", "ПС 110/10 Северная", 60.47123, 56.77456)
|
||||
|
||||
keys = {rw._stable_external_id(f, f["properties"]) for f in (base, moved, renamed)}
|
||||
assert len(keys) == 3
|
||||
assert rw._stable_external_id(moved, moved["properties"]) == "h:8c5b5fa5f35c651b"
|
||||
assert rw._stable_external_id(renamed, renamed["properties"]) == "h:297e00f6e0d4f3b8"
|
||||
|
||||
|
||||
def test_stable_external_id_no_geometry() -> None:
|
||||
"""Фича без геометрии: координатные компоненты пустые, ключ всё равно стабилен."""
|
||||
f: dict[str, Any] = {"id": "fid--x", "properties": {"sc_name": "ПС 110/10 Уктус"}}
|
||||
assert rw._stable_external_id(f, f["properties"]) == "h:eb91917f35aff23f"
|
||||
|
||||
|
||||
def test_coord_e5_rounds_on_binary_double_not_decimal() -> None:
|
||||
"""Округление по ДВОИЧНОМУ double, не по десятичному представлению.
|
||||
|
||||
64.423605*1e5 == 6442360.499999999 → 6442360; «по десятичному» вышло бы 6442361
|
||||
(так считал бы round(ST_X(geom)::numeric*100000) — расхождение на 0.19% реальных
|
||||
координат). 99c обязана давать те же цифры, поэтому семантика закреплена тестом.
|
||||
"""
|
||||
assert rw._coord_e5(64.423605) == "6442360"
|
||||
assert rw._coord_e5(-64.423605) == "-6442360"
|
||||
# 60.123455*1e5 == ровно 6012345.5 → полукругление ОТ нуля, симметрично знаку.
|
||||
assert rw._coord_e5(60.123455) == "6012346"
|
||||
assert rw._coord_e5(-60.123455) == "-6012346"
|
||||
assert rw._coord_e5(60.6) == "6060000"
|
||||
assert rw._coord_e5(None) == ""
|
||||
|
||||
|
||||
# ── sanitize_tp_capacity_mva (кВА-санитайз) ───────────────────────────────────
|
||||
|
||||
|
||||
|
|
|
|||
137
data/sql/99c_power_supply_centers_dedup.sql
Normal file
137
data/sql/99c_power_supply_centers_dedup.sql
Normal file
|
|
@ -0,0 +1,137 @@
|
|||
-- 99c_power_supply_centers_dedup.sql
|
||||
-- Issue #3322 — power_supply_centers раздут ×10: 4880 строк на 481 уникальный ЦП.
|
||||
--
|
||||
-- Причина. rosseti_wfs_loader брал external_id из feature['id'] WFS-ответа, а
|
||||
-- GeoServer отдаёт СЕССИОННЫЙ fid (новый на каждый GetFeature) → ON CONFLICT
|
||||
-- (source, external_id) не срабатывал ни разу, каждый weekly-прогон добавлял
|
||||
-- полный набор ~488 фич заново. Починка разбора сама старые строки не убирает
|
||||
-- (ON CONFLICT ничего не перезапишет, ключи не совпадут) → нужен этот бэкфилл.
|
||||
--
|
||||
-- Что делает файл:
|
||||
-- (а) пересчитывает external_id по НОВОЙ формуле (см. ниже) для всех строк
|
||||
-- source='rosseti_wfs';
|
||||
-- (б) схлопывает копии: победитель группы — свежайший снапшот
|
||||
-- (fetched_at DESC NULLS LAST, id DESC — DESC в PG это NULLS FIRST,
|
||||
-- поэтому NULLS LAST задан ЯВНО);
|
||||
-- (в) печатает числа: строк до / после, удалено, переключено на новый ключ.
|
||||
-- Ожидание после прогона — ~481-488 строк (столько ЦП отдаёт источник).
|
||||
-- (г) идемпотентен: на повторном прогоне ключи уже совпадают → 0 удалений,
|
||||
-- 0 обновлений, «до» = «после».
|
||||
--
|
||||
-- ФОРМУЛА КЛЮЧА (дублирует rosseti_wfs_loader._stable_external_id — менять только
|
||||
-- синхронно, иначе следующий weekly-прогон вставит второй комплект строк):
|
||||
-- seed = sc_name_norm || '|' || voltage_class || '|' || lon_e5 || '|' || lat_e5
|
||||
-- external_id = 'h:' || left(hex(sha256(utf8(seed))), 16)
|
||||
-- где lon_e5/lat_e5 — координата в единицах 1e-5 градуса (~1 м), округление
|
||||
-- floor(|v|*1e5 + 0.5) со знаком — ДВОИЧНОЕ, ровно как в питоне; пустая строка,
|
||||
-- если geom отсутствует. Целые, а не форматированный float — текстовое
|
||||
-- представление double в питоне и в PG различается.
|
||||
--
|
||||
-- ПОЧЕМУ НЕ round(...::numeric): каст float8→numeric берёт кратчайшее десятичное
|
||||
-- представление, и округление идёт по нему, а не по двоичному double. На реальных
|
||||
-- координатах расходится в 0.19% случаев (замер: 761 из 400000), напр. 64.423605
|
||||
-- → питон 6442360 (двоичное 6442360.499999999), numeric-путь 6442361. Каждое
|
||||
-- расхождение = вечный дубль ЦП, который сам не зарастёт: миграция применяется
|
||||
-- один раз (_schema_migrations). Поэтому в SQL считаем ТЕМ ЖЕ double: floor/abs/
|
||||
-- sign над float8 — это IEEE754, бит в бит как math.floor в питоне.
|
||||
-- sha256, а не sha1: sha256 встроен в PG16, sha1 потребовал бы pgcrypto.
|
||||
--
|
||||
-- Байт-в-байт совпадение с питоном держится на том, что SQL НИЧЕГО не нормализует
|
||||
-- сам: sc_name_norm и voltage_class — уже готовые колонки, их записал тот же
|
||||
-- normalize_sc_name / parse_voltage_class. Если normalize_sc_name когда-нибудь
|
||||
-- изменится, старые sc_name_norm разъедутся с новыми ключами — тогда нужен
|
||||
-- повторный прогон логики этого файла (он идемпотентен, ре-apply безопасен).
|
||||
--
|
||||
-- Порядок: миграция ПЕРЕД деплоем кода (schema-first) — новый код после неё
|
||||
-- попадает ON CONFLICT-ом в уже схлопнутые строки.
|
||||
--
|
||||
-- Naming: deploy.yml применяет файлы по `ls -1 data/sql/*.sql | sort`;
|
||||
-- '99c_' идёт после '99b_grant_quarter_price_index_fdw.sql' ('b' < 'c').
|
||||
|
||||
BEGIN;
|
||||
|
||||
SET LOCAL lock_timeout = '5s';
|
||||
|
||||
DO $$
|
||||
DECLARE
|
||||
rows_before bigint;
|
||||
names_before bigint;
|
||||
rows_after bigint;
|
||||
names_after bigint;
|
||||
deleted bigint;
|
||||
rekeyed bigint;
|
||||
BEGIN
|
||||
SELECT count(*), count(DISTINCT sc_name_norm)
|
||||
INTO rows_before, names_before
|
||||
FROM power_supply_centers
|
||||
WHERE source = 'rosseti_wfs';
|
||||
|
||||
CREATE TEMP TABLE psc_new_key ON COMMIT DROP AS
|
||||
SELECT
|
||||
id,
|
||||
fetched_at,
|
||||
'h:' || substring(
|
||||
encode(
|
||||
sha256(convert_to(
|
||||
sc_name_norm
|
||||
|| '|' || coalesce(voltage_class, '')
|
||||
|| '|' || CASE WHEN geom IS NULL THEN ''
|
||||
ELSE (sign(ST_X(geom))
|
||||
* floor(abs(ST_X(geom)) * 100000 + 0.5))::bigint::text END
|
||||
|| '|' || CASE WHEN geom IS NULL THEN ''
|
||||
ELSE (sign(ST_Y(geom))
|
||||
* floor(abs(ST_Y(geom)) * 100000 + 0.5))::bigint::text END,
|
||||
'UTF8'
|
||||
)),
|
||||
'hex'
|
||||
) FROM 1 FOR 16
|
||||
) AS new_key
|
||||
FROM power_supply_centers
|
||||
WHERE source = 'rosseti_wfs';
|
||||
|
||||
-- (б) схлопывание: оставляем свежайший снапшот каждой группы.
|
||||
-- Резервы (reserve_mva и пр.) не теряются: rosseti/eesk-лоадеры пишут их
|
||||
-- UPDATE-ом по sc_name_norm, т.е. во ВСЕ копии сразу, победитель их несёт.
|
||||
WITH ranked AS (
|
||||
SELECT
|
||||
id,
|
||||
row_number() OVER (
|
||||
PARTITION BY new_key
|
||||
ORDER BY fetched_at DESC NULLS LAST, id DESC
|
||||
) AS rn
|
||||
FROM psc_new_key
|
||||
)
|
||||
DELETE FROM power_supply_centers p
|
||||
USING ranked r
|
||||
WHERE p.id = r.id
|
||||
AND r.rn > 1;
|
||||
GET DIAGNOSTICS deleted = ROW_COUNT;
|
||||
|
||||
-- (а) пересчёт ключа у выживших. После DELETE каждый new_key принадлежит
|
||||
-- ровно одной строке → UNIQUE (source, external_id) не нарушается.
|
||||
-- IS DISTINCT FROM даёт идемпотентность: второй прогон обновит 0 строк.
|
||||
UPDATE power_supply_centers p
|
||||
SET external_id = k.new_key
|
||||
FROM psc_new_key k
|
||||
WHERE p.id = k.id
|
||||
AND p.external_id IS DISTINCT FROM k.new_key;
|
||||
GET DIAGNOSTICS rekeyed = ROW_COUNT;
|
||||
|
||||
SELECT count(*), count(DISTINCT sc_name_norm)
|
||||
INTO rows_after, names_after
|
||||
FROM power_supply_centers
|
||||
WHERE source = 'rosseti_wfs';
|
||||
|
||||
RAISE NOTICE '#3322 power_supply_centers: было % строк / % имён -> стало % строк / % имён (удалено %, переключено на стабильный ключ %)',
|
||||
rows_before, names_before, rows_after, names_after, deleted, rekeyed;
|
||||
|
||||
-- 700 — потолок здравого смысла: источник отдаёт ~488 ЦП по области.
|
||||
-- Превышение = формула ключа не схлопнула дубли. EXCEPTION, а не WARNING:
|
||||
-- иначе файл пометится applied навсегда, а дубли останутся. Откат всей
|
||||
-- транзакции ничего не теряет и оставляет миграцию непринятой до разбора.
|
||||
IF rows_after > 700 THEN
|
||||
RAISE EXCEPTION '#3322: после дедупа осталось % строк (ожидалось ~481-488) — формула ключа не схлопнула дубли, транзакция откачена', rows_after;
|
||||
END IF;
|
||||
END $$;
|
||||
|
||||
COMMIT;
|
||||
Loading…
Add table
Reference in a new issue