`UNIQUE (doc_group, doc_num)` вводился, чтобы склеивать ОДИН документ,
пришедший из двух схем портала. Замер 20.08.2026 показал, что задача,
ради которой ключ введён, почти отсутствует, а побочный эффект огромен:
docNum у ГИСОГД НЕ уникален — разрешение и изменения к нему носят один
номер.
группа документов различных key различных docNum схлопывается
DocRS 6098 6096 4305 1793
DocRV 5419 5415 4969 450
DocIZ 548 547 393 155
общих docNum между схемами (DocRS): 2 ← ради этого ключ и вводился
общих key между схемами (DocRS): 2 ← те же два
На проде 9182 строки против 12 065 документов на портале — нет 23.9 %
реестра. Пример 66-06-06-2026: портал отдаёт два документа (key …719586 —
само разрешение, key …752293 — изменения к нему), а UPSERT с
предпочтением позднего date_reg оставлял только изменение. Так вытеснено
598 из 4320 строк РНС (13.8 %) — в §6 на месте разрешения показывается
изменение к нему, без признака подмены.
Ключ стал `UNIQUE (source_key)`: разделяет разрешение и изменения (разные
key) и по-прежнему склеивает настоящие межсхемные дубли (у них key
ОБЩИЙ — ровно 7 записей по всем группам). Дедуп перед сменой не нужен:
source_key на проде уже уникален (9182 из 9182, NOT NULL).
Заодно группа DocIZ добавлена в GROUP_CODE — её не было вовсе, 548
документов не грузились. CHECK расширен значением 'IZ'.
§6 сужена до РНС/РВЭ ЯВНО: агрегат обещает total_count = rs_count +
rv_count, а строки 'IZ' попадали бы в total и ни в один счётчик.
Показывать ли изменения отдельной строкой — вопрос продуктовый (#2986);
до его решения сужение стоит в запросе, а не держится на том, что таких
строк «пока нет».
Проверки:
- два гейта на лоадер (GROUP_CODE и цель ON CONFLICT) — БЕЗ базы,
двусторонние: на origin/main дают конкретные неверные значения
({'DocRS','DocRV'} и старый ON CONFLICT в тексте запроса);
- гейт на §6 и контроль инварианта total = rs + rv на данных — красные
на origin/main;
- герметичная репетиция миграции на временной копии: со старым ключом
разрешение и изменение схлопываются в одну строку (и остаётся именно
изменение — как на проде), после миграции живут раздельно; межсхемный
дубль по-прежнему склеивается; CHECK принимает 'IZ' и отвергает мусор.
Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
512 lines
26 KiB
Python
512 lines
26 KiB
Python
"""Загрузчик реестра РНС/РВЭ Свердловской области из открытого API ГИСОГД-СО (#2367).
|
||
|
||
Источник: ``gisogd66.midural.ru`` — региональная ГИС органов гос. власти Свердловской
|
||
области. Открытый JSON-API раздела 13 («Сведения о документах в деле о земельном
|
||
участке»), БЕЗ авторизации (bare curl 200, разведка 2026-07-04).
|
||
|
||
Три эндпоинта (все GET JSON, base ``https://gisogd66.midural.ru/api/v1/{schema}``):
|
||
* ``/gisogddocgroups/razdel13`` — группы документов раздела (DocRS/DocRV/…);
|
||
* ``/gisogddocgroupdocs/razdel13/{group}`` — весь список группы одним ответом (без
|
||
пагинации): key, name, docNum, dateDoc, numReg, dateReg;
|
||
* ``/gisogddocs/razdel13/{key}`` — карточка документа: spatialUnit (кад.№,
|
||
может быть несколько через запятую / отсутствовать), infosetGeoJson (GeoJSON Polygon
|
||
WGS84 lon/lat, может отсутствовать — hasGeometry=false), approvedOrganization, docNum,
|
||
dateDoc, dateReg.
|
||
|
||
Две схемы (тянем обе, sverdregion первым): ``agate_sverdregion`` (регион-агрегат,
|
||
стабильные карточки, покрывает и ЕКБ) и ``agate_ekbgo`` (ЕКБ, свежее на ~неделю, но
|
||
часть старых ключей битая на портале). Один документ может встретиться в обеих →
|
||
бизнес-ключ ``(doc_group, doc_num)``; при конфликте предпочитаем запись с более
|
||
поздним ``date_reg``.
|
||
|
||
Инкрементальность: список группы сверяется с БД по (source_schema, source_key, date_reg);
|
||
карточки (второй запрос — дорогой) тянутся ТОЛЬКО для новых / изменившихся документов.
|
||
|
||
Politeness: единый httpx-клиент с limits(max_connections=5), timeout=20с, 2 ретрая на
|
||
сетевых ошибках. Карточки (5xx-битые записи портала) НЕ ретраятся — deterministically
|
||
broken, не транзиент. Circuit breaker: 50 карточек подряд с ошибкой → прерываем группу
|
||
(битый сегмент). UPSERT — per-row SAVEPOINT (битая строка не валит батч). Периодический
|
||
db.commit() (каждые 200 записей + после каждой schema/group), чтобы партиал пережил
|
||
soft_time_limit-kill.
|
||
"""
|
||
|
||
from __future__ import annotations
|
||
|
||
import json
|
||
import logging
|
||
import re
|
||
import time
|
||
from datetime import date, datetime
|
||
from typing import Any
|
||
|
||
import httpx
|
||
from sqlalchemy import text
|
||
from sqlalchemy.orm import Session
|
||
|
||
from app.core.db import SessionLocal
|
||
|
||
logger = logging.getLogger(__name__)
|
||
|
||
# ── Константы источника ───────────────────────────────────────────────────────
|
||
|
||
GISOGD_BASE = "https://gisogd66.midural.ru/api/v1"
|
||
_SECTION = "razdel13"
|
||
|
||
# Схемы источника — обе тянем. Порядок ВАЖЕН: sverdregion первым.
|
||
# sverdregion (регион-агрегат) отдаёт карточки стабильно и покрывает в т.ч. ЕКБ-доки
|
||
# через агрегацию, поэтому он — базовый источник. ekbgo вторым как freshness-энхансер
|
||
# (свежее на ~неделю), НО у него большой сегмент старых ключей (префикс ~100001007xxx)
|
||
# детерминированно отдаёт 500 на карточках (сломано на стороне портала, 2026-07-04) —
|
||
# гоняем его после того, как регион уже покрыл основную массу (см. circuit breaker ниже).
|
||
SCHEMAS: tuple[str, ...] = ("agate_sverdregion", "agate_ekbgo")
|
||
|
||
# Группа источника (doc group key) → наш doc_group-код в БД (CHECK IN ('RS','RV','IZ')).
|
||
#
|
||
# DocIZ («Изменение в Разрешение на строительство») есть на портале с самого начала,
|
||
# но в этом словаре его не было — 548 документов не грузились вовсе (#2986).
|
||
GROUP_CODE: dict[str, str] = {"DocRS": "RS", "DocRV": "RV", "DocIZ": "IZ"}
|
||
|
||
_HTTP_TIMEOUT = 20.0
|
||
_MAX_CONNECTIONS = 5
|
||
_MAX_RETRIES = 2 # ретраи (helper делает суммарно до 1 + _MAX_RETRIES попыток)
|
||
_RETRY_BACKOFF_S = 1.5
|
||
_CARD_SLEEP_S = 0.05 # вежливая пауза между карточками (второй, «дорогой» запрос)
|
||
_COMMIT_EVERY = 200 # периодический commit внутри группы (партиал переживает time-limit)
|
||
# Circuit breaker: столько карточек ПОДРЯД с ошибкой fetch → прерываем группу (портал
|
||
# отдаёт битый сегмент старых ключей детерминированным 5xx; остаток догоним будущими ранами).
|
||
_CARD_FAIL_BREAKER = 50
|
||
|
||
USER_AGENT = "GenDesign/1.0 (+https://gendsgn.ru) Site Finder gisogd permits loader"
|
||
|
||
# Валидация кадастрового номера: NN:NN:NNNNNN(N):N (6-7 знаков в третьем блоке).
|
||
_CAD_RE = re.compile(r"^\d{2}:\d{2}:\d{6,7}:\d+$")
|
||
|
||
|
||
# ── HTTP helpers (module-level — мокаются в тестах) ───────────────────────────
|
||
|
||
|
||
def _make_client() -> httpx.Client:
|
||
"""httpx.Client с пулом (max_connections=5), таймаутом и UA. Закрывать вызывающим."""
|
||
return httpx.Client(
|
||
base_url=GISOGD_BASE,
|
||
timeout=_HTTP_TIMEOUT,
|
||
follow_redirects=True,
|
||
headers={"User-Agent": USER_AGENT},
|
||
limits=httpx.Limits(max_connections=_MAX_CONNECTIONS),
|
||
)
|
||
|
||
|
||
def _get_json(client: httpx.Client, path: str, *, retries_5xx: int = _MAX_RETRIES) -> Any:
|
||
"""GET path → JSON с ретраями на 5xx / сетевых ошибках.
|
||
|
||
Сетевые Timeout/TransportError ретраятся всегда (до 1 + _MAX_RETRIES попыток) —
|
||
это транзиент. 5xx ретраятся ``retries_5xx`` раз: списки группы — с дефолтом
|
||
(транзиентная нагрузка портала), карточки — с ``retries_5xx=0`` (5xx карточки
|
||
детерминированно битые записи сегмента портала, не транзиент). 4xx → сразу raise.
|
||
"""
|
||
last_exc: Exception | None = None
|
||
for attempt in range(_MAX_RETRIES + 1):
|
||
try:
|
||
resp = client.get(path)
|
||
if resp.status_code >= 500:
|
||
last_exc = httpx.HTTPStatusError(
|
||
f"{resp.status_code} on {path}", request=resp.request, response=resp
|
||
)
|
||
# 5xx: ретраим только пока attempt в пределах retries_5xx (0 → сразу raise).
|
||
if attempt >= retries_5xx:
|
||
raise last_exc
|
||
logger.warning(
|
||
"gisogd66: %s (попытка %d) на %s", resp.status_code, attempt + 1, path
|
||
)
|
||
time.sleep(_RETRY_BACKOFF_S * (attempt + 1))
|
||
continue
|
||
resp.raise_for_status()
|
||
return resp.json()
|
||
except (httpx.TimeoutException, httpx.TransportError) as exc:
|
||
last_exc = exc
|
||
logger.warning(
|
||
"gisogd66: сетевая ошибка (попытка %d) на %s: %s", attempt + 1, path, exc
|
||
)
|
||
time.sleep(_RETRY_BACKOFF_S * (attempt + 1))
|
||
assert last_exc is not None
|
||
raise last_exc
|
||
|
||
|
||
def fetch_group_list(client: httpx.Client, schema: str, group: str) -> list[dict[str, Any]]:
|
||
"""Весь список документов группы (без пагинации) → list of dict.
|
||
|
||
Поля элемента: key, name, docNum, dateDoc, numReg, dateReg. Пустой список при
|
||
неожиданном (не-list) ответе.
|
||
"""
|
||
data = _get_json(client, f"/{schema}/gisogddocgroupdocs/{_SECTION}/{group}")
|
||
if not isinstance(data, list):
|
||
logger.warning("gisogd66: %s/%s список не-list (%s)", schema, group, type(data).__name__)
|
||
return []
|
||
logger.info("gisogd66: %s/%s список документов=%d", schema, group, len(data))
|
||
return data
|
||
|
||
|
||
def fetch_doc_card(client: httpx.Client, schema: str, key: str | int) -> dict[str, Any] | None:
|
||
"""Карточка документа → dict (или None при не-dict / ошибке карточки).
|
||
|
||
Поля: spatialUnit, infosetGeoJson, hasGeometry, approvedOrganization, docNum,
|
||
dateDoc, dateReg, name.
|
||
|
||
5xx карточки НЕ ретраятся (retries_5xx=0): битый сегмент старых ключей портала
|
||
отдаёт детерминированный 500 — 3 попытки с backoff = ~12с потерянного времени/карточку.
|
||
"""
|
||
data = _get_json(client, f"/{schema}/gisogddocs/{_SECTION}/{key}", retries_5xx=0)
|
||
if not isinstance(data, dict):
|
||
logger.warning("gisogd66: карточка %s/%s не-dict (%s)", schema, key, type(data).__name__)
|
||
return None
|
||
return data
|
||
|
||
|
||
# ── Парсинг полей ─────────────────────────────────────────────────────────────
|
||
|
||
|
||
def _parse_date(value: Any) -> date | None:
|
||
"""ISO-datetime / date-строка источника → date. None если пусто / не парсится.
|
||
|
||
Источник отдаёт «2026-06-02T00:00:00» (dateDoc) и «2026-06-19T15:55:11» (dateReg).
|
||
"""
|
||
if not value:
|
||
return None
|
||
if isinstance(value, datetime):
|
||
return value.date()
|
||
if isinstance(value, date):
|
||
return value
|
||
s = str(value).strip()
|
||
if not s:
|
||
return None
|
||
# Обрезаем таймзону-суффикс 'Z' если вдруг придёт, парсим ISO.
|
||
try:
|
||
return datetime.fromisoformat(s.replace("Z", "+00:00")).date()
|
||
except ValueError:
|
||
pass
|
||
for fmt in ("%Y-%m-%d", "%d.%m.%Y"):
|
||
try:
|
||
return datetime.strptime(s[:10], fmt).date()
|
||
except ValueError:
|
||
continue
|
||
logger.debug("gisogd66: не распарсил дату %r", value)
|
||
return None
|
||
|
||
|
||
def parse_cad_nums(spatial_unit: Any) -> list[str]:
|
||
"""spatialUnit → список валидных кадастровых номеров (может быть несколько через запятую).
|
||
|
||
Split по запятой, strip, валидация по _CAD_RE. Невалидные (мусор) отбрасываются с
|
||
warning-логом. Отсутствие / пусто → []. Дубликаты схлопываются (сохраняя порядок).
|
||
"""
|
||
if not spatial_unit:
|
||
return []
|
||
out: list[str] = []
|
||
seen: set[str] = set()
|
||
for part in str(spatial_unit).split(","):
|
||
cad = part.strip()
|
||
if not cad:
|
||
continue
|
||
if not _CAD_RE.match(cad):
|
||
logger.warning(
|
||
"gisogd66: невалидный кадастровый номер %r в spatialUnit — отброшен", cad
|
||
)
|
||
continue
|
||
if cad not in seen:
|
||
seen.add(cad)
|
||
out.append(cad)
|
||
return out
|
||
|
||
|
||
def extract_geojson(card: dict[str, Any]) -> str | None:
|
||
"""infosetGeoJson карточки → GeoJSON-строка (для ST_GeomFromGeoJSON) или None.
|
||
|
||
Источник кладёт GeoJSON КАК СТРОКУ в поле infosetGeoJson. Пусто / hasGeometry=false
|
||
/ невалидный JSON → None (лоадер запишет geom=NULL). Приведение Polygon→MultiPolygon
|
||
делается в SQL через ST_Multi — здесь возвращаем сырой GeoJSON как есть.
|
||
"""
|
||
raw = card.get("infosetGeoJson")
|
||
if not raw:
|
||
return None
|
||
if isinstance(raw, dict):
|
||
# На случай если источник когда-нибудь отдаст объект, а не строку.
|
||
return json.dumps(raw)
|
||
s = str(raw).strip()
|
||
if not s:
|
||
return None
|
||
try:
|
||
json.loads(s) # sanity-парс: битую строку не отправляем в PostGIS
|
||
except (ValueError, TypeError) as exc:
|
||
logger.warning("gisogd66: битый infosetGeoJson (%s) — geom=NULL", exc)
|
||
return None
|
||
return s
|
||
|
||
|
||
# ── Инкрементальность ─────────────────────────────────────────────────────────
|
||
|
||
|
||
def _existing_reg_dates(db: Session, schema: str) -> dict[str, date | None]:
|
||
"""{source_key: date_reg} для уже загруженных из этой схемы записей.
|
||
|
||
Используется для инкрементальной сверки: карточку тянем только если key новый ИЛИ
|
||
date_reg в списке отличается от сохранённого.
|
||
"""
|
||
rows = db.execute(
|
||
text("SELECT source_key, date_reg FROM gisogd_permits WHERE source_schema = :schema"),
|
||
{"schema": schema},
|
||
).all()
|
||
return {str(r[0]): r[1] for r in rows}
|
||
|
||
|
||
# ── UPSERT ────────────────────────────────────────────────────────────────────
|
||
|
||
|
||
def _upsert_permit(db: Session, rec: dict[str, Any]) -> str:
|
||
"""UPSERT одной записи в gisogd_permits по source_key. Per-row SAVEPOINT.
|
||
|
||
Ключ — идентификатор документа на портале (#2986). Раньше ключом был
|
||
(doc_group, doc_num), но docNum у ГИСОГД НЕ уникален: разрешение и изменения к
|
||
нему носят один номер, и UPSERT оставлял только одно из них. Замер 20.08.2026:
|
||
так схлопывалось 2243 документа, а межсхемных дублей — ради которых ключ и
|
||
вводился — всего 7, и у них key ОБЩИЙ, то есть новый ключ их тоже склеивает.
|
||
|
||
Конфликт-резолв: при коллизии обновляем ТОЛЬКО если у новой записи date_reg НЕ
|
||
старше сохранённой (EXCLUDED.date_reg >= existing) — предпочитаем более позднюю
|
||
регистрацию (или запись без даты не затирает датированную). Осмыслен ровно для
|
||
тех 7 межсхемных совпадений.
|
||
|
||
Returns: 'inserted' | 'updated' | 'skipped_unchanged'.
|
||
"""
|
||
result = db.execute(
|
||
text("""
|
||
INSERT INTO gisogd_permits
|
||
(doc_group, doc_num, doc_name, date_doc, date_reg,
|
||
approved_organization, source_schema, source_key,
|
||
cad_nums, geom, fetched_at, updated_at)
|
||
VALUES (
|
||
:doc_group, :doc_num, :doc_name, :date_doc, :date_reg,
|
||
:approved_org, :source_schema, :source_key,
|
||
CAST(:cad_nums AS text[]),
|
||
-- CAST обязателен: без него ни одно вхождение :geojson не даёт Postgres
|
||
-- тип параметра (ST_GeomFromGeoJSON перегружена text/json/jsonb) →
|
||
-- psycopg.errors.AmbiguousParameter на КАЖДОЙ строке (прод-прогон #2367).
|
||
CASE WHEN CAST(:geojson AS text) IS NULL THEN NULL
|
||
ELSE ST_Multi(ST_SetSRID(
|
||
ST_GeomFromGeoJSON(CAST(:geojson AS text)), 4326)) END,
|
||
NOW(), NOW()
|
||
)
|
||
ON CONFLICT (source_key) DO UPDATE
|
||
SET doc_group = EXCLUDED.doc_group,
|
||
doc_num = EXCLUDED.doc_num,
|
||
doc_name = EXCLUDED.doc_name,
|
||
date_doc = EXCLUDED.date_doc,
|
||
date_reg = EXCLUDED.date_reg,
|
||
approved_organization = EXCLUDED.approved_organization,
|
||
source_schema = EXCLUDED.source_schema,
|
||
cad_nums = EXCLUDED.cad_nums,
|
||
geom = EXCLUDED.geom,
|
||
updated_at = NOW()
|
||
WHERE EXCLUDED.date_reg IS NOT DISTINCT FROM gisogd_permits.date_reg
|
||
OR EXCLUDED.date_reg > gisogd_permits.date_reg
|
||
OR gisogd_permits.date_reg IS NULL
|
||
RETURNING (xmax = 0) AS is_insert
|
||
"""),
|
||
{
|
||
"doc_group": rec["doc_group"],
|
||
"doc_num": rec["doc_num"],
|
||
"doc_name": rec["doc_name"],
|
||
"date_doc": rec["date_doc"],
|
||
"date_reg": rec["date_reg"],
|
||
"approved_org": rec["approved_organization"],
|
||
"source_schema": rec["source_schema"],
|
||
"source_key": rec["source_key"],
|
||
"cad_nums": rec["cad_nums"],
|
||
"geojson": rec["geojson"],
|
||
},
|
||
).scalar()
|
||
if result is None:
|
||
# WHERE-гейт не пропустил апдейт (у нас более старый date_reg) → без изменений.
|
||
return "skipped_unchanged"
|
||
return "inserted" if result else "updated"
|
||
|
||
|
||
# ── Оркестрация ───────────────────────────────────────────────────────────────
|
||
|
||
|
||
def load_gisogd_permits(db: Session | None = None) -> dict[str, int]:
|
||
"""Инкрементально загрузить РНС/РВЭ ГИСОГД-СО (обе схемы, обе группы) в gisogd_permits.
|
||
|
||
Для каждой (schema, group): список → сверка с БД по (source_key, date_reg) →
|
||
карточки только для новых/изменённых → UPSERT по (doc_group, doc_num) с
|
||
предпочтением позднего date_reg. ``db`` для совместимости сигнатуры; None → своя
|
||
SessionLocal.
|
||
|
||
Returns: счётчики {listed, fetched_details, inserted, updated, skipped_unchanged,
|
||
no_cad, no_geom, card_failed, group_aborted}.
|
||
"""
|
||
owns_session = db is None
|
||
if db is None:
|
||
db = SessionLocal()
|
||
|
||
counts = {
|
||
"listed": 0,
|
||
"fetched_details": 0,
|
||
"inserted": 0,
|
||
"updated": 0,
|
||
"skipped_unchanged": 0,
|
||
"no_cad": 0,
|
||
"no_geom": 0,
|
||
"card_failed": 0, # каждая упавшая карточка (5xx / сеть / не-dict)
|
||
"group_aborted": 0, # групп прервано circuit breaker'ом (битый сегмент портала)
|
||
}
|
||
client = _make_client()
|
||
try:
|
||
for schema in SCHEMAS:
|
||
existing = _existing_reg_dates(db, schema)
|
||
for group_key, doc_group in GROUP_CODE.items():
|
||
try:
|
||
items = fetch_group_list(client, schema, group_key)
|
||
except Exception as exc:
|
||
logger.exception("gisogd66: список %s/%s упал: %s", schema, group_key, exc)
|
||
continue
|
||
counts["listed"] += len(items)
|
||
_sync_group(db, client, schema, doc_group, items, existing, counts)
|
||
# Commit после КАЖДОЙ (schema, group), а не один раз в конце: партиал
|
||
# должен пережить soft_time_limit-kill (иначе теряем весь прогон).
|
||
db.commit()
|
||
except Exception as exc:
|
||
db.rollback()
|
||
logger.exception("gisogd66: outer tx rolled back: %s", exc)
|
||
raise
|
||
finally:
|
||
client.close()
|
||
if owns_session:
|
||
db.close()
|
||
|
||
logger.info("load_gisogd_permits done: %s", counts)
|
||
return counts
|
||
|
||
|
||
def _sync_group(
|
||
db: Session,
|
||
client: httpx.Client,
|
||
schema: str,
|
||
doc_group: str,
|
||
items: list[dict[str, Any]],
|
||
existing: dict[str, date | None],
|
||
counts: dict[str, int],
|
||
) -> None:
|
||
"""Обработать одну (schema, group): инкрементальная сверка + карточки + UPSERT.
|
||
|
||
Периодический commit каждые ``_COMMIT_EVERY`` обработанных записей (партиал переживёт
|
||
time-limit kill). Circuit breaker: ``_CARD_FAIL_BREAKER`` карточек ПОДРЯД с ошибкой
|
||
fetch → прерываем группу (битый сегмент портала — остаток догоним будущими ранами).
|
||
Любая успешная карточка сбрасывает счётчик подряд-фейлов.
|
||
"""
|
||
processed = 0 # успешно обработанных (upsert/skip) с последнего commit
|
||
consecutive_fail = 0 # карточек подряд с ошибкой fetch (для circuit breaker)
|
||
for item in items:
|
||
key = item.get("key")
|
||
if key is None:
|
||
continue
|
||
source_key = str(key)
|
||
list_reg = _parse_date(item.get("dateReg"))
|
||
|
||
# Инкрементальность: пропускаем карточку, если key уже есть И date_reg совпал.
|
||
if source_key in existing and existing[source_key] == list_reg:
|
||
counts["skipped_unchanged"] += 1
|
||
consecutive_fail = 0 # skip — не фейл карточки, сбрасываем breaker
|
||
processed += 1
|
||
if processed % _COMMIT_EVERY == 0:
|
||
db.commit()
|
||
continue
|
||
|
||
card: dict[str, Any] | None
|
||
try:
|
||
card = fetch_doc_card(client, schema, source_key)
|
||
except Exception as exc:
|
||
logger.warning("gisogd66: карточка %s/%s упала: %s", schema, source_key, exc)
|
||
card = None
|
||
finally:
|
||
time.sleep(_CARD_SLEEP_S)
|
||
if card is None:
|
||
# Неуспех fetch (5xx / сеть / не-dict): считаем и двигаем circuit breaker.
|
||
counts["card_failed"] += 1
|
||
consecutive_fail += 1
|
||
if consecutive_fail >= _CARD_FAIL_BREAKER:
|
||
logger.error(
|
||
"gisogd66: %s/%s: %d карточек подряд 5xx — прерываю группу "
|
||
"(битый сегмент портала)",
|
||
schema,
|
||
doc_group,
|
||
_CARD_FAIL_BREAKER,
|
||
)
|
||
counts["group_aborted"] += 1
|
||
return
|
||
continue
|
||
consecutive_fail = 0 # успешная карточка сбрасывает breaker
|
||
counts["fetched_details"] += 1
|
||
|
||
rec = _build_record(item, card, schema, doc_group, source_key, list_reg)
|
||
if rec is None:
|
||
continue
|
||
if not rec["cad_nums"]:
|
||
counts["no_cad"] += 1
|
||
if rec["geojson"] is None:
|
||
counts["no_geom"] += 1
|
||
|
||
try:
|
||
with db.begin_nested(): # SAVEPOINT — откат только этой строки
|
||
outcome = _upsert_permit(db, rec)
|
||
except Exception as exc:
|
||
logger.warning("gisogd66: upsert %s %r упал: %s", doc_group, rec["doc_num"], exc)
|
||
continue
|
||
counts[outcome] += 1
|
||
processed += 1
|
||
if processed % _COMMIT_EVERY == 0:
|
||
db.commit()
|
||
|
||
|
||
def _build_record(
|
||
item: dict[str, Any],
|
||
card: dict[str, Any],
|
||
schema: str,
|
||
doc_group: str,
|
||
source_key: str,
|
||
list_reg: date | None,
|
||
) -> dict[str, Any] | None:
|
||
"""Список-элемент + карточка → нормализованная запись для UPSERT. None если нет doc_num.
|
||
|
||
doc_num берём из карточки, fallback на список; strip обязателен (в списке ~25%
|
||
номеров с ведущими/хвостовыми пробелами). Даты предпочитаем из карточки, fallback
|
||
на список. cad_nums из spatialUnit карточки, geojson из infosetGeoJson.
|
||
"""
|
||
raw_num = card.get("docNum") or item.get("docNum")
|
||
doc_num = str(raw_num).strip() if raw_num else ""
|
||
if not doc_num:
|
||
logger.warning("gisogd66: карточка %s/%s без docNum — пропуск", schema, source_key)
|
||
return None
|
||
|
||
return {
|
||
"doc_group": doc_group,
|
||
"doc_num": doc_num,
|
||
"doc_name": card.get("name") or item.get("name"),
|
||
"date_doc": _parse_date(card.get("dateDoc")) or _parse_date(item.get("dateDoc")),
|
||
"date_reg": _parse_date(card.get("dateReg")) or list_reg,
|
||
"approved_organization": card.get("approvedOrganization"),
|
||
"source_schema": schema,
|
||
"source_key": source_key,
|
||
"cad_nums": parse_cad_nums(card.get("spatialUnit")),
|
||
"geojson": extract_geojson(card),
|
||
}
|
||
|
||
|
||
__all__ = [
|
||
"GROUP_CODE",
|
||
"SCHEMAS",
|
||
"extract_geojson",
|
||
"fetch_doc_card",
|
||
"fetch_group_list",
|
||
"load_gisogd_permits",
|
||
"parse_cad_nums",
|
||
]
|