gendesign/backend/app/services/scrapers/domrf_catalog_object.py
bot-backend f7d4b5bccf
All checks were successful
Deploy / changes (push) Successful in 9s
Deploy / build-frontend (push) Has been skipped
Deploy / deploy-caddy (push) Has been skipped
Deploy / build-backend (push) Successful in 2m28s
Deploy / build-worker (push) Successful in 4m9s
Deploy / deploy (push) Successful in 1m25s
Deploy / deploy-status (push) Successful in 1s
Deploy / perimeter-smoke (push) Successful in 9s
fix(ptica): «объекта нет в БД» считается пропуском, а не сбоем (#2464) (#2974)
2026-08-20 11:59:39 +00:00

551 lines
25 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.

"""DOM.РФ catalog-OBJECT scraper (issue #297 sub-task 22d).
Fills ~25 NULL columns in domrf_kn_objects from public SSR catalog page:
https://наш.дом.рф/сервисы/каталог-новостроек/объект/{obj_id}
kn-API не возвращает: wall_type, energy_eff, ceiling_height_m, parking_*,
playground_*, finishing_variants_count, etc. — все эти поля есть в
__NEXT_DATA__ JSON блоке на странице каталога (Next.js SSR).
Uses BrowserSession from app.services.scrapers.stealth (Playwright + WAF bypass).
Fetches HTML, extracts __NEXT_DATA__ JSON, maps to DB columns,
UPDATE domrf_kn_objects WHERE obj_id = :id (не перетирает kn-API данные).
"""
from __future__ import annotations
import asyncio
import json
import logging
import re
from datetime import date
from typing import Any
from sqlalchemy import text
from sqlalchemy.orm import Session
from app.core.config import settings
from app.services.scrapers.stealth import BASE_URL, BrowserSession, WafBlockedError, jitter_sleep
logger = logging.getLogger(__name__)
# Сколько WAF-блоков ПОДРЯД прерывают батч (#2464). Одиночный блок бывает
# переходным (сессия перегреет cookies и восстановится), три подряд — стена.
_WAF_BREAKER_THRESHOLD = 3
# URL шаблон страницы объекта в каталоге DOM.РФ.
# Человекочитаемый вид: https://наш.дом.рф/сервисы/каталог-новостроек/объект/{obj_id}
CATALOG_OBJECT_PATH = "/сервисы/каталог-новостроек/объект/{obj_id}"
# JS snippet — аналог _FETCH_HTML_JS из domrf_catalog.py.
# Выполняется внутри живой Playwright-страницы, возвращает HTML текст.
_FETCH_HTML_JS = """
async ({url}) => {
try {
const r = await fetch(url, {credentials: 'include'});
const ctype = r.headers.get('content-type') || '';
const body = await r.text();
return {ok: r.ok, status: r.status, body, contentType: ctype};
} catch (e) {
return {ok: false, status: 0, body: String(e), contentType: ''};
}
}
"""
# UPDATE SQL — обновляет только catalog-derived поля.
# COALESCE гарантирует что NULL-значение не перетирает уже заполненное поле.
UPDATE_OBJECT_CATALOG_SQL = text(
"""
UPDATE domrf_kn_objects SET
obj_class = COALESCE(:obj_class, obj_class),
wall_type = COALESCE(:wall_type, wall_type),
energy_eff = COALESCE(:energy_eff, energy_eff),
section_count = COALESCE(:section_count, section_count),
parking_total_slots = COALESCE(:parking_total_slots, parking_total_slots),
guest_parking_inside_count = COALESCE(
:guest_parking_inside_count, guest_parking_inside_count
),
guest_parking_outside_count = COALESCE(
:guest_parking_outside_count, guest_parking_outside_count
),
ceiling_height_m = COALESCE(:ceiling_height_m, ceiling_height_m),
finishing_variants_count = COALESCE(:finishing_variants_count, finishing_variants_count),
has_free_planning = COALESCE(:has_free_planning, has_free_planning),
avg_flat_area_m2 = COALESCE(:avg_flat_area_m2, avg_flat_area_m2),
elevators_passenger_count = COALESCE(
:elevators_passenger_count, elevators_passenger_count
),
elevators_cargo_count = COALESCE(:elevators_cargo_count, elevators_cargo_count),
playground_kids_count = COALESCE(:playground_kids_count, playground_kids_count),
playground_sport_count = COALESCE(:playground_sport_count, playground_sport_count),
has_bike_paths = COALESCE(:has_bike_paths, has_bike_paths),
trash_areas_count = COALESCE(:trash_areas_count, trash_areas_count),
has_ramp = COALESCE(:has_ramp, has_ramp),
has_low_platforms = COALESCE(:has_low_platforms, has_low_platforms),
has_wheelchair_lift = COALESCE(:has_wheelchair_lift, has_wheelchair_lift),
first_floor_type = COALESCE(:first_floor_type, first_floor_type),
parking_provision_pct = COALESCE(:parking_provision_pct, parking_provision_pct),
project_published_at = COALESCE(:project_published_at, project_published_at),
project_declaration_num = COALESCE(:project_declaration_num, project_declaration_num),
domrf_score_infrastructure = COALESCE(
:domrf_score_infrastructure, domrf_score_infrastructure
),
domrf_score_transport = COALESCE(:domrf_score_transport, domrf_score_transport),
catalog_scraped_at = NOW()
WHERE obj_id = :obj_id
AND snapshot_date = :snapshot_date
"""
)
# ── Value helpers ─────────────────────────────────────────────────────────────
def _to_numeric_comma(s: Any) -> float | None:
"""Конвертировать строку с запятой-десятичным разделителем в float.
Примеры: "2,7" → 2.7; "2.7" → 2.7; "" → None; None → None.
"""
if s is None:
return None
raw = str(s).strip().replace(",", ".")
if not raw:
return None
try:
return float(raw)
except ValueError:
return None
def _to_date_ddmmyyyy(s: Any) -> date | None:
"""Конвертировать строку "DD.MM.YYYY" в date.
Примеры: "31.03.2025" → date(2025, 3, 31); "" → None; invalid → None.
"""
if not s:
return None
raw = str(s).strip()
if not raw:
return None
try:
parts = raw.split(".")
if len(parts) == 3:
return date(int(parts[2]), int(parts[1]), int(parts[0]))
except (ValueError, IndexError):
pass
return None
def _to_bool_int(v: Any) -> bool | None:
"""Конвертировать 0/1 (или любое int-like) в bool.
Примеры: 1 → True; 0 → False; None → None; 3 → True (>0).
"""
if v is None:
return None
try:
return int(v) > 0
except (ValueError, TypeError):
return None
def _to_bool_da_net(s: Any) -> bool | None:
"""Конвертировать "Да"/"Нет" строку в bool.
Примеры: "Да" → True; "Нет" → False; "" → None; None → None.
"""
if s is None:
return None
raw = str(s).strip().lower()
if raw == "да":
return True
if raw == "нет":
return False
return None
def _safe_int(v: Any) -> int | None:
"""Безопасная конвертация в int, None при ошибке."""
if v is None:
return None
try:
return int(v)
except (ValueError, TypeError):
return None
# ── HTML fetching ─────────────────────────────────────────────────────────────
async def fetch_catalog_object_html(session: BrowserSession, obj_id: int) -> str:
"""Получить SSR-HTML страницы объекта в каталоге DOM.РФ.
Использует тот же паттерн что fetch_catalog_html из domrf_catalog.py:
fetch() внутри живой Playwright-страницы — WAF-fingerprint идентичен браузеру.
Raises:
WafBlockedError: если вернулся не-HTML (JS-challenge или JSON).
RuntimeError: при 404 или исчерпании попыток.
"""
if session._page is None:
raise RuntimeError("BrowserSession not bootstrapped")
url = BASE_URL + CATALOG_OBJECT_PATH.format(obj_id=obj_id)
last_err: Exception | None = None
for attempt in range(5):
async with session._sem:
await jitter_sleep(300, 700)
try:
session._request_count += 1
result = await session._page.evaluate(_FETCH_HTML_JS, {"url": url})
except Exception as exc:
last_err = exc
logger.warning(
"catalog_object html evaluate err attempt=%d obj_id=%d: %r",
attempt,
obj_id,
exc,
)
await asyncio.sleep(2**attempt)
continue
status: int = result.get("status", 0)
body: str = result.get("body", "")
ctype: str = result.get("contentType", "")
if status in (429,) or status >= 500 or status == 0:
last_err = RuntimeError(f"transient status={status}")
logger.warning(
"catalog_object transient status=%d attempt=%d obj_id=%d, backing off",
status,
attempt,
obj_id,
)
await asyncio.sleep(2**attempt)
continue
if status == 404:
raise RuntimeError(f"catalog_object 404 for obj_id={obj_id}")
if status != 200:
raise RuntimeError(f"catalog_object http {status}: {body[:200]} obj_id={obj_id}")
# Проверяем что вернулся HTML, а не WAF JS-challenge.
is_html = "text/html" in ctype or "<!doctype" in body[:100].lower()
if body and not is_html:
raise WafBlockedError(
f"non-HTML response for obj_id={obj_id}: status={status} ctype={ctype!r}"
f" body[:120]={body[:120]!r}"
)
if not body:
raise RuntimeError(f"catalog_object empty body for obj_id={obj_id}")
return body
raise RuntimeError(f"catalog_object html max retries exhausted obj_id={obj_id}: {last_err!r}")
# ── __NEXT_DATA__ extraction ──────────────────────────────────────────────────
def extract_next_data(html: str) -> dict[str, Any]:
"""Извлечь JSON из тега <script id="__NEXT_DATA__"> в SSR HTML.
Raises:
ValueError: если тег не найден или JSON не парсится.
"""
match = re.search(
r'<script\s+id=["\']__NEXT_DATA__["\'][^>]*>(.+?)</script>',
html,
re.DOTALL,
)
if not match:
raise ValueError("__NEXT_DATA__ script tag not found in HTML")
raw_json = match.group(1).strip()
try:
return json.loads(raw_json) # type: ignore[no-any-return]
except json.JSONDecodeError as exc:
raise ValueError(f"__NEXT_DATA__ JSON parse error: {exc}") from exc
# ── Field mapping ─────────────────────────────────────────────────────────────
def parse_catalog_object(next_data: dict[str, Any]) -> dict[str, Any]:
"""Извлечь поля объекта из __NEXT_DATA__ и вернуть dict для UPDATE.
Все .get() безопасны — partial responses OK, отсутствующие поля = None.
Возвращает dict с bind-параметрами для UPDATE_OBJECT_CATALOG_SQL.
"""
pp: dict[str, Any] = next_data.get("props", {}).get("pageProps", {})
ai: dict[str, Any] = pp.get("additionalInfo") or {}
quart: dict[str, Any] = pp.get("quartography") or {}
indexes: dict[str, Any] = pp.get("indexes") or {}
decl: dict[str, Any] = pp.get("projectDeclaration") or {}
# first_floor_type: 1 = нежилой, 0 = жилой
first_floor_raw = quart.get("nonLivFirstFloor")
first_floor_type: str | None = None
if first_floor_raw is not None:
try:
first_floor_type = "нежилой" if int(first_floor_raw) == 1 else "жилой"
except (ValueError, TypeError):
pass
# elevators_cargo_count = cargoElevatorsCount + cargoPassengerElevatorCount
cargo = _safe_int(ai.get("cargoElevatorsCount"))
cargo_pass = _safe_int(ai.get("cargoPassengerElevatorCount"))
if cargo is not None or cargo_pass is not None:
elevators_cargo_count: int | None = (cargo or 0) + (cargo_pass or 0)
else:
elevators_cargo_count = None
return {
"obj_class": pp.get("buildingClass"),
"wall_type": pp.get("wallMaterial"),
"energy_eff": pp.get("objEnergyEfficiency"),
"section_count": _safe_int(quart.get("objLivElemEntrCnt")),
"parking_total_slots": _safe_int(pp.get("parkingCount")),
"guest_parking_inside_count": _safe_int(ai.get("objectParkingPlaces")),
"guest_parking_outside_count": _safe_int(ai.get("nearbyParkingPlaces")),
"ceiling_height_m": _to_numeric_comma(ai.get("ceilingHeight")),
"finishing_variants_count": _safe_int(pp.get("finishTypeCount")),
"has_free_planning": _to_bool_da_net(pp.get("freePlan")),
"avg_flat_area_m2": _to_numeric_comma(quart.get("objLivElemSqAvg")),
"elevators_passenger_count": _safe_int(ai.get("passengerElevatorsCount")),
"elevators_cargo_count": elevators_cargo_count,
"playground_kids_count": _safe_int(ai.get("playgroundsCount")),
"playground_sport_count": _safe_int(ai.get("sportsgroundCount")),
"has_bike_paths": _to_bool_int(ai.get("bicycleLane")),
"trash_areas_count": _safe_int(ai.get("trashAreaCount")),
"has_ramp": _to_bool_int(ai.get("ramp")),
"has_low_platforms": _to_bool_int(ai.get("curbLowering")),
"has_wheelchair_lift": _to_bool_int(ai.get("wheelchairElevatorsCount")),
"first_floor_type": first_floor_type,
"parking_provision_pct": _to_numeric_comma(ai.get("parkingAvailabilityPerc")),
"project_published_at": _to_date_ddmmyyyy(pp.get("publicationDate")),
"project_declaration_num": decl.get("number"),
"domrf_score_infrastructure": _safe_int(indexes.get("infrastructure")),
"domrf_score_transport": _safe_int(indexes.get("transport")),
# TODO: obj_checks (6 detailed checks) — separate investigation (task #21).
# pageProps.isChecked (bool), verificationId, verificationFlg available here
# but detailed per-check breakdown requires separate API investigation.
}
# ── DB write ──────────────────────────────────────────────────────────────────
async def scrape_catalog_object(
db: Session,
session: BrowserSession,
obj_id: int,
snapshot_date: date,
) -> bool | None:
"""Scrape одного объекта: fetch HTML → extract __NEXT_DATA__ → parse → UPDATE.
Использует SAVEPOINT (begin_nested) для изоляции per-row ошибок.
Логирует результат через logger.info.
Returns:
True — UPDATE затронул строку;
None — ПРОПУСК: строки (obj_id, snapshot_date) в БД нет. Это не сбой:
obj_id берутся из БД, но снимок мог смениться между выборкой и
UPDATE'ом. Раньше этот случай возвращал False и попадал в
счётчик failed вместе с настоящими ошибками, а объявленный в
контракте счётчик skipped всегда оставался нулём (#2464);
False — сбой: не скачалось, не распарсилось, упал UPDATE.
Третье состояние сделано через None, а не через новый Literal, намеренно:
прежние True/False сохраняют смысл, поэтому вызывающие и тесты, полагающиеся
на них, не меняются.
"""
logger.info("catalog_object scrape start obj_id=%d snapshot_date=%s", obj_id, snapshot_date)
try:
html = await fetch_catalog_object_html(session, obj_id)
except WafBlockedError:
# #2464: WAF-блок — не «этот объект не дошёл», а закрытая дверь. Раньше он
# гасился здесь и возвращался как обычная неудача, поэтому батч-цикл шёл
# дальше и слал ЖИВОЙ запрос на каждый оставшийся obj_id в уже забаненную
# сессию. Пробрасываем: решение принимает предохранитель в батче.
raise
except Exception as exc:
logger.warning("catalog_object fetch failed obj_id=%d: %s", obj_id, exc)
return False
try:
next_data = extract_next_data(html)
except ValueError as exc:
logger.warning("catalog_object extract_next_data failed obj_id=%d: %s", obj_id, exc)
return False
try:
data = parse_catalog_object(next_data)
except Exception as exc:
logger.warning("catalog_object parse failed obj_id=%d: %s", obj_id, exc)
return False
fields_extracted = len([v for v in data.values() if v is not None])
params: dict[str, Any] = {
"obj_id": obj_id,
"snapshot_date": snapshot_date,
**data,
}
try:
with db.begin_nested():
result = db.execute(UPDATE_OBJECT_CATALOG_SQL, params)
rows_affected: int = result.rowcount or 0
except Exception as exc:
logger.warning("catalog_object UPDATE failed obj_id=%d: %s", obj_id, exc)
return False
if rows_affected == 0:
logger.warning(
"catalog_object UPDATE 0 rows obj_id=%d snapshot_date=%s — not in DB?",
obj_id,
snapshot_date,
)
return None
logger.info(
"catalog_object scraped obj_id=%d fields=%d rows_updated=%d",
obj_id,
fields_extracted,
rows_affected,
)
return True
async def scrape_catalog_objects(
db: Session,
obj_ids: list[int],
snapshot_date: date,
region_code: int = 66,
) -> dict[str, int]:
"""Scrape списка объектов через один BrowserSession.
Запускает один BrowserSession на весь batch; jitter_sleep (300700 мс)
встроен в fetch_catalog_object_html для защиты от rate-limit.
Returns:
{"processed": N, "succeeded": N, "failed": N, "skipped": N}
"""
stats: dict[str, int] = {
"processed": 0,
"succeeded": 0,
"failed": 0,
"skipped": 0,
}
if not obj_ids:
logger.info("scrape_catalog_objects: empty list, nothing to do")
return stats
logger.info(
"scrape_catalog_objects: starting %d objects region=%d snapshot_date=%s",
len(obj_ids),
region_code,
snapshot_date,
)
async with BrowserSession(
region_code=region_code,
# Страницы каталога публичные — Basic auth не нужен
auth=None,
# #2445 D2 anti-ban: этот scraper ходит по тому же /сервисы/* path family,
# что вызвал DOM.РФ WAF hard-ban 2026-05-24 (#2443). Без явного override
# BrowserSession тихо наследует модульный дефолт stealth.py (concurrency=8,
# jitter 600-1500ms) — тот самый throttle-less профиль, что и раньше
# триггерил volume-ban. Переиспользуем throttled настройки, принятые
# scrape_kn.py после #1945 (settings.scrape_kn_* — общий anti-ban рычаг,
# НЕ специфичный для KN-sweep, несмотря на имя).
concurrency=settings.scrape_kn_browser_concurrency,
jitter_min_ms=settings.scrape_kn_request_jitter_min_ms,
jitter_max_ms=settings.scrape_kn_request_jitter_max_ms,
) as session:
# Warm-up: visit /сервисы/каталог-новостроек/ to obtain WAF cookies.
# DOM.РФ WAF (2026-05-24) блокирует SSR-страницы объектов без этих cookies.
# Idempotent — один вызов покрывает весь batch через этот BrowserSession.
await session.warm_up()
# #2464: предохранитель на серию WAF-блоков. Замер 20.08: в очереди 13200
# объектов из 13801, а DOM.РФ отдаёт страницу «Доступ заблокирован [403]»
# с капчей (#2443). Без предохранителя один прогон «Загрузить все» выдал бы
# 13200 живых запросов в забаненную сессию — ровно то, что углубляет бан
# (анти-бан-комментарий к BrowserSession выше про тот же path family).
# Порог не единица: одиночный блок бывает переходным, три подряд — стена.
consecutive_waf = 0
for obj_id in obj_ids:
stats["processed"] += 1
try:
ok = await scrape_catalog_object(db, session, obj_id, snapshot_date)
except WafBlockedError as exc:
consecutive_waf += 1
stats["failed"] += 1
logger.warning(
"catalog_object WAF blocked obj_id=%d (подряд %d/%d): %s",
obj_id,
consecutive_waf,
_WAF_BREAKER_THRESHOLD,
exc,
)
if consecutive_waf >= _WAF_BREAKER_THRESHOLD:
stats["aborted_on_waf"] = 1
logger.error(
"scrape_catalog_objects: %d WAF-блока подряд — прерываю батч,"
" обработано %d из %d",
consecutive_waf,
stats["processed"],
len(obj_ids),
)
break
continue
consecutive_waf = 0
if ok is None:
# Строки в БД нет — это пропуск, а не сбой (см. контракт выше).
stats["skipped"] += 1
continue
if ok:
stats["succeeded"] += 1
# Фиксируем сразу, а не одним commit'ом в конце (#2464). Раньше весь
# батч жил в одной незакоммиченной транзакции, и любой отказ ПОСЛЕ
# цикла — исключение в BrowserSession.__aexit__, снятие Celery-таски,
# перезапуск контейнера — обнулял все уже успешные UPDATE'ы.
# Это не теория: беговой режим здесь force=True («Загрузить все»),
# то есть SQL без LIMIT. На 20.08.2026 в очереди 13200 объектов из
# 13801 — многочасовой прогон, где отказ в конце стоил бы всего.
# SAVEPOINT внутри scrape_catalog_object к этому моменту уже снят,
# поэтому commit здесь корректен.
try:
db.commit()
except Exception:
db.rollback()
raise
else:
stats["failed"] += 1
# Финальный commit. Успешные строки зафиксированы по ходу цикла (см. выше), но
# этот вызов остаётся: он закрывает транзакцию, которую могли autobegin'ить
# неудачные итерации (их SAVEPOINT откатился, а внешняя транзакция открыта),
# и сохраняет прежнее поведение для вызывающих, которые на него полагались.
try:
db.commit()
except Exception:
db.rollback()
raise
logger.info(
"scrape_catalog_objects done: processed=%d succeeded=%d failed=%d skipped=%d",
stats["processed"],
stats["succeeded"],
stats["failed"],
stats["skipped"],
)
return stats