"""ДОМ.РФ капремонт open data loader: houses.year_built/material_walls (issue #2013). CONTEXT: houses.year_built заполнен только на 38%, houses.material_walls — на 0%. Существующий ГИС-ЖКХ loader (zhkh_flats_loader.py, мигр. 146/149) заполняет houses.zhkh_year (70%) и houses.zhkh_floors (70%), НО никогда не копирует их в houses.year_built/total_floors — это основной пробел. material_walls вообще ни разу не заполнялся никаким источником. Estimator (app/services/estimator.py) фильтрует когорту по `year_built BETWEEN ...` — реальные годы напрямую двигают точность оценки. ИСТОЧНИК: ДОМ.РФ капремонт open data (free, no auth), region 66 (Свердловская обл.): КР1.1 house registry (export/190) — zip → CSV, delimiter ';', UTF-8 BOM. Колонки: mkd_code, houseguid (ФИАС GUID дома), address, commission_year (год ввода, int), total_sq (decimal, запятая — «909,80»), number_floors_max (int), + служебные. КР1.2 constructive elements (export/275) — zip → CSV, тот же delimiter/encoding. LONG FORMAT: одна строка на конструктивный элемент на mkd_code; wall_material заполнен ТОЛЬКО на строке construction_element_type='фасад'. Чтобы получить материал стен на дом — группируем по mkd_code, берём первую строку с непустым wall_material. TLS: домен домрф.рф отдаёт RU-сертификат, который httpx с default trust store не верифицирует ("unable to get local issuer certificate") — та же ситуация, что и у sber_index.py. Открытые неавторизованные данные, поэтому используем verify=False, как остальные RU-gov loader'ы в этом репо. Матч staging→houses: domrf_kapremont.houseguid = COALESCE(houses.gar_house_guid, houses.house_fias_id, houses.zhkh_house_guid) — простой приоритетный COALESCE-джойн (без multi-attempt fallback-если-нет-матча — see PR discussion, "keep it simple/safe"). psycopg v3: SQL через `text(...)` использует `CAST(:x AS type)`, НИКОГДА `:x::type`. """ from __future__ import annotations import csv import io import logging import tempfile import zipfile from collections.abc import Iterator from dataclasses import dataclass from pathlib import Path import httpx from sqlalchemy import text from sqlalchemy.orm import Session logger = logging.getLogger(__name__) # ───────────────────────────────────────────────────────────────────────────── # Константы # ───────────────────────────────────────────────────────────────────────────── # export/190 (КР1.1) и export/275 (КР1.2) без доп. параметров отдают Свердловскую # обл. (region 66) — подтверждено инспекцией скачанных файлов. Для других регионов # понадобится region-gid параметр (не реализовано — вне скоупа #2013). KR11_URL = "https://xn--80adsazqn.xn--p1aee.xn--p1ai/opendata/export/190" KR12_URL = "https://xn--80adsazqn.xn--p1aee.xn--p1ai/opendata/export/275" DOWNLOAD_TIMEOUT_SEC = 180 # UPSERT в staging чанками (в одном SAVEPOINT), как zhkh_flats_loader.SAVEPOINT_CHUNK — # сбойный чанк откатывается изолированно, остальные доезжают. UPSERT_CHUNK_SIZE = 500 # ───────────────────────────────────────────────────────────────────────────── # Чистые хелперы парсинга (юнит-тестируются без сети/БД) # ───────────────────────────────────────────────────────────────────────────── def parse_decimal_comma(raw: str | None) -> float | None: """«909,80» / «909.80» / «» / None → float|None. Запятая — decimal separator ДОМ.РФ CSV.""" if raw is None: return None s = raw.strip() if not s: return None try: return float(s.replace(",", ".")) except ValueError: return None def parse_int_field(raw: str | None) -> int | None: """Устойчивый str → int|None (commission_year / number_floors_max). Пусто/None/нечисло → None. Терпит decimal-строки («5,0») через float-фоллбек — ДОМ.РФ CSV в принципе может так отдать целочисленные поля. """ if raw is None: return None s = raw.strip() if not s: return None try: return int(s) except ValueError: pass try: return int(float(s.replace(",", "."))) except ValueError: return None @dataclass(slots=True) class Kr11Row: """Одна строка КР1.1 house registry (house-per-row).""" mkd_code: str houseguid: str | None address: str | None commission_year: int | None number_floors_max: int | None total_sq: float | None @dataclass(slots=True) class DomrfHouse: """Объединённая КР1.1 + КР1.2(wall_material) строка — готова к UPSERT в staging.""" mkd_code: str houseguid: str | None address: str | None commission_year: int | None number_floors_max: int | None total_sq: float | None wall_material: str | None def _open_csv(path: str | Path) -> Iterator[dict[str, str]]: """csv.DictReader по ДОМ.РФ CSV: delimiter ';', UTF-8 BOM (utf-8-sig).""" with open(path, encoding="utf-8-sig", newline="") as f: reader = csv.DictReader(f, delimiter=";") yield from reader def parse_kr11_csv(path: str | Path) -> dict[str, Kr11Row]: """КР1.1 → dict по mkd_code. Строки без mkd_code пропускаются. Дубликат mkd_code (не ожидается — mkd_code PK в реестре ДОМ.РФ) — последняя строка побеждает (совпадает с семантикой ON CONFLICT DO UPDATE ниже). """ rows: dict[str, Kr11Row] = {} for raw in _open_csv(path): mkd_code = (raw.get("mkd_code") or "").strip() if not mkd_code: continue rows[mkd_code] = Kr11Row( mkd_code=mkd_code, houseguid=(raw.get("houseguid") or "").strip() or None, address=(raw.get("address") or "").strip() or None, commission_year=parse_int_field(raw.get("commission_year")), number_floors_max=parse_int_field(raw.get("number_floors_max")), total_sq=parse_decimal_comma(raw.get("total_sq")), ) return rows def parse_kr12_wall_materials(path: str | Path) -> dict[str, str]: """КР1.2 (long-format) → dict mkd_code → wall_material. wall_material непустой ТОЛЬКО на строке construction_element_type='фасад' — группировка не по этому полю, а просто "первая непустая wall_material на mkd_code" (устойчиво даже если разметка типа элемента когда-то изменится). """ materials: dict[str, str] = {} for raw in _open_csv(path): mkd_code = (raw.get("mkd_code") or "").strip() wall = (raw.get("wall_material") or "").strip() if mkd_code and wall and mkd_code not in materials: materials[mkd_code] = wall return materials def build_domrf_houses( kr11: dict[str, Kr11Row], wall_materials: dict[str, str] ) -> list[DomrfHouse]: """КР1.1 rows + КР1.2 wall_material lookup → список DomrfHouse (для UPSERT).""" return [ DomrfHouse( mkd_code=row.mkd_code, houseguid=row.houseguid, address=row.address, commission_year=row.commission_year, number_floors_max=row.number_floors_max, total_sq=row.total_sq, wall_material=wall_materials.get(row.mkd_code), ) for row in kr11.values() ] # ───────────────────────────────────────────────────────────────────────────── # HTTP: скачивание zip + извлечение CSV # ───────────────────────────────────────────────────────────────────────────── def _download_zip(url: str, *, client: httpx.Client) -> bytes: resp = client.get(url, timeout=DOWNLOAD_TIMEOUT_SEC) resp.raise_for_status() return resp.content def _extract_csv_from_zip(data: bytes, dest_dir: Path) -> Path: """Распаковывает первый *.csv из zip-байтов в dest_dir, возвращает путь.""" with zipfile.ZipFile(io.BytesIO(data)) as zf: csv_names = [n for n in zf.namelist() if n.lower().endswith(".csv")] if not csv_names: raise ValueError(f"zip не содержит .csv записей: {zf.namelist()!r}") extracted = zf.extract(csv_names[0], dest_dir) return Path(extracted) def fetch_domrf_csvs(dest_dir: Path, *, client: httpx.Client) -> tuple[Path, Path]: """Скачивает КР1.1 + КР1.2 zip'ы, распаковывает в dest_dir. Возвращает (kr11, kr12).""" kr11_path = _extract_csv_from_zip(_download_zip(KR11_URL, client=client), dest_dir) kr12_path = _extract_csv_from_zip(_download_zip(KR12_URL, client=client), dest_dir) return kr11_path, kr12_path # ───────────────────────────────────────────────────────────────────────────── # UPSERT staging (domrf_kapremont, мигр. 176) # ───────────────────────────────────────────────────────────────────────────── _UPSERT_SQL = text( """ INSERT INTO domrf_kapremont ( mkd_code, houseguid, address, commission_year, number_floors_max, total_sq, wall_material, loaded_at ) VALUES ( CAST(:mkd_code AS text), CAST(:houseguid AS text), CAST(:address AS text), CAST(:commission_year AS smallint), CAST(:number_floors_max AS int), CAST(:total_sq AS numeric), CAST(:wall_material AS text), now() ) ON CONFLICT (mkd_code) DO UPDATE SET houseguid = EXCLUDED.houseguid, address = EXCLUDED.address, commission_year = EXCLUDED.commission_year, number_floors_max = EXCLUDED.number_floors_max, total_sq = EXCLUDED.total_sq, wall_material = EXCLUDED.wall_material, loaded_at = now() WHERE domrf_kapremont.houseguid IS DISTINCT FROM EXCLUDED.houseguid OR domrf_kapremont.address IS DISTINCT FROM EXCLUDED.address OR domrf_kapremont.commission_year IS DISTINCT FROM EXCLUDED.commission_year OR domrf_kapremont.number_floors_max IS DISTINCT FROM EXCLUDED.number_floors_max OR domrf_kapremont.total_sq IS DISTINCT FROM EXCLUDED.total_sq OR domrf_kapremont.wall_material IS DISTINCT FROM EXCLUDED.wall_material """ ) def _chunk_houses(items: list[DomrfHouse], size: int) -> Iterator[list[DomrfHouse]]: for i in range(0, len(items), size): yield items[i : i + size] def upsert_domrf_kapremont( db: Session, houses: list[DomrfHouse], *, chunk_size: int = UPSERT_CHUNK_SIZE ) -> int: """UPSERT списка DomrfHouse в domrf_kapremont, чанками по SAVEPOINT. Не коммитит (caller). Идемпотентно (IS DISTINCT FROM gate в _UPSERT_SQL — повторный прогон с теми же данными не трогает уже актуальные строки). Сбойный чанк откатывается изолированно (та же схема, что zhkh_flats_loader.match_houses_to_zhkh). """ upserted = 0 for chunk in _chunk_houses(houses, chunk_size): try: with db.begin_nested(): for h in chunk: res = db.execute( _UPSERT_SQL, { "mkd_code": h.mkd_code, "houseguid": h.houseguid, "address": h.address, "commission_year": h.commission_year, "number_floors_max": h.number_floors_max, "total_sq": h.total_sq, "wall_material": h.wall_material, }, ) upserted += res.rowcount except Exception: logger.warning( "domrf_kapremont upsert: чанк из %d строк сбойнул (откат savepoint)", len(chunk), exc_info=True, ) return upserted def load_domrf_kapremont( db: Session, *, kr11_path: str | Path | None = None, kr12_path: str | Path | None = None, work_dir: str | Path | None = None, chunk_size: int = UPSERT_CHUNK_SIZE, dry_run: bool = False, ) -> dict[str, int]: """Скачивает (если пути не заданы) КР1.1+КР1.2, парсит, UPSERT в domrf_kapremont. kr11_path/kr12_path — локальные CSV (пропустить скачивание; используется тестами и ops-дебагом). work_dir — куда распаковывать скачанные zip (по умолчанию — временный каталог, удаляется после загрузки). dry_run — парс происходит, но НИ ОДНОЙ записи в БД не делается. Не коммитит (caller). """ tmp_ctx: tempfile.TemporaryDirectory[str] | None = None if kr11_path is None or kr12_path is None: if work_dir is not None: dest = Path(work_dir) dest.mkdir(parents=True, exist_ok=True) else: tmp_ctx = tempfile.TemporaryDirectory() dest = Path(tmp_ctx.name) # verify=False: домрф.рф отдаёт RU-сертификат, не верифицируемый дефолтным # trust store'ом (та же ситуация, что sber_index.py). Публичные open data, # без auth/PII — приемлемо skip'нуть TLS-верификацию. with httpx.Client(timeout=DOWNLOAD_TIMEOUT_SEC, verify=False) as client: if kr11_path is None: kr11_path = _extract_csv_from_zip(_download_zip(KR11_URL, client=client), dest) if kr12_path is None: kr12_path = _extract_csv_from_zip(_download_zip(KR12_URL, client=client), dest) try: kr11 = parse_kr11_csv(kr11_path) wall_materials = parse_kr12_wall_materials(kr12_path) houses = build_domrf_houses(kr11, wall_materials) upserted = 0 if dry_run else upsert_domrf_kapremont(db, houses, chunk_size=chunk_size) finally: if tmp_ctx is not None: tmp_ctx.cleanup() result = { "kr11_rows": len(kr11), "kr12_wall_rows": len(wall_materials), "houses_built": len(houses), "upserted": upserted, } logger.info("domrf_kapremont load DONE (dry_run=%s): %s", dry_run, result) return result # ───────────────────────────────────────────────────────────────────────────── # Backfill houses.year_built/material_walls/total_floors # ───────────────────────────────────────────────────────────────────────────── # Шаг 1: дома с domrf-матчем — year_built/material_walls/total_floors из domrf, # С ЗАФОЛЖЕННЫМ zhkh_year/zhkh_floors фоллбеком внутри ТОЙ ЖЕ строки COALESCE # (COALESCE(year_built, domrf.commission_year, zhkh_year) — ровно как в issue). # Только заполняет NULL — никогда не перезаписывает существующий non-null year_built. _BACKFILL_HOUSES_FROM_DOMRF_SQL = text( """ UPDATE houses h SET year_built = COALESCE(h.year_built, d.commission_year, h.zhkh_year), material_walls = COALESCE(h.material_walls, d.wall_material), total_floors = COALESCE(h.total_floors, d.number_floors_max, h.zhkh_floors) FROM domrf_kapremont d WHERE d.houseguid = COALESCE(h.gar_house_guid, h.house_fias_id, h.zhkh_house_guid) AND ( (h.year_built IS NULL AND COALESCE(d.commission_year, h.zhkh_year) IS NOT NULL) OR (h.material_walls IS NULL AND d.wall_material IS NOT NULL) OR (h.total_floors IS NULL AND COALESCE(d.number_floors_max, h.zhkh_floors) IS NOT NULL) ) """ ) _BACKFILL_HOUSES_FROM_DOMRF_COUNT_SQL = text( """ SELECT count(*) FROM houses h JOIN domrf_kapremont d ON d.houseguid = COALESCE(h.gar_house_guid, h.house_fias_id, h.zhkh_house_guid) WHERE (h.year_built IS NULL AND COALESCE(d.commission_year, h.zhkh_year) IS NOT NULL) OR (h.material_walls IS NULL AND d.wall_material IS NOT NULL) OR (h.total_floors IS NULL AND COALESCE(d.number_floors_max, h.zhkh_floors) IS NOT NULL) """ ) # Шаг 2: ЛЮБОЙ дом (включая без domrf-матча вообще — шаг 1 INNER JOIN их не трогает) # добивается zhkh_year/zhkh_floors, если ещё NULL. Домам, уже тронутым шагом 1, шаг 2 # ничего не меняет: их COALESCE в шаге 1 уже включал zhkh_* как fallback, поэтому # гейт `IS NULL AND zhkh_* IS NOT NULL` здесь для них ложен — шаги взаимоисключающи # (не двойной счёт при суммировании houses_updated). _BACKFILL_HOUSES_ZHKH_FALLBACK_SQL = text( """ UPDATE houses SET year_built = COALESCE(year_built, zhkh_year), total_floors = COALESCE(total_floors, zhkh_floors) WHERE (year_built IS NULL AND zhkh_year IS NOT NULL) OR (total_floors IS NULL AND zhkh_floors IS NOT NULL) """ ) _BACKFILL_HOUSES_ZHKH_FALLBACK_COUNT_SQL = text( """ SELECT count(*) FROM houses WHERE (year_built IS NULL AND zhkh_year IS NOT NULL) OR (total_floors IS NULL AND zhkh_floors IS NOT NULL) """ ) def backfill_houses_from_domrf(db: Session, *, dry_run: bool = False) -> dict[str, int]: """COALESCE(year_built, domrf.commission_year, zhkh_year) + material_walls + total_floors. Двухшаговый UPDATE (мирроринг zhkh guid-match → cadastre-fallback паттерна из zhkh_flats_loader.py): 1) дома с domrf-матчем (COALESCE(gar_house_guid, house_fias_id, zhkh_house_guid) = domrf_kapremont.houseguid) — заполняются из domrf, с zhkh_year/zhkh_floors фоллбеком внутри той же COALESCE. 2) ЛЮБОЙ дом (в т.ч. без domrf-матча) — добивает year_built/total_floors из zhkh_year/zhkh_floors, если шаг 1 их не тронул. Только COALESCE (заполняет NULL, никогда не перезаписывает existing non-null). Не коммитит (caller). dry_run — ноль записей, только SELECT count(*) по тем же предикатам. """ if dry_run: domrf_matched = db.execute(_BACKFILL_HOUSES_FROM_DOMRF_COUNT_SQL).scalar_one() zhkh_fallback = db.execute(_BACKFILL_HOUSES_ZHKH_FALLBACK_COUNT_SQL).scalar_one() result = { "domrf_matched": domrf_matched, "zhkh_fallback": zhkh_fallback, "houses_updated": 0, } logger.info("backfill_houses_from_domrf DRY-RUN: %s", result) return result domrf_updated = db.execute(_BACKFILL_HOUSES_FROM_DOMRF_SQL).rowcount zhkh_updated = db.execute(_BACKFILL_HOUSES_ZHKH_FALLBACK_SQL).rowcount result = { "domrf_matched": domrf_updated, "zhkh_fallback": zhkh_updated, "houses_updated": domrf_updated + zhkh_updated, } logger.info("backfill_houses_from_domrf DONE: %s", result) return result # ───────────────────────────────────────────────────────────────────────────── # Propagate houses.year_built → listings.year_built # ───────────────────────────────────────────────────────────────────────────── # Только NULL listings.year_built, только когда houses.year_built уже известен — # никогда не перезаписывает существующий listings.year_built (source-provided данные # приоритетнее houses-агрегата). _PROPAGATE_LISTINGS_YEAR_SQL = text( """ UPDATE listings l SET year_built = h.year_built FROM houses h WHERE l.house_id_fk = h.id AND l.year_built IS NULL AND h.year_built IS NOT NULL """ ) _PROPAGATE_LISTINGS_YEAR_COUNT_SQL = text( """ SELECT count(*) FROM listings l JOIN houses h ON l.house_id_fk = h.id WHERE l.year_built IS NULL AND h.year_built IS NOT NULL """ ) def propagate_listings_year_from_houses(db: Session, *, dry_run: bool = False) -> dict[str, int]: """UPDATE listings.year_built = houses.year_built где listing.year_built ещё NULL. Не коммитит (caller). dry_run — ноль записей, только SELECT count(*) по тому же предикату (ключ результата — would_update вместо listings_updated). """ if dry_run: would_update = db.execute(_PROPAGATE_LISTINGS_YEAR_COUNT_SQL).scalar_one() logger.info("propagate_listings_year_from_houses DRY-RUN: would_update=%d", would_update) return {"listings_updated": 0, "would_update": would_update} updated = db.execute(_PROPAGATE_LISTINGS_YEAR_SQL).rowcount logger.info("propagate_listings_year_from_houses DONE: listings_updated=%d", updated) return {"listings_updated": updated}