"""Celery task: harvest РИАСУРТ Свердл (folderId 1224) → ``riasurt_sverdl`` (#108). 14 ключевых слоёв РИАСУРТ Свердл (ПЗЗ/планирование, инженерные ЗОУИТ, risk, opportunity) по агломерации ЕКБ — ОКРАИНЫ (Берёзовский, В.Пышма, Среднеуральск, Арамиль, Сысерть), НЕ ЕКБ-сити (тот не интегрирован в ФГИС ТП). Foundation для multi-city scaling Site Finder. Источник — NSPD aeggis/v4 WMS (тот же, что федеральные слои), layerId из регионального каталога РИАСУРТ (см. ``nspd_client.RIASURT_SVERDL_LAYERS``). Bulk-метод клиента ``get_riasurt_sverdl_in_bbox`` grid-walk'ит каждый слой в bbox одного МО. Mirror conventions (ird_harvest.py / backend.md): • SessionLocal() + try/finally close; logger (не print). • SAVEPOINT per-row (``with db.begin_nested():``) — одна битая фича не валит батч; commit в конце. • CAST(:x AS type) — НИКОГДА :x::type (psycopg v3). Геометрия 3857 → MULTIPOLYGON. • geom_msk66 — generated STORED в БД (миграция 163), здесь НЕ пишем. • Сбой одного слоя/МО не валит весь прогон. Расписание — ежеквартально (beat_schedule.py: riasurt-sverdl-harvest-quarterly): региональные градостроительные зоны меняются медленно. ВНИМАНИЕ: bbox 5 МО — ПЛЕЙСХОЛДЕРЫ (см. MO_BBOXES TODO). Точные bbox резолвятся ПОСЛЕ деплоя через NSPD search по границам МО / ручную выгрузку. Текущие значения — грубые приближения в EPSG:3857 (Web Mercator), достаточные для smoke, НО НЕ для прод-harvest. """ from __future__ import annotations import json import logging from datetime import UTC, datetime from typing import Any from sqlalchemy import text from sqlalchemy.orm import Session from app.core.db import SessionLocal from app.services.scrapers.nspd_client import ( RIASURT_SVERDL_LAYERS, NSPDClient, riasurt_layer_topic, ) from app.workers.celery_app import celery_app logger = logging.getLogger(__name__) # ── bbox 5 МО агломерации ЕКБ (EPSG:3857) ──────────────────────────────────── # TODO(#108): ПЛЕЙСХОЛДЕРЫ — грубые приближения вокруг центров МО (±~6 км в 3857). # Точные административные bbox резолвятся ПОСЛЕ деплоя (NSPD search по границе МО / # ручная выгрузка границ из РИАСУРТ). НЕ запускать harvest_all_riasurt_sverdl на проде # до уточнения — иначе grid-walk промахнётся мимо застройки. Формат: (xmin, ymin, xmax, ymax). # # Центры (приблизительно, WGS84 → 3857): # Берёзовский ~56.91N 60.81E В.Пышма ~56.97N 60.58E # Среднеуральск ~56.99N 60.48E Арамиль ~56.69N 60.83E # Сысерть ~56.50N 60.82E _PLACEHOLDER_HALF_M = 6000.0 # ±6 км вокруг центра — заведомо грубо, заменить на реальные границы MO_BBOXES: dict[str, tuple[float, float, float, float]] = { # TODO(#108): заменить плейсхолдеры на реальные административные bbox МО. "Берёзовский": (6770000.0, 7710000.0, 6782000.0, 7722000.0), "Верхняя Пышма": (6745000.0, 7720000.0, 6757000.0, 7732000.0), "Среднеуральск": (6734000.0, 7723000.0, 6746000.0, 7735000.0), "Арамиль": (6772000.0, 7677000.0, 6784000.0, 7689000.0), "Сысерть": (6771000.0, 7639000.0, 6783000.0, 7651000.0), } # UPSERT фичи РИАСУРТ. Геометрия НСПД (GeoJSON 3857) → ST_GeomFromGeoJSON → SetSRID 3857 → # ST_Multi (приводим Polygon к MULTIPOLYGON под колонку). geom_msk66 — generated в БД. # Нет стабильного natural-ключа дедупа у РИАСУРТ-фич → INSERT без ON CONFLICT; перед # harvest_all очищаем строки по mo_name (re-harvest квартала идемпотентен на уровне МО). _INSERT_SQL = text( """ INSERT INTO riasurt_sverdl ( source_layer_id, layer_topic, mo_name, obshnz, description, raw_props, geom, fetched_at ) VALUES ( CAST(:source_layer_id AS integer), CAST(:layer_topic AS text), CAST(:mo_name AS text), CAST(:obshnz AS text), CAST(:description AS text), CAST(:raw_props AS jsonb), ST_Multi(ST_SetSRID(ST_GeomFromGeoJSON(CAST(:geojson AS text)), 3857)), CAST(:fetched_at AS timestamptz) ) """ ) _DELETE_MO_SQL = text("DELETE FROM riasurt_sverdl WHERE mo_name = CAST(:mo_name AS text)") def _extract_obshnz(props: dict[str, Any]) -> str | None: """reg_numb_border / индекс зоны из properties фичи (когда есть).""" opts = props.get("options") or {} val = ( props.get("reg_numb_border") or opts.get("reg_numb_border") or props.get("obshnz") or opts.get("obshnz") ) return str(val) if val is not None else None def _extract_description(props: dict[str, Any]) -> str | None: """Человекочитаемое название/описание фичи.""" opts = props.get("options") or {} val = ( props.get("label") or props.get("name") or props.get("descr") or opts.get("label") or opts.get("descr") ) return str(val) if val is not None else None def _insert_feature( db: Session, *, layer_id: int, mo_name: str, feature: Any, fetched_at: str, ) -> bool: """INSERT одной РИАСУРТ-фичи под SAVEPOINT. True если записана. Пропускает фичи без geometry — одна битая не валит батч. """ geom = feature.geometry if not geom: return False props = feature.properties or {} try: with db.begin_nested(): db.execute( _INSERT_SQL, { "source_layer_id": layer_id, "layer_topic": riasurt_layer_topic(layer_id), "mo_name": mo_name, "obshnz": _extract_obshnz(props), "description": _extract_description(props), "raw_props": json.dumps(props, ensure_ascii=False), "geojson": json.dumps(geom), "fetched_at": fetched_at, }, ) return True except Exception as exc: logger.warning( "riasurt_sverdl: insert failed layer=%s mo=%s: %s", layer_id, mo_name, exc, ) return False @celery_app.task( name="tasks.riasurt_sverdl_harvest.harvest_riasurt_sverdl_for_mo", rate_limit="1/s", ) def harvest_riasurt_sverdl_for_mo( mo_name: str, mo_bbox: list[float] | tuple[float, float, float, float], layers: list[int] | None = None, ) -> dict[str, int]: """Harvest 14 слоёв РИАСУРТ Свердл в bbox одного МО → riasurt_sverdl. Перед записью очищает прошлые строки этого МО (re-harvest идемпотентен на уровне МО — у РИАСУРТ-фич нет стабильного natural-ключа дедупа). Args: mo_name: имя МО (Берёзовский / Верхняя Пышма / ...). Денормализуется в строки. mo_bbox: (xmin, ymin, xmax, ymax) в EPSG:3857. layers: список layerId РИАСУРТ. None → все 14 ключевых. Returns: {'features': M, 'layers': L} — записано фич / обработано слоёв. """ fetched_at = datetime.now(UTC).isoformat() bbox = ( float(mo_bbox[0]), float(mo_bbox[1]), float(mo_bbox[2]), float(mo_bbox[3]), ) client = NSPDClient() db = SessionLocal() n_features = 0 n_layers = 0 try: per_layer = client.get_riasurt_sverdl_in_bbox(bbox, layers) # Идемпотентность на уровне МО: чистим прошлый проход перед записью нового. with db.begin_nested(): db.execute(_DELETE_MO_SQL, {"mo_name": mo_name}) for layer_id, feats in per_layer.items(): n_layers += 1 for feat in feats: if _insert_feature( db, layer_id=layer_id, mo_name=mo_name, feature=feat, fetched_at=fetched_at, ): n_features += 1 db.commit() logger.info( "riasurt_sverdl: МО=%s слоёв=%d фич=%d", mo_name, n_layers, n_features, ) except Exception: db.rollback() logger.exception("riasurt_sverdl: harvest_riasurt_sverdl_for_mo упал, mo=%s", mo_name) raise finally: db.close() return {"features": n_features, "layers": n_layers} @celery_app.task(name="tasks.riasurt_sverdl_harvest.harvest_all_riasurt_sverdl") def harvest_all_riasurt_sverdl(layers: list[int] | None = None) -> dict[str, int]: """Прогон по всем 5 МО агломерации ЕКБ (MO_BBOXES) → riasurt_sverdl. Вызывает harvest_riasurt_sverdl_for_mo синхронно по каждому МО (rate_limit самой per-MO задачи действует при .delay()-вызове; здесь синхронный прогон — каждый МО последовательно, чтобы один beat-тик не запускал 5 параллельных grid-walk'ов). Args: layers: список layerId. None → все 14 ключевых (RIASURT_SVERDL_LAYERS). Returns: {'mo': N, 'features': M} — обработано МО / суммарно записано фич. """ n_mo = 0 n_features = 0 for mo_name, bbox in MO_BBOXES.items(): try: res = harvest_riasurt_sverdl_for_mo(mo_name, bbox, layers) n_features += int(res.get("features", 0)) n_mo += 1 except Exception: logger.exception("riasurt_sverdl: МО %s упал — продолжаем остальные", mo_name) continue logger.info( "riasurt_sverdl: harvest_all завершён, МО=%d (из %d) фич=%d", n_mo, len(MO_BBOXES), n_features, ) return {"mo": n_mo, "features": n_features} __all__ = [ "MO_BBOXES", "harvest_all_riasurt_sverdl", "harvest_riasurt_sverdl_for_mo", ] # Layers catalog re-export (для удобства тестов / интроспекции). _RIASURT_LAYER_COUNT = len(RIASURT_SVERDL_LAYERS)