"""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 (свой детект по
) или вообще без
статуса (сайдкар не дошёл до навигации), поэтому по статусу она диагностировалась
как «не разобрали»: рос только `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