gendesign/tradein-mvp/backend/app/tasks/cian_history_backfill.py
bot-backend ac5b044f7e
Some checks failed
CI Trade-In / changes (pull_request) Successful in 14s
CI / changes (pull_request) Successful in 18s
CI Trade-In / browser-tests (pull_request) Has been skipped
CI Trade-In / frontend-checks (pull_request) Has been skipped
CI / backend-tests (pull_request) Has been skipped
CI / frontend-tests (pull_request) Has been skipped
CI / openapi-codegen-check (pull_request) Has been skipped
CI Trade-In / backend-tests (pull_request) Failing after 6m17s
fix(scrapers): ожидаемые исходы сбора (бан, пустой пул, капча) больше не error
Скрапер один давал 3006 error-строк в сутки из ~3700 по всему Trade-In —
ленту перестали читать, и настоящая поломка терялась в ней.

Принцип: ожидаемый исход сбора (площадка забанила, пул прокси пуст,
капча/недогруз, серия блоков перевалила порог circuit breaker) — это
состояние работы против недружелюбного источника, а не инцидент. В error
остаётся только неожиданное: изменившаяся вёрстка/схема (Cian markup
change), просроченный токен ротации прокси (ASocks 401), неразобранное
исключение.

Переведено error -> warning в 9 файлах, 12 мест: "СТОП — пул прокси пуст"
(avito/domclick/cian_history/yandex_newbuilding_sweep x2), "пул прокси
исчерпан" (cian_session, cian_price_history, yandex_address_backfill),
FAIL-CLOSED без здорового узла для source (proxy_egress), ABORT по счётчику
подтверждённых блоков площадки (avito, domclick, yandex_detail_backfill).

Оставлено error намеренно: cookie-алерты Cian/DomClick (#2658, #2674) —
они рассчитаны именно на LoggingIntegration(event_level=ERROR) в
scheduler_main.py и без него молчат по 37 дней; ABORT по смешанным/soft
причинам без единого подтверждённого блока площадки (#3272, #2674/#3196) —
это может быть наш баг, а не бан, сигнал сознательно не приглушали.

GlitchTip: сентри-интеграция скрапера уже настроена как
LoggingIntegration(level=INFO, event_level=ERROR) в scheduler_main.py —
отдельной правки sentry_scrub.py не требуется, понижение уровня logger
само убирает эти записи из GlitchTip.

Итоговая FINISHED-строка со счётчиками (attempted/enriched/blocked/failed)
уже существует в каждом detail_backfill — новую не добавлял.

Refs #3471
2026-09-12 14:15:59 +03:00

503 lines
26 KiB
Python
Raw Permalink 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.

"""Batch backfill для Cian historical data.
Запускается:
- Manual через POST /admin/scrape/cian-backfill-history
- Future: in-app scheduler (после bootstrap'а Celery beat)
Что делает:
1. SELECT listings WHERE source='cian' AND source_url IS NOT NULL
AND NOT EXISTS (SELECT 1 FROM offer_price_history WHERE listing_id=...)
LIMIT batch_size
→ fetch_detail(browser_fetcher=bf) + save_detail_enrichment per listing
Использует BrowserFetcher (camoufox) для получения JS-rendered HTML, потому что
curl_cffi не рендерит priceChanges (они hidden в raw HTML).
2. Houses block: SELECT houses WHERE cian_zhk_url IS NOT NULL
AND NOT EXISTS (SELECT 1 FROM houses_price_dynamics WHERE house_id=...)
LIMIT batch_size
→ fetch_newbuilding(cian_zhk_url) + save_newbuilding_enrichment per house
Requires migration 071_houses_cian_zhk_url.sql (cian_zhk_url column).
Rate limit: scraper_settings.get_scraper_delay('cian') between requests.
Сигнал живости (#2725): батч дёргает `on_progress` на КАЖДОЙ сущности — caller
переливает это в scrape_runs.heartbeat_at. Пока колбэка не было, планировщик слал
heartbeat один раз ДО батча, а `reap_zombies` меряет ровно heartbeat с порогом 6 ч —
и добивал живые прогоны строго на 6-м часу (6 прод-прогонов, у пятерых внутри окна
писались строки offer_price_history, у одного — до 5.4 ч после старта). Ослаблять
критерий нельзя: пометка 'zombie' снимает running-блокировку источника
(`has_running_run`), без неё зависший прогон запер бы источник навсегда.
"""
from __future__ import annotations
import asyncio
import logging
import time
from collections import Counter
from collections.abc import Callable
from dataclasses import dataclass, field
from scraper_kit.browser_fetcher import BrowserFetcher, ban_kind_from_status
from scraper_kit.orchestration.pipeline import ban_kind_of_exception
from scraper_kit.orchestration.runs import BAN_KIND_UNKNOWN
from scraper_kit.providers.cian.detail import fetch_detail, save_detail_enrichment
from scraper_kit.providers.cian.valuation import estimate_via_cian_valuation
from scraper_kit.proxy_errors import caused_by_no_proxy
from sqlalchemy import text
from sqlalchemy.orm import Session
from app.core.config import settings
from app.services.scraper_adapters import (
RealMatcherAdapter,
RealProxyProvider,
RealScraperConfig,
)
from app.services.scraper_settings import get_scraper_delay
logger = logging.getLogger(__name__)
@dataclass
class CianBackfillResult:
listings_total: int = 0
listings_processed: int = 0
listings_succeeded: int = 0
listings_failed_fetch: int = 0
listings_failed_save: int = 0
price_changes_attempted: int = 0
houses_total: int = 0
houses_processed: int = 0
houses_succeeded: int = 0
houses_failed_fetch: int = 0
houses_failed_save: int = 0
valuations_total: int = 0
valuations_processed: int = 0
valuations_succeeded: int = 0
valuations_failed: int = 0
duration_sec: float = field(default=0.0)
# #3196: отказы detail-фетча и их перепись (kind -> сколько раз). До этой правки
# циановский прогон отдавал наверх только "не смогли обогатить", и scrape_runs.ban_kind
# у него не проставлялся вовсе.
listings_blocked: int = 0
ban_kinds: Counter[str] = field(default_factory=Counter)
# #3197: прогон оборван, потому что пул прокси пуст — к площадке не ходили вовсе.
# Не блок и не отказ Циана: caller (scheduler) обязан пометить прогон failed, а не
# banned, иначе запись утверждает про площадку то, чего не было.
no_proxy_stop: bool = False
@property
def ban_kind(self) -> str:
"""Доминирующий диагноз отказов прогона (#3196), пригоден для scrape_runs.ban_kind.
Считает тот же `_dominant_ban_kind`, что и остальные backfill-и: один вид →
он; строгое большинство → оно; иначе (и при пустой переписи) — 'unknown'.
Импорт локальный — таск не должен тянуть services на уровне модуля.
"""
from app.services.scrape_runs import _dominant_ban_kind
return _dominant_ban_kind(self.ban_kinds)
def _note_refusal(
result: CianBackfillResult, status: int | None, exc: BaseException | None = None
) -> str | None:
"""Записать отказ detail-фетча, если его природа установлена — типом или статусом.
Инвариант (#3196) прежний: непустой `ban_kinds` ⟺ отказ был ПОКАЗАН, а не назначен.
Установить его можно двумя способами, и статуса одного мало:
* ТИП исключения (#3402 follow-up) — `CianBlockedError` (капча Циана, отказ
сайдкара `SidecarBanPageError`, WAF-403) наследует `ProxyBanError`, то есть по
построению означает «площадка себя показала»: `ban_kind_of_exception` даёт
'platform'. Капча приезжает с HTTP 200 (свой детект по <title>) или вообще без
статуса (сайдкар не дошёл до навигации), поэтому по статусу она диагностировалась
как «не разобрали»: рос только `listings_failed_fetch`, `ban_kinds` оставался
пустым — и следующая капча-волна снова выглядела бы дрейфом нашей разметки;
* HTTP-статус ответа (403/429/5xx) — прежний путь для всего остального.
Всё, что не установлено ни тем, ни другим ('unknown' по типу И None по статусу), НЕ
инкрементит ни `listings_blocked`, ни `ban_kinds`: недиагностируемые случаи, записанные
как 'unknown', в scrape_runs.mark_banned (scheduler.py) превращали наши собственные
сбои в фиктивный бан площадки (#2764).
"""
kind = ban_kind_of_exception(exc) if exc is not None else BAN_KIND_UNKNOWN
if kind == BAN_KIND_UNKNOWN:
# Тип ничего не доказал — спрашиваем статус (прежнее поведение).
kind = ban_kind_from_status(status)
if kind is None:
return None
result.listings_blocked += 1
result.ban_kinds[kind] += 1
return kind
# ── Выборки «что ещё не добрано» ─────────────────────────────────────────────
# Две РАЗНЫЕ цели, которые до #3284 были склеены в одну. Историческая выборка
# ключуется по offer_price_history, поэтому объявление, у которого история уже
# есть, а карточки нет, не вернётся ей НИКОГДА (на 30.08 таких 1697). Вторая
# выборка закрывает ровно этот пробел и заодно даёт Циану то, что у avito /
# domclick / yandex есть давно, — добор по признаку «нет карточки».
_LISTINGS_PENDING_SQL: dict[str, str] = {
# Прежнее поведение, побайтово. Менять его правкой про добор карточек нельзя:
# по этой выборке живёт суточный cian_history_backfill.
"history": """
SELECT l.id, l.source_url
FROM listings l
LEFT JOIN offer_price_history oph ON oph.listing_id = l.id
WHERE l.source = 'cian'
AND l.source_url IS NOT NULL
AND oph.listing_id IS NULL
LIMIT :lim
""",
# #3284. ORDER BY здесь есть, а в "history" нет, и это намеренно: очередь
# карточек (на 30.08 — 20 728 объявлений) заведомо длиннее любого батча, и
# порядок решает, что мы успеем добрать. Свежие важнее: по ним считается
# оценка. У "history" очередь того же порядка, но её сортировку трогать —
# отдельное решение с отдельной проверкой, не побочный эффект этой правки.
"detail": """
SELECT l.id, l.source_url
FROM listings l
WHERE l.source = 'cian'
AND l.source_url IS NOT NULL
AND l.detail_enriched_at IS NULL
ORDER BY l.last_seen_at DESC NULLS LAST
LIMIT :lim
""",
}
async def backfill_cian_history(
db: Session,
*,
batch_size: int = 50,
do_listings: bool = True,
do_houses: bool = True,
do_valuations: bool = False,
dry_run: bool = False,
on_progress: Callable[[CianBackfillResult], None] | None = None,
listings_pending: str = "history",
) -> CianBackfillResult:
"""Iterate Cian listings + houses with missing history, fetch+save.
Args:
db: SQLAlchemy session (caller-owned; each entity commits internally via
save_detail_enrichment / save_newbuilding_enrichment).
batch_size: max entities to process per call (listings + houses + valuations counted
separately).
do_listings: process listings batch (offer_price_history backfill).
do_houses: process houses batch (houses_price_dynamics backfill via
fetch_newbuilding). Requires migration 071_houses_cian_zhk_url.sql applied.
do_valuations: process Cian Valuation Calculator batch (external_valuations backfill).
Default False — opt-in because each call hits Cian auth-gated API.
dry_run: skip all fetch+save; only count and log pending rows.
listings_pending: какую очередь разбирает блок listings (#3284) --
"history" (дефолт, прежнее поведение): нет строки в offer_price_history;
"detail": нет карточки (detail_enriched_at IS NULL), свежие первыми.
Неизвестное значение -- ValueError, а не молчаливый дефолт: пустой
батч из-за опечатки в default_params выглядел бы как «всё добрано».
on_progress: колбэк живости (#2725) — вызывается на каждой сущности ЛЮБОГО из
трёх блоков, до её обработки, с текущим (мутируемым) result. Caller пишет
heartbeat; исключения колбэка — на его совести (планировщик глушит их сам),
здесь они прервали бы батч.
Returns:
CianBackfillResult with per-domain counters + total wall-clock duration.
"""
result = CianBackfillResult()
t0 = time.time()
delay = get_scraper_delay("cian") # seconds; default 5.0
# ── 1. Listings: pending-выборка (см. _LISTINGS_PENDING_SQL) ──────────────
if do_listings:
try:
pending_sql = _LISTINGS_PENDING_SQL[listings_pending]
except KeyError:
raise ValueError(
f"listings_pending={listings_pending!r} неизвестен; "
f"допустимы {sorted(_LISTINGS_PENDING_SQL)}"
) from None
rows = db.execute(text(pending_sql), {"lim": batch_size}).mappings().all()
result.listings_total = len(rows)
if dry_run:
logger.info("dry_run: would process %d cian listings", result.listings_total)
else:
# One BrowserFetcher instance shared across all listings in this batch.
# priceChanges requires JS rendering — curl_cffi returns empty list (#1574).
# proxy_provider/use_pool/environment (#3197, шаг A2): без этих трёх сайдкар
# берёт свой env-прокси (SCRAPER_PROXY_URL) — прогон шёл мимо пула из 4 узлов
# целиком (ни выбора узла, ни scrape_proxy_source_bans, ни ротации), а
# прод-отказ «пул пуст → не ходить на env/direct» (#2616) на этом пути был
# мёртв: он смотрит на environment, который сюда не доезжал. Образец —
# domclick_detail_backfill.py:403 и house_imv_backfill.py:749 (#2698/#3197).
_cfg = RealScraperConfig()
async with BrowserFetcher(
source="cian",
endpoint=settings.browser_http_endpoint,
proxy_provider=RealProxyProvider(),
use_pool=_cfg.use_proxy_pool_browser,
environment=_cfg.environment,
) as bf:
for row in rows:
listing_id: int = row["id"]
source_url: str = row["source_url"]
result.listings_processed += 1
if on_progress is not None:
on_progress(result)
enrichment = None
try:
enrichment = await fetch_detail(source_url, browser_fetcher=bf)
except Exception as exc:
# #3197: пустой пул — не отказ площадки: запрос не уходил вовсе,
# следующее объявление упрётся ровно в то же самое (иначе батч
# крутит впустую весь список). Опора — тип в цепочке причин, а не
# текст: fetch_detail заворачивает сбой фетча в своё исключение.
if caused_by_no_proxy(exc):
result.no_proxy_stop = True
result.listings_failed_fetch += 1
logger.warning(
"cian_history_backfill: СТОП — пул прокси пуст, к площадке "
"не ходили. listing_id=%s processed=%d succeeded=%d",
listing_id,
result.listings_processed,
result.listings_succeeded,
)
break
# exc, а не только статус: капча Циана — подтверждённый отказ
# площадки (CianBlockedError), но приходит с HTTP 200/без статуса
# (#3402 follow-up). Без типа волна капчи писалась в счётчик
# «не разобрали» и прогон отдавал пустой ban_kinds.
kind = _note_refusal(result, bf.last_response_status, exc)
logger.warning(
"cian_detail fetch failed for listing_id=%s url=%s: %s "
"(http=%s ban_kind=%s)",
listing_id,
source_url,
exc,
bf.last_response_status,
kind,
)
result.listings_failed_fetch += 1
await asyncio.sleep(delay)
continue
if enrichment is None:
kind = _note_refusal(result, bf.last_response_status)
logger.warning(
"cian_detail fetch returned None for listing_id=%s url=%s "
"(http=%s ban_kind=%s)",
listing_id,
source_url,
bf.last_response_status,
kind,
)
result.listings_failed_fetch += 1
await asyncio.sleep(delay)
continue
try:
# matcher (#2435): резолвит канонический дом из bti_data (address/geo
# листинга) — иначе BTI-снапшот распарсен, но выбрасывается.
save_detail_enrichment(
db, listing_id, enrichment, matcher=RealMatcherAdapter()
)
result.listings_succeeded += 1
result.price_changes_attempted += len(enrichment.price_changes or [])
except Exception as exc:
logger.warning(
"cian_detail save failed for listing_id=%s: %s", listing_id, exc
)
result.listings_failed_save += 1
# Roll back to clean session state so next listing can proceed.
# save_detail_enrichment commits on success; on failure the
# transaction is left open/dirty — rollback to avoid session poison
# (same class of defect as the houses block below).
try:
db.rollback()
except Exception as rb_exc:
logger.warning(
"cian_detail rollback failed for listing_id=%s: %s",
listing_id,
rb_exc,
)
await asyncio.sleep(delay)
# ── 2. Houses: missing houses_price_dynamics ──────────────────────────────
# no_proxy_stop (#3197): пул пуст — дома идут через тот же пул (fetch_newbuilding с
# RealProxyProvider ниже), крутить их незачем.
if do_houses and not result.no_proxy_stop:
# Kit's fetch_newbuilding() now accepts config= (issue #2322 fixed) — pass
# RealScraperConfig() at the call site below so BrowserFetcher gets a real
# endpoint instead of degrading to endpoint=None (#2397 Part D2).
from scraper_kit.providers.cian.newbuilding import (
fetch_newbuilding,
save_newbuilding_enrichment,
)
house_rows = (
db.execute(
text("""
SELECT h.id, h.cian_zhk_url
FROM houses h
WHERE h.cian_zhk_url IS NOT NULL
AND NOT EXISTS (
SELECT 1 FROM houses_price_dynamics hpd
WHERE hpd.house_id = h.id
)
LIMIT :lim
"""),
{"lim": batch_size},
)
.mappings()
.all()
)
result.houses_total = len(house_rows)
if dry_run:
logger.info("dry_run: would process %d cian houses", result.houses_total)
else:
for hrow in house_rows:
house_id: int = hrow["id"]
zhk_url: str = hrow["cian_zhk_url"]
result.houses_processed += 1
if on_progress is not None:
on_progress(result)
enrichment = None
try:
# proxy_provider (#2767): тот же сожжённый env-узел бил и сюда —
# это второй вызывающий fetch_newbuilding, чинить надо оба.
enrichment = await fetch_newbuilding(
zhk_url,
config=RealScraperConfig(),
proxy_provider=RealProxyProvider(),
)
except Exception as exc:
logger.warning(
"cian_newbuilding fetch failed for house_id=%s url=%s: %s",
house_id,
zhk_url,
exc,
)
result.houses_failed_fetch += 1
await asyncio.sleep(delay)
continue
if enrichment is None:
logger.warning(
"cian_newbuilding fetch returned None for house_id=%s url=%s",
house_id,
zhk_url,
)
result.houses_failed_fetch += 1
await asyncio.sleep(delay)
continue
try:
save_newbuilding_enrichment(db, house_id, enrichment)
result.houses_succeeded += 1
except Exception as exc:
logger.warning(
"cian_newbuilding save failed for house_id=%s: %s", house_id, exc
)
result.houses_failed_save += 1
# Roll back to clean session state so next house can proceed.
# save_newbuilding_enrichment commits on success; on failure the
# transaction is left open/dirty — rollback to avoid session poison.
try:
db.rollback()
except Exception as rb_exc:
logger.warning(
"cian_newbuilding rollback failed for house_id=%s: %s",
house_id,
rb_exc,
)
await asyncio.sleep(delay)
# ── 3. Cian listings без external_valuations (price prediction backfill) ──
if do_valuations and not result.no_proxy_stop: # #3197: пул пуст — см. блок домов
rows = (
db.execute(
text("""
SELECT l.id, l.address, l.area_m2, l.rooms, l.floor, l.total_floors
FROM listings l
WHERE l.source = 'cian'
AND l.address IS NOT NULL
AND l.area_m2 IS NOT NULL
AND l.rooms IS NOT NULL
AND l.floor IS NOT NULL
AND l.total_floors IS NOT NULL
AND NOT EXISTS (
SELECT 1 FROM external_valuations ev
WHERE ev.source = 'cian_valuation'
AND ev.listing_id = l.id
AND ev.expires_at > NOW()
)
LIMIT :lim
"""),
{"lim": batch_size},
)
.mappings()
.all()
)
result.valuations_total = len(rows)
if dry_run:
logger.info(
"dry_run: would process %d cian listings for valuation", result.valuations_total
)
else:
for row in rows:
result.valuations_processed += 1
if on_progress is not None:
on_progress(result)
try:
cval = await estimate_via_cian_valuation(
db,
config=RealScraperConfig(),
address=row["address"],
total_area=float(row["area_m2"]),
rooms_count=int(row["rooms"]),
floor=int(row["floor"]),
total_floors=int(row["total_floors"]),
repair_type="cosmetic",
deal_type="sale",
use_cache=False, # force fresh fetch — populate cache
listing_id=int(row["id"]), # canonical link via mig 044
)
if cval is not None and cval.sale_price_rub:
result.valuations_succeeded += 1
else:
result.valuations_failed += 1
except Exception as exc:
logger.warning(
"cian_valuation backfill failed for listing_id=%s: %s", row["id"], exc
)
result.valuations_failed += 1
await asyncio.sleep(delay)
result.duration_sec = time.time() - t0
logger.info(
"cian_backfill done: listings=%d/%d (ok=%d fetch_fail=%d save_fail=%d "
"price_rows=%d) houses=%d/%d (ok=%d fail=%d) "
"valuations=%d/%d (ok=%d fail=%d) %.1fs",
result.listings_processed,
result.listings_total,
result.listings_succeeded,
result.listings_failed_fetch,
result.listings_failed_save,
result.price_changes_attempted,
result.houses_processed,
result.houses_total,
result.houses_succeeded,
result.houses_failed_fetch + result.houses_failed_save,
result.valuations_processed,
result.valuations_total,
result.valuations_succeeded,
result.valuations_failed,
result.duration_sec,
)
return result