gendesign/tradein-mvp/backend/app/tasks/yandex_newbuilding_sweep.py
bot-backend eb564cfc01 feat(tradein/scraper-kit): migrate backfill-task imports to kit, Group C (#2310)
Migrates legacy app.services.scrapers.* imports to scraper_kit equivalents for
house_imv_backfill.py, avito_detail_backfill.py, cian_history_backfill.py,
ekb_geoportal_ingest.py, and yandex_detail_backfill.py, proving parity via
tests/support/parity.assert_parity per the epic's gate (#2304).

newbuilding_enrich_backfill.py and yandex_newbuilding_sweep.py are left fully
on legacy imports: both call into scraper_kit.providers.{cian,yandex}.newbuilding,
which construct BrowserFetcher(source=...) without the now-mandatory endpoint=
kwarg (issue #2322, verified still open against the actual provider source, not
just issue status) -- no caller-side fix is possible, matching Group A's (#2305)
precedent for the same bug in admin.py.

Config-gated kit function footguns found and fixed (Group B #2306 pattern):
  - avito fetch_detail's backconnect-on-403 retry silently drops when config=
    is omitted -- now passes config=RealScraperConfig() explicitly, with a
    regression test proving the gate.
  - kit AvitoScraper's constructor now requires ScraperConfig positionally --
    wired via RealScraperConfig(), which _rotate_ip() reads for
    avito_proxy_rotate_url.
  - BrowserFetcher(source=...) call sites (house_imv_backfill, avito/cian
    detail backfills) now pass the mandatory endpoint=settings.browser_http_endpoint.

New footgun discovered (NOT #2322, flagged for follow-up): kit's
build_warmed_session() builds its curl_cffi session via _build_detail_session()
with no config parameter at all, unlike fetch_detail -- migrating it would
silently drop the sticky MGTS-proxy egress on avito_detail_backfill's
warm-batch path (the prod default). Left build_warmed_session/
_AVITO_WARM_SEARCH_URL on legacy imports, documented inline.

Refs #2310
2026-07-04 01:19:34 +03:00

432 lines
16 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.

"""Yandex Newbuilding sweep-task (#974): enrichment ЖК в market.yandex_jk_enrichment.
Context
-------
309 house_sources строк с ext_source='yandex_realty_nb' имеют ext_id (= yandex_jk_id),
но у большинства NULL yandex_jk_slug и не заполнены данные в market.yandex_jk_enrichment.
Для каждого дома цепочка:
1. resolve_yandex_jk_slug(jk_id) → slug (через BrowserFetcher SERP)
Результат сохраняется в houses.yandex_jk_slug (SAVEPOINT) — resumable.
2. YandexNewbuildingScraper.fetch_jk(slug, jk_id) → YandexNewbuildingInfo
Через BrowserFetcher (tradein-browser camoufox) — единственный рабочий путь.
3. UPSERT в market.yandex_jk_enrichment (ON CONFLICT (ext_id) DO UPDATE).
4. UPDATE houses.yandex_jk_id WHERE yandex_jk_slug = slug.
Idempotency
-----------
- force=False: пропускает дома у которых уже есть строка в market.yandex_jk_enrichment.
- UPSERT через ON CONFLICT (ext_id) DO UPDATE — безопасен при повторном запуске.
- SAVEPOINT per house — один сбойный fetch не прерывает батч.
- dry_run: подсчёт популяции без fetch.
Anti-bot / resilience
---------------------
- request_delay_sec (default из get_scraper_delay('yandex_realty_nb')) с ±20% jitter.
- SAVEPOINT per house: resolve-фаза и enrich-фаза — отдельные savepoint'ы.
Execution
---------
- Чистый async callable — NO Celery. Триггер: admin endpoint или scheduler.
- psycopg v3 conventions: CAST(:x AS type) в SQL (никогда ::`); logger (никогда print).
"""
from __future__ import annotations
import asyncio
import json
import logging
import random
import time
from dataclasses import dataclass, field, fields
from sqlalchemy import text
from sqlalchemy.orm import Session
from app.services.scraper_settings import get_scraper_delay
logger = logging.getLogger(__name__)
__all__ = [
"YandexNewbuildingSweepResult",
"count_yandex_newbuilding_houses",
"enrich_yandex_newbuilding_sweep",
]
_ENRICHMENT_SOURCE = "yandex_realty_nb"
@dataclass
class YandexNewbuildingSweepResult:
"""Per-run счётчики yandex newbuilding sweep."""
# Sizing (независимо от limit).
total: int = 0 # дома с ext_source='yandex_realty_nb'
fetchable: int = 0 # из них: с yandex_jk_slug ИЛИ с ext_id
pending: int = 0 # fetchable И ещё не обогащены (force=False)
# Processing (ограничен limit).
processed: int = 0
skipped_already_enriched: int = 0
succeeded: int = 0
resolved_slug: int = 0 # ext_id → slug разрезолвлен + сохранён
failed_resolve: int = 0 # slug не удалось разрезолвить
failed_fetch: int = 0 # fetch_jk вернул None / упал
rows_inserted: int = 0 # строк в market.yandex_jk_enrichment (новых/обновлённых)
duration_sec: float = field(default=0.0)
def to_dict(self) -> dict[str, int | float]:
return {f.name: getattr(self, f.name) for f in fields(self)}
# ── SQL ────────────────────────────────────────────────────────────────────────
_COUNT_TOTAL = """
SELECT COUNT(DISTINCT h.id)
FROM houses h
JOIN house_sources hs ON hs.house_id = h.id
WHERE hs.ext_source = 'yandex_realty_nb'
"""
_COUNT_FETCHABLE = """
SELECT COUNT(DISTINCT h.id)
FROM houses h
JOIN house_sources hs ON hs.house_id = h.id
WHERE hs.ext_source = 'yandex_realty_nb'
AND (h.yandex_jk_slug IS NOT NULL OR hs.ext_id IS NOT NULL)
"""
_COUNT_PENDING = """
SELECT COUNT(DISTINCT h.id)
FROM houses h
JOIN house_sources hs ON hs.house_id = h.id
WHERE hs.ext_source = 'yandex_realty_nb'
AND (h.yandex_jk_slug IS NOT NULL OR hs.ext_id IS NOT NULL)
AND NOT EXISTS (
SELECT 1 FROM market.yandex_jk_enrichment e WHERE e.ext_id = hs.ext_id
)
"""
# DISTINCT ON (h.id) — у дома может быть >1 house_sources строки; берём любую ext_id.
_SELECT_PENDING_HOUSES = """
SELECT DISTINCT ON (h.id)
h.id AS house_id,
h.yandex_jk_slug,
h.yandex_jk_id,
hs.ext_id
FROM houses h
JOIN house_sources hs ON hs.house_id = h.id
WHERE hs.ext_source = 'yandex_realty_nb'
AND (h.yandex_jk_slug IS NOT NULL OR hs.ext_id IS NOT NULL)
AND (
CAST(:force AS boolean) = TRUE
OR NOT EXISTS (
SELECT 1 FROM market.yandex_jk_enrichment e WHERE e.ext_id = hs.ext_id
)
)
ORDER BY h.id, hs.ext_id NULLS LAST
LIMIT :lim
"""
_UPDATE_SLUG = """
UPDATE houses
SET yandex_jk_slug = CAST(:slug AS text)
WHERE id = CAST(:hid AS bigint)
"""
_UPSERT_ENRICHMENT = """
INSERT INTO market.yandex_jk_enrichment (
ext_id, name, developer_name, address,
lat, lon, rating, ratings_count, text_reviews_count,
raw_payload, updated_at
) VALUES (
CAST(:ext_id AS text),
CAST(:name AS text),
CAST(:developer_name AS text),
CAST(:address AS text),
CAST(:lat AS float8),
CAST(:lon AS float8),
CAST(:rating AS float4),
CAST(:ratings_count AS int),
CAST(:text_reviews_count AS int),
CAST(:raw_payload AS jsonb),
NOW()
)
ON CONFLICT (ext_id) DO UPDATE SET
name = EXCLUDED.name,
developer_name = EXCLUDED.developer_name,
address = EXCLUDED.address,
lat = EXCLUDED.lat,
lon = EXCLUDED.lon,
rating = EXCLUDED.rating,
ratings_count = EXCLUDED.ratings_count,
text_reviews_count = EXCLUDED.text_reviews_count,
raw_payload = EXCLUDED.raw_payload,
updated_at = NOW()
"""
def count_yandex_newbuilding_houses(db: Session) -> dict[str, int]:
"""Sizing: сколько yandex_realty_nb домов существует / fetchable / pending.
Дешёвый (3 COUNT) — безопасен для dry_run.
"""
total = int(db.execute(text(_COUNT_TOTAL)).scalar_one())
fetchable = int(db.execute(text(_COUNT_FETCHABLE)).scalar_one())
pending = int(db.execute(text(_COUNT_PENDING)).scalar_one())
return {"total": total, "fetchable": fetchable, "pending": pending}
async def enrich_yandex_newbuilding_sweep(
db: Session,
*,
limit: int = 5,
force: bool = False,
request_delay_sec: float | None = None,
dry_run: bool = False,
city: str = "ekaterinburg",
) -> YandexNewbuildingSweepResult:
"""Enrichment sweep для yandex_realty_nb домов → market.yandex_jk_enrichment.
Per house:
- resolve slug если NULL (persist в houses.yandex_jk_slug, SAVEPOINT)
- fetch_jk через BrowserFetcher
- UPSERT market.yandex_jk_enrichment (ON CONFLICT ext_id)
- UPDATE houses.yandex_jk_id (если изменился)
Args:
db: tradein-mvp SQLAlchemy session (НЕ gendesign DB).
limit: max домов за один прогон. Default 5 — bounded start.
force: переобработать уже обогащённые (UPSERT-safe). Default False.
request_delay_sec: пауза между домами. None → get_scraper_delay('yandex_realty_nb').
dry_run: считать популяцию без fetch/write.
city: город для Yandex Realty URL (ekaterinburg).
Returns:
YandexNewbuildingSweepResult со счётчиками.
"""
# NOT migrated to scraper_kit (issue #2310, Group C): both
# scraper_kit.providers.yandex.newbuilding.YandexNewbuildingScraper.fetch_jk()
# AND resolve_yandex_jk_slug() construct BrowserFetcher(source="yandex")
# WITHOUT the now-mandatory endpoint= kwarg — every call raises TypeError,
# with no caller-side fix possible (the kit provider functions don't expose
# a config/endpoint hook at all). This is issue #2322, verified STILL OPEN
# by reading the actual provider source (not just the issue's open/closed
# status) at the time of this migration. Left on legacy entirely, exactly
# like Group A (#2305) did for the same bug in admin.py's debug endpoints.
from app.services.scrapers.yandex_newbuilding import (
YandexNewbuildingScraper,
resolve_yandex_jk_slug,
)
result = YandexNewbuildingSweepResult()
t0 = time.time()
delay = (
request_delay_sec
if request_delay_sec is not None
else get_scraper_delay(_ENRICHMENT_SOURCE)
)
# ── Sizing ──────────────────────────────────────────────────────────────
sizing = count_yandex_newbuilding_houses(db)
result.total = sizing["total"]
result.fetchable = sizing["fetchable"]
result.pending = sizing["pending"]
rows = db.execute(text(_SELECT_PENDING_HOUSES), {"force": force, "lim": limit}).mappings().all()
logger.info(
"yandex-nb-sweep: total=%d fetchable=%d pending=%d; "
"selected=%d (limit=%d force=%s dry_run=%s delay=%.1fs city=%s)",
result.total,
result.fetchable,
result.pending,
len(rows),
limit,
force,
dry_run,
delay,
city,
)
if dry_run:
logger.info(
"dry_run: would process %d yandex_realty_nb houses (ids=%s)",
len(rows),
[r["house_id"] for r in rows],
)
result.duration_sec = time.time() - t0
return result
for idx, row in enumerate(rows):
house_id: int = int(row["house_id"])
jk_slug: str | None = row["yandex_jk_slug"]
ext_id: str | None = row["ext_id"]
result.processed += 1
# Idempotency fast-path (belt-and-braces под concurrent writer)
if not force and jk_slug:
already = db.execute(
text(
"SELECT 1 FROM market.yandex_jk_enrichment WHERE ext_id = CAST(:ext_id AS text)"
),
{"ext_id": ext_id},
).fetchone()
if already is not None:
result.skipped_already_enriched += 1
logger.debug("skip house_id=%s — already enriched ext_id=%s", house_id, ext_id)
continue
# ── Resolve slug ─────────────────────────────────────────────────
if not jk_slug:
if not ext_id:
logger.warning("house_id=%s: нет yandex_jk_slug и нет ext_id — skip", house_id)
result.failed_resolve += 1
continue
try:
resolved = await resolve_yandex_jk_slug(ext_id, city=city)
except Exception as exc:
logger.warning(
"resolve_yandex_jk_slug house_id=%s ext_id=%s raised: %s",
house_id,
ext_id,
exc,
)
resolved = None
# Anti-bot sleep после resolve-fetch
await _sleep_with_jitter(delay, idx, len(rows), force=True)
if not resolved:
logger.warning("slug unresolved house_id=%s ext_id=%s — skip", house_id, ext_id)
result.failed_resolve += 1
continue
# Persist slug под SAVEPOINT
sp = db.begin_nested()
try:
db.execute(text(_UPDATE_SLUG), {"slug": resolved, "hid": house_id})
sp.commit()
db.commit()
except Exception as exc:
sp.rollback()
logger.warning(
"persist slug failed house_id=%s ext_id=%s: %s", house_id, ext_id, exc
)
result.failed_resolve += 1
continue
jk_slug = resolved
result.resolved_slug += 1
# ── Guard: ext_id обязателен для fetch_jk (формирует URL /{slug}-{id}/) ──
if not ext_id:
logger.warning(
"house_id=%s: jk_slug=%s известен, но ext_id пустой — "
"fetch_jk пропущен (пустой jk_id даёт битый URL)",
house_id,
jk_slug,
)
result.failed_fetch += 1
continue
# ── Fetch через BrowserFetcher ────────────────────────────────────
info = None
try:
scraper = YandexNewbuildingScraper()
info = await scraper.fetch_jk(jk_slug=jk_slug, jk_id=ext_id, city=city)
except Exception as exc:
logger.warning(
"fetch_jk failed house_id=%s jk_slug=%s ext_id=%s: %s",
house_id,
jk_slug,
ext_id,
exc,
)
result.failed_fetch += 1
await _sleep_with_jitter(delay, idx, len(rows))
continue
if info is None:
logger.warning(
"fetch_jk returned None house_id=%s jk_slug=%s (anti-bot / parse miss?)",
house_id,
jk_slug,
)
result.failed_fetch += 1
await _sleep_with_jitter(delay, idx, len(rows))
continue
# ── UPSERT market.yandex_jk_enrichment под SAVEPOINT ─────────────
sp = db.begin_nested()
try:
db.execute(
text(_UPSERT_ENRICHMENT),
{
"ext_id": info.ext_id,
"name": info.name,
"developer_name": info.developer_name,
"address": info.address,
"lat": info.lat,
"lon": info.lon,
"rating": info.rating,
"ratings_count": info.ratings_count,
"text_reviews_count": info.text_reviews_count,
"raw_payload": json.dumps(info.raw_payload or {}, ensure_ascii=False),
},
)
sp.commit()
db.commit()
result.rows_inserted += 1
result.succeeded += 1
logger.info(
"enriched house_id=%s ext_id=%s slug=%s name=%r",
house_id,
info.ext_id,
jk_slug,
info.name,
)
except Exception as exc:
sp.rollback()
logger.warning(
"UPSERT yandex_jk_enrichment failed house_id=%s ext_id=%s: %s",
house_id,
ext_id,
exc,
)
try:
db.rollback()
except Exception as rb_exc:
logger.warning("rollback failed house_id=%s: %s", house_id, rb_exc)
await _sleep_with_jitter(delay, idx, len(rows))
result.duration_sec = time.time() - t0
logger.info(
"yandex-nb-sweep done: processed=%d ok=%d skip=%d resolved=%d "
"resolve_fail=%d fetch_fail=%d rows_inserted=%d | %.1fs",
result.processed,
result.succeeded,
result.skipped_already_enriched,
result.resolved_slug,
result.failed_resolve,
result.failed_fetch,
result.rows_inserted,
result.duration_sec,
)
return result
async def _sleep_with_jitter(delay: float, idx: int, total: int, *, force: bool = False) -> None:
"""Polite anti-bot sleep с ±20% jitter.
Пропускается после последнего элемента (нет следующего fetch) — ЕСЛИ force=False.
force=True используется для resolve→enrich gap внутри одного дома.
"""
if delay <= 0:
return
if not force and idx >= total - 1:
return
await asyncio.sleep(delay * random.uniform(0.8, 1.2))