"""scrape_pipeline.py — Avito full orchestrator (Stage 2e + city sweep).
Координирует 5 шагов одного pipeline run:
1. SEARCH — AvitoScraper().fetch_around(lat, lon, radius_m, pages=N)
2. SAVE — save_listings(db, lots) → (inserted, updated)
3. GROUP — dedupe house_url из lots → unique set
4. ENRICH_HOUSES — fetch_house_catalog + save_house_catalog_enrichment per unique house
5. ENRICH_DETAIL — top-N listings (detail_enriched_at IS NULL) → fetch_detail + save
City sweep: run_avito_city_sweep — итерирует EKB_ANCHORS × pages_per_anchor,
с cooperative cancel через scrape_runs.status и heartbeat update каждый anchor.
Anti-bot hardening (2023-05-23):
- Shared AsyncSession на весь sweep (один TLS handshake)
- asyncio.sleep с random +-20% jitter между detail requests
- AvitoBlockedError/AvitoRateLimitedError propagation → early abort
- После 3 consecutive blocks — abort sweep + mark_banned (status='banned')
Graceful degradation на каждом step: exception в одной house/listing
не валит весь pipeline (кроме blocked exceptions — они abort весь sweep).
IMV-фаза (финальная, после всех anchor'ов): если enrich_imv=True —
для домов, тронутых в этом sweep'е (touched house_id из houses-enrichment),
вызывается house_imv_backfill.process_houses_imv_batch. Ошибки per-house
логируются и считаются, но не валят sweep.
"""
from __future__ import annotations
import asyncio
import logging
import random
from contextlib import AsyncExitStack
from dataclasses import dataclass, field, fields
from datetime import date, timedelta
from urllib.parse import urlparse
from curl_cffi.requests import AsyncSession
from sqlalchemy import text
from sqlalchemy.orm import Session
from app.core.config import settings
from app.services import scrape_runs
from app.services.scrapers.avito import AvitoScraper
from app.services.scrapers.avito_detail import fetch_detail, save_detail_enrichment
from app.services.scrapers.avito_exceptions import AvitoBlockedError, AvitoRateLimitedError
from app.services.scrapers.avito_houses import fetch_house_catalog, save_house_catalog_enrichment
from app.services.scrapers.base import ScrapedLot, save_listings
from app.services.scrapers.browser_fetcher import BrowserFetcher
from app.services.scrapers.yandex_realty import YandexRealtyScraper
from app.services.yandex_price_history import record_yandex_price_history
logger = logging.getLogger(__name__)
async def _rotate_proxy_ip(*, reason: str, rotations_done: int) -> bool:
"""Вызвать changeip-ссылку mobileproxy и подождать settle (#1790).
Аналог AvitoScraper._rotate_ip, но на уровне pipeline — используется в
enrichment-фазах (houses + detail) где нет доступа к экземпляру scraper'а.
Returns True при успешной смене IP, False если rotate_url не задан или ошибка.
Логирует каждую ротацию (причина, порядковый номер, остаток лимита).
"""
rotate_url = settings.avito_proxy_rotate_url
max_rot = settings.avito_proxy_max_rotations
if not rotate_url:
return False
sep = "&" if "?" in rotate_url else "?"
try:
async with AsyncSession(timeout=30) as rot:
resp = await rot.get(f"{rotate_url}{sep}format=json")
new_ip: str = ""
try:
data = resp.json()
new_ip = str(data.get("new_ip", ""))
except Exception:
pass
await asyncio.sleep(settings.avito_proxy_rotate_settle_s)
logger.info(
"pipeline: IP rotated via changeip — reason=%s rotation=#%d remaining=%d new_ip=%s",
reason,
rotations_done + 1,
max_rot - rotations_done - 1,
new_ip or "unknown",
)
return True
except Exception:
logger.warning(
"pipeline: IP rotation failed (reason=%s, rotation=#%d)",
reason,
rotations_done + 1,
exc_info=True,
)
return False
# Per-anchor watchdog timeout (секунды). Если обработка одного anchor'а занимает
# дольше этого значения (зависший HTTP-запрос, прокси hang), asyncio.wait_for
# прерывает её с TimeoutError, sweep продолжается со следующим anchor'ом.
# Диапазон 180-300 с; 240 — разумный середина (4 мин > любой нормальный anchor).
ANCHOR_TIMEOUT_SEC: int = 240
# Константы для расчёта watchdog-таймаута combos-режима Yandex sweep.
# В combos-режиме единственный "anchor" делает num_combos × max_pages fetch'ей;
# каждый fetch занимает request_delay_sec + сетевой overhead. ADDRESS_ENRICH_BUDGET_S
# добавляет буфер на address-enrich фазу (per-listing HTTP + sleep × N листингов).
# Пример: 30 combos × 3 pages × (9+12) + 300 ≈ 2190s (~37 мин) — надёжно укладывается
# в окно 02:00-05:00 без риска перекрытия следующего sweep'а.
_YANDEX_COMBOS_PER_FETCH_S: float = 12.0 # network + parse margin (секунды на страницу)
_YANDEX_ADDRESS_ENRICH_BUDGET_S: float = 300.0 # бюджет на address-enrich фазу
# Константы для расчёта watchdog-таймаута Avito city sweep.
# Avito detail-fetch идёт через браузерный сервис с page.goto timeout 60s, и
# anti-bot страницы нередко доходят до этого лимита. Фиксированный ANCHOR_TIMEOUT_SEC=240
# guillotine'ит anchor ещё в detail-фазе → SERP-счётчики теряются, run показывает
# lots_fetched=0 при реально сохранённых листингах. Таймаут масштабируется:
# _avito_anchor_timeout = ANCHOR_TIMEOUT_SEC
# + detail_top_n × _AVITO_PER_DETAIL_S
# + (_AVITO_HOUSES_BUDGET_S if enrich_houses else 0)
# Пример: detail_top_n=20 + houses → 240 + 1000 + 180 = 1420s/anchor (worst-case).
# 5 anchor'ов × 1420s ≈ 2ч worst-case (нормальный прогон быстрее: не каждый detail
# достигает 60s timeout). Окно лишь ограничивает старт следующего sweep'а — scheduler
# ждёт финиша задачи и не запустит параллельный run.
_AVITO_PER_DETAIL_S: float = 50.0 # секунд на один detail-fetch (incl. occasional 60s timeout)
_AVITO_HOUSES_BUDGET_S: float = 180.0 # бюджет на house-enrich фазу
# Fail-fast: если detail phase накапливает подряд N timeout/ошибок, прерываем её
# (аналог consecutive_blocks для AvitoBlockedError в step 5). Срабатывает только при
# явной деградации browser-сервиса, не мешает нормальным transient-ошибкам.
_AVITO_DETAIL_CONSECUTIVE_TIMEOUT_ABORT: int = 5
# #1820: порог ПОСЛЕДОВАТЕЛЬНЫХ блоков (AvitoBlockedError/AvitoRateLimitedError) внутри
# enrichment-фаз run_avito_pipeline. Один заблокированный фетч (house/detail) больше НЕ
# роняет весь прогон — блок ловится пер-item (как обычная ошибка), и лишь N подряд
# прерывают enrichment-фазу. Listings из SEARCH+SAVE сохраняются всегда. Зеркалит
# паттерн avito_detail_backfill (max_consecutive_blocks).
_AVITO_PIPELINE_CONSECUTIVE_BLOCK_ABORT: int = 3
# Default anchors ЕКБ — 5 точек покрытия города
EKB_ANCHORS: list[tuple[float, float, str]] = [
(56.8400, 60.6050, "Центр"),
(56.7950, 60.5300, "ЮЗ"),
(56.8970, 60.6100, "Уралмаш"),
(56.7700, 60.5500, "Академический"),
(56.8650, 60.6200, "Пионерский"),
]
_CHROME_HEADERS = {
"Accept": "text/html,application/xhtml+xml,application/xml;q=0.9,*/*;q=0.8",
"Accept-Language": "ru-RU,ru;q=0.9,en;q=0.8",
"Cache-Control": "max-age=0",
"Sec-Fetch-Dest": "document",
"Sec-Fetch-Mode": "navigate",
"Sec-Fetch-Site": "none",
"Sec-Fetch-User": "?1",
"Upgrade-Insecure-Requests": "1",
}
def _avito_proxies() -> dict[str, str] | None:
"""#623/#806: mobile-proxy egress для shared-сессии sweep'а.
КРИТИЧНО: sweep создаёт собственную shared AsyncSession и присваивает её
`scraper._cffi`, перезатирая проксированную сессию из `AvitoScraper.__aenter__`.
Через эту же сессию идёт и SERP, и detail-энричмент (fetch_detail). Без прокси
detail-страницы летят с datacenter-IP → HTTP 429 → весь run mark_banned.
Пусто (env не задан) → прямое подключение (dev/no-op).
Читает settings.scraper_proxy_url (SCRAPER_PROXY_URL > AVITO_PROXY_URL fallback).
"""
url = settings.scraper_proxy_url
return {"http": url, "https": url} if url else None
# ── Result dataclasses ──────────────────────────────────────────
@dataclass
class PipelineCounters:
"""Counters одного pipeline run для логов/админки."""
lots_fetched: int = 0
lots_inserted: int = 0
lots_updated: int = 0
unique_houses: int = 0
houses_enriched: int = 0
houses_failed: int = 0
detail_attempted: int = 0
detail_enriched: int = 0
detail_failed: int = 0
errors: list[str] = field(default_factory=list)
@dataclass
class PipelineResult:
"""Полный результат full Avito pipeline."""
anchor_lat: float
anchor_lon: float
radius_m: int
counters: PipelineCounters
enrich_houses: bool
enrich_detail_top_n: int
touched_house_ids: set[int] = field(default_factory=set)
# ── Main orchestrator ───────────────────────────────────────────
async def run_avito_pipeline(
db: Session,
lat: float,
lon: float,
radius_m: int = 1500,
*,
enrich_houses: bool = True,
enrich_detail_top_n: int = 10,
pages: int = 1,
request_delay_sec: float | None = None,
shared_session: AsyncSession | None = None,
shared_browser: BrowserFetcher | None = None,
) -> PipelineResult:
"""Full Avito search → houses → detail enrichment pipeline.
Steps:
1. SEARCH: AvitoScraper().fetch_around(lat, lon, radius_m, pages=pages)
2. SAVE: save_listings(db, lots) → (inserted, updated)
3. GROUP: dedupe house_url из новых lots → unique house path set
4. ENRICH_HOUSES: fetch_house_catalog + save per house (with sleep between)
5. ENRICH_DETAIL: top-N listings → fetch_detail + save (abort on 3 blocks)
shared_session — если передан, используется без закрытия (lifecycle у вызывающего).
shared_browser — если передан в browser-mode, используется без закрытия.
Блокировки (AvitoBlockedError / AvitoRateLimitedError) в SEARCH-фазе (шаг 1)
по-прежнему propagate наружу (без листингов прогон бессмыслен). В enrichment-фазах
(houses шаг 4 + detail шаг 5 + house-detail шаг 5b) блок ловится ПЕР-ITEM как обычная
ошибка (#1820): один заблокированный фетч НЕ роняет весь прогон. Лишь
_AVITO_PIPELINE_CONSECUTIVE_BLOCK_ABORT блоков ПОДРЯД прерывают enrichment — но
PipelineResult всё равно возвращается с уже собранными listings/обогащениями.
"""
counters = PipelineCounters()
lots: list[ScrapedLot] = []
detail_delay = request_delay_sec if request_delay_sec is not None else 7.0
# #1820: общий счётчик ПОСЛЕДОВАТЕЛЬНЫХ блоков на все enrichment-фазы (4/5/5b).
# Сбрасывается на любом успешном fetch; при достижении порога — ротация IP (#1790)
# или abort если бюджет ротаций исчерпан.
enrich_consecutive_blocks = 0
enrich_rotations_done = 0 # #1790: сколько раз уже сменили IP в этом run
enrichment_aborted = False
browser_mode = settings.scraper_fetch_mode == "browser"
session: AsyncSession | None = None
browser_fetcher: BrowserFetcher | None = None
own_session = False
own_browser = False
scraper = AvitoScraper()
if browser_mode:
browser_fetcher = shared_browser
if browser_fetcher is None:
browser_fetcher = BrowserFetcher(source="avito")
await browser_fetcher.__aenter__()
own_browser = True
scraper._browser = browser_fetcher
else:
own_session = shared_session is None
session = shared_session or AsyncSession(
impersonate="chrome120",
timeout=25,
headers=_CHROME_HEADERS,
proxies=_avito_proxies(),
)
scraper._cffi = session
try:
# ── Step 1: search ──────────────────────────────────────
try:
lots = await scraper.fetch_around(
lat,
lon,
radius_m,
pages=pages,
delay_override_sec=request_delay_sec,
)
counters.lots_fetched = len(lots)
logger.info(
"pipeline:search anchor=(%.4f,%.4f) r=%dm fetched=%d",
lat,
lon,
radius_m,
len(lots),
)
except (AvitoBlockedError, AvitoRateLimitedError):
logger.error("pipeline:search BLOCKED by Avito at anchor=(%.4f,%.4f)", lat, lon)
raise
except Exception as e:
logger.exception("pipeline:search failed")
counters.errors.append(f"search: {e}")
return PipelineResult(lat, lon, radius_m, counters, enrich_houses, enrich_detail_top_n)
# ── Step 2: save listings ───────────────────────────────
if lots:
try:
counters.lots_inserted, counters.lots_updated = save_listings(db, lots)
except Exception as e:
logger.exception("pipeline:save_listings failed")
counters.errors.append(f"save_listings: {e}")
# ── Step 3: group by house ──────────────────────────────
unique_house_paths: set[str] = set()
for lot in lots:
if lot.house_url:
try:
parsed = urlparse(lot.house_url)
path = parsed.path if parsed.path else lot.house_url
unique_house_paths.add(path)
except Exception:
continue
counters.unique_houses = len(unique_house_paths)
logger.info("pipeline:group_by_house unique=%d", counters.unique_houses)
# ── Step 4: enrich houses с shared session + sleep ──────
touched_house_ids: set[int] = set()
if enrich_houses and unique_house_paths and not enrichment_aborted:
house_paths_list = list(unique_house_paths)
for idx, house_path in enumerate(house_paths_list):
try:
enrichment = await fetch_house_catalog(
house_path, cffi_session=session, browser_fetcher=browser_fetcher
)
house_counters = save_house_catalog_enrichment(db, enrichment)
counters.houses_enriched += 1
enrich_consecutive_blocks = 0
hid = house_counters.get("house_id")
if hid:
touched_house_ids.add(int(hid))
except (AvitoBlockedError, AvitoRateLimitedError) as e:
# #1820: блок ловим пер-item как обычную ошибку — НЕ роняем прогон.
enrich_consecutive_blocks += 1
counters.houses_failed += 1
counters.errors.append(f"BLOCKED_HOUSE {house_path}: {e}")
logger.warning(
"pipeline:house BLOCKED at %s (consecutive=%d)",
house_path,
enrich_consecutive_blocks,
)
if enrich_consecutive_blocks >= _AVITO_PIPELINE_CONSECUTIVE_BLOCK_ABORT:
# #1790: Перед abort'ом пробуем сменить IP (changeip).
# После ротации сбрасываем счётчик и продолжаем — НЕ abort.
# Если ротации исчерпаны — abort enrichment (listings уже сохранены).
if enrich_rotations_done < settings.avito_proxy_max_rotations:
rotated = await _rotate_proxy_ip(
reason=f"house consecutive={enrich_consecutive_blocks}",
rotations_done=enrich_rotations_done,
)
enrich_rotations_done += 1
if rotated:
enrich_consecutive_blocks = 0
logger.info(
"pipeline:enrich ROTATE — IP changed, "
"reset consecutive blocks (rotation #%d/%d)",
enrich_rotations_done,
settings.avito_proxy_max_rotations,
)
continue
logger.error(
"pipeline:enrich ABORT — %d consecutive blocks (IP rate-limited), "
"listings preserved",
enrich_consecutive_blocks,
)
counters.errors.append(
f"AVITO_BLOCKED — enrichment aborted after "
f"{enrich_consecutive_blocks} consecutive blocks (house phase)"
)
enrichment_aborted = True
break
except Exception as e:
logger.warning("pipeline:house_enrich failed for %s: %s", house_path, e)
counters.houses_failed += 1
counters.errors.append(f"house {house_path}: {e}")
# #1368: save_house_catalog_enrichment не делает внутренний rollback,
# поэтому упавший промежуточный statement оставляет общий Session в
# InFailedSqlTransaction. Без сброса — каскад «current transaction is
# aborted» на все оставшиеся дома/listings этого якоря.
try:
db.rollback()
except Exception:
pass
if idx < len(house_paths_list) - 1:
await asyncio.sleep(detail_delay)
# ── Step 5: enrich detail для top-N priority listings ───
if enrich_detail_top_n > 0 and not enrichment_aborted:
priority_rows = (
db.execute(
text("""
SELECT source_url
FROM listings
WHERE source = 'avito'
AND source_url IS NOT NULL
AND (
(
detail_enriched_at IS NULL
AND price_rub > 0
AND ST_DWithin(
geom::geography,
ST_MakePoint(:lon, :lat)::geography,
:radius
)
)
OR (
detail_enriched_at IS NULL
AND scraped_at > NOW() - INTERVAL '2 hours'
)
)
ORDER BY scraped_at DESC NULLS LAST
LIMIT :limit
"""),
{
"lat": lat,
"lon": lon,
"radius": radius_m * 2,
"limit": enrich_detail_top_n,
},
)
.mappings()
.all()
)
detail_house_paths: set[str] = set()
for idx, row in enumerate(priority_rows):
source_url: str = row["source_url"]
counters.detail_attempted += 1
try:
item_url = (
urlparse(source_url).path if source_url.startswith("http") else source_url
)
enrichment_detail = await fetch_detail(
item_url, cffi_session=session, browser_fetcher=browser_fetcher
)
if save_detail_enrichment(db, enrichment_detail):
counters.detail_enriched += 1
# #871: house-link есть в detail SSR (в SERP-карточке — JS-only).
# Копим catalog-пути для house-enrich ниже (Step 5b).
if enrichment_detail.house_catalog_url:
hp = urlparse(enrichment_detail.house_catalog_url).path
if hp:
detail_house_paths.add(hp)
enrich_consecutive_blocks = 0
except (AvitoBlockedError, AvitoRateLimitedError) as e:
# #1820: блок ловим пер-item, НЕ роняем прогон. Счётчик общий с house-фазой.
enrich_consecutive_blocks += 1
counters.detail_failed += 1
counters.errors.append(f"BLOCKED_DETAIL #{idx}: {e}")
logger.warning(
"pipeline:detail BLOCKED #%d/%d (consecutive=%d): %s",
idx + 1,
len(priority_rows),
enrich_consecutive_blocks,
e,
)
if enrich_consecutive_blocks >= _AVITO_PIPELINE_CONSECUTIVE_BLOCK_ABORT:
# #1790: Перед abort'ом пробуем сменить IP (changeip).
# После ротации сбрасываем счётчик — продолжаем detail-обход.
# Если ротации исчерпаны — abort detail-фазы (listings уже сохранены).
if enrich_rotations_done < settings.avito_proxy_max_rotations:
rotated = await _rotate_proxy_ip(
reason=f"detail consecutive={enrich_consecutive_blocks}",
rotations_done=enrich_rotations_done,
)
enrich_rotations_done += 1
if rotated:
enrich_consecutive_blocks = 0
logger.info(
"pipeline:enrich ROTATE — IP changed, "
"reset consecutive blocks (rotation #%d/%d)",
enrich_rotations_done,
settings.avito_proxy_max_rotations,
)
continue
logger.error(
"pipeline:detail ABORT — %d consecutive blocks (IP rate-limited), "
"listings preserved",
enrich_consecutive_blocks,
)
counters.errors.append(
f"AVITO_BLOCKED — enrichment aborted after "
f"{idx + 1}/{len(priority_rows)} details (consecutive blocks)"
)
# #1820: НЕ raise — прерываем только detail-фазу (и Step 5b),
# PipelineResult возвращается с уже сохранёнными listings/houses.
enrichment_aborted = True
break
except Exception as e:
counters.detail_failed += 1
counters.errors.append(f"detail {source_url}: {e}")
logger.warning("pipeline:detail_enrich failed for %s: %s", source_url, e)
# #1386: НЕ сбрасываем consecutive_blocks здесь. Транзиентная
# сетевая/parse-ошибка (5xx, timeout, connection reset) — правдоподобный
# побочный эффект активного IP-блока (Avito не отдаёт 403/429 единообразно).
# Сброс маскировал бы продолжающийся блок (block, block, timeout→reset, ...)
# и не давал бы сработать early-abort на 3 consecutive blocks.
# Сброс на успехе (см. выше) остаётся.
# #1368: save_detail_enrichment без внутреннего rollback мог оставить
# общий Session в aborted-state — сбрасываем, иначе каскад на оставшиеся
# detail-строки и Step 5b.
try:
db.rollback()
except Exception:
pass
if idx < len(priority_rows) - 1:
jitter = random.uniform(0.8, 1.2)
await asyncio.sleep(detail_delay * jitter)
# ── Step 5b: enrich houses, найденные через detail-страницы (#871) ──
# SERP-карточки больше не несут house-link (JS-render) → дома берём из
# detail SSR. Дедуп против обработанных в Step 4 (unique_house_paths).
new_house_paths = detail_house_paths - unique_house_paths
if enrich_houses and new_house_paths and not enrichment_aborted:
counters.unique_houses += len(new_house_paths)
nh_list = list(new_house_paths)
for h_idx, house_path in enumerate(nh_list):
try:
enrichment = await fetch_house_catalog(
house_path, cffi_session=session, browser_fetcher=browser_fetcher
)
house_counters = save_house_catalog_enrichment(db, enrichment)
counters.houses_enriched += 1
enrich_consecutive_blocks = 0
hid = house_counters.get("house_id")
if hid:
touched_house_ids.add(int(hid))
except (AvitoBlockedError, AvitoRateLimitedError) as e:
# #1820: блок пер-item, общий счётчик; не роняем прогон.
enrich_consecutive_blocks += 1
counters.houses_failed += 1
counters.errors.append(f"BLOCKED_HOUSE(detail) {house_path}: {e}")
logger.warning(
"pipeline:house(detail) BLOCKED at %s (consecutive=%d)",
house_path,
enrich_consecutive_blocks,
)
if enrich_consecutive_blocks >= _AVITO_PIPELINE_CONSECUTIVE_BLOCK_ABORT:
# #1790: попытка ротации перед abort'ом (общий бюджет run).
if enrich_rotations_done < settings.avito_proxy_max_rotations:
rotated = await _rotate_proxy_ip(
reason=f"house-detail consecutive={enrich_consecutive_blocks}",
rotations_done=enrich_rotations_done,
)
enrich_rotations_done += 1
if rotated:
enrich_consecutive_blocks = 0
logger.info(
"pipeline:enrich ROTATE — IP changed, "
"reset consecutive blocks (rotation #%d/%d)",
enrich_rotations_done,
settings.avito_proxy_max_rotations,
)
continue
logger.error(
"pipeline:enrich ABORT — %d consecutive blocks "
"(IP rate-limited), listings preserved",
enrich_consecutive_blocks,
)
counters.errors.append(
f"AVITO_BLOCKED — enrichment aborted after "
f"{enrich_consecutive_blocks} consecutive blocks "
f"(house-detail phase)"
)
enrichment_aborted = True
break
except Exception as e:
logger.warning(
"pipeline:house(detail) enrich failed for %s: %s", house_path, e
)
counters.houses_failed += 1
counters.errors.append(f"house(detail) {house_path}: {e}")
# #1368: см. Step 4 — сбрасываем aborted-state общего Session,
# иначе каскад на оставшиеся дома detail-фазы.
try:
db.rollback()
except Exception:
pass
if h_idx < len(nh_list) - 1:
await asyncio.sleep(detail_delay)
logger.info(
"pipeline:done anchor=(%.4f,%.4f) lots=%d (ins=%d/upd=%d) "
"houses=%d/%d detail=%d/%d touched_house_ids=%d errors=%d",
lat,
lon,
counters.lots_fetched,
counters.lots_inserted,
counters.lots_updated,
counters.houses_enriched,
counters.unique_houses,
counters.detail_enriched,
counters.detail_attempted,
len(touched_house_ids),
len(counters.errors),
)
return PipelineResult(
lat,
lon,
radius_m,
counters,
enrich_houses,
enrich_detail_top_n,
touched_house_ids=touched_house_ids,
)
finally:
if own_session and session is not None:
await session.close()
if own_browser and browser_fetcher is not None:
await browser_fetcher.__aexit__(None, None, None)
@dataclass
class CitySweepCounters:
"""Aggregate counters для full city sweep run."""
anchors_total: int = 0
anchors_done: int = 0
lots_fetched: int = 0
lots_inserted: int = 0
lots_updated: int = 0
unique_houses: int = 0
houses_enriched: int = 0
houses_failed: int = 0
detail_attempted: int = 0
detail_enriched: int = 0
detail_failed: int = 0
imv_attempted: int = 0
imv_enriched: int = 0
imv_failed: int = 0
errors_count: int = 0
def to_dict(self) -> dict[str, int]:
return {f.name: getattr(self, f.name) for f in fields(self)}
async def run_avito_city_sweep(
db: Session,
*,
run_id: int,
anchors: list[tuple[float, float, str]] | None = None,
radius_m: int = 1500,
pages_per_anchor: int = 3,
enrich_houses: bool = True,
detail_top_n: int = 10,
request_delay_sec: float = 7.0,
enrich_imv: bool = True,
) -> CitySweepCounters:
"""Full city sweep: iterate anchors × pages → save → enrich houses + detail → IMV.
- Single shared AsyncSession на весь sweep (один TLS fingerprint)
- AvitoBlockedError/AvitoRateLimitedError → ABORT sweep + mark_banned (status='banned')
- Прочие errors per-anchor логируются, не валят весь sweep
- После всех anchor'ов: если enrich_imv=True — IMV-оценка тронутых домов
(cooperative cancel + per-house graceful error handling)
SERP-счётчики (lots_fetched/inserted/updated) записываются в sweep-level counters
НЕМЕДЛЕННО после save_listings — до начала detail/houses фазы. Даже если anchor
превышает watchdog-таймаут в detail-фазе, SERP-результаты не теряются.
"""
from app.services.house_imv_backfill import process_houses_imv_batch
_anchors = anchors if anchors is not None else EKB_ANCHORS
counters = CitySweepCounters(anchors_total=len(_anchors))
all_touched_house_ids: set[int] = set()
# Масштабируемый watchdog-таймаут для Avito anchor'а. В отличие от Yandex/Cian,
# Avito detail-fetch идёт через browser service с page.goto timeout 60s. При
# detail_top_n=20 detail-фаза может занять до N×60s — фиксированный ANCHOR_TIMEOUT_SEC
# обрезает её раньше SERP, теряя счётчики. Масштабирование по detail_top_n + houses
# гарантирует, что SERP+save всегда успевает до watchdog.
_avito_anchor_timeout = (
float(ANCHOR_TIMEOUT_SEC)
+ detail_top_n * _AVITO_PER_DETAIL_S
+ (_AVITO_HOUSES_BUDGET_S if enrich_houses else 0.0)
)
logger.info(
"city-sweep run_id=%d: avito anchor_timeout=%.0fs "
"(base=%d + detail_top_n=%d×%.0fs + houses=%.0fs)",
run_id,
_avito_anchor_timeout,
ANCHOR_TIMEOUT_SEC,
detail_top_n,
_AVITO_PER_DETAIL_S,
_AVITO_HOUSES_BUDGET_S if enrich_houses else 0.0,
)
browser_mode = settings.scraper_fetch_mode == "browser"
# #1790: общий бюджет IP-ротаций на весь sweep (разделяется по всем anchor'ам).
sweep_rotations_done = 0
async with AsyncExitStack() as stack:
session: AsyncSession | None = None
shared_bf: BrowserFetcher | None = None
if browser_mode:
shared_bf = await stack.enter_async_context(BrowserFetcher(source="avito"))
else:
session = await stack.enter_async_context(
AsyncSession(
impersonate="chrome120",
timeout=25,
headers=_CHROME_HEADERS,
proxies=_avito_proxies(),
)
)
try:
for idx, (lat, lon, name) in enumerate(_anchors, start=1):
if scrape_runs.is_cancelled(db, run_id):
logger.info(
"city-sweep run_id=%d: cancelled at anchor #%d/%d (%s)",
run_id,
idx,
len(_anchors),
name,
)
scrape_runs.update_heartbeat(db, run_id, counters.to_dict())
return counters
logger.info(
"city-sweep run_id=%d anchor #%d/%d (%s, %.4f, %.4f)",
run_id,
idx,
len(_anchors),
name,
lat,
lon,
)
# Capture loop variables in default args (B023): prevents stale binding
# if the coroutine is scheduled after the loop variable changes.
_a_lat, _a_lon, _a_name = lat, lon, name
async def _avito_anchor_phases(
_lat: float = _a_lat,
_lon: float = _a_lon,
_name: str = _a_name,
) -> None:
"""Все фазы одного Avito anchor'а (SERP+save → houses → detail).
Обновляет sweep-level counters и all_touched_house_ids через nonlocal
НЕМЕДЛЕННО после SERP+save (до detail-фазы). Это гарантирует, что
TimeoutError из wait_for теряет только detail-счётчики, но не SERP.
"""
nonlocal all_touched_house_ids, sweep_rotations_done
# ── Phase 1+2: SERP + save ─────────────────────────────────
scraper = AvitoScraper()
if browser_mode:
scraper._browser = shared_bf
else:
scraper._cffi = session
try:
anchor_lots: list[ScrapedLot] = await scraper.fetch_around(
_lat,
_lon,
radius_m,
pages=pages_per_anchor,
delay_override_sec=request_delay_sec,
)
except (AvitoBlockedError, AvitoRateLimitedError):
logger.error(
"city-sweep run_id=%d: SERP BLOCKED at anchor (%s, %.4f, %.4f)",
run_id,
_name,
_lat,
_lon,
)
raise
counters.lots_fetched += len(anchor_lots)
if anchor_lots:
try:
ins, upd = save_listings(db, anchor_lots)
counters.lots_inserted += ins
counters.lots_updated += upd
except Exception as save_exc:
logger.exception(
"city-sweep run_id=%d: save_listings failed at anchor %s: %s",
run_id,
_name,
save_exc,
)
counters.errors_count += 1
logger.info(
"city-sweep run_id=%d anchor %s: SERP fetched=%d ins=%d upd=%d",
run_id,
_name,
len(anchor_lots),
counters.lots_inserted,
counters.lots_updated,
)
# SERP-счётчики зафиксированы в sweep-level counters.
# Дальнейший TimeoutError из wait_for их не сотрёт.
# ── Phase 3: group by house ────────────────────────────────
unique_house_paths: set[str] = set()
for lot in anchor_lots:
if lot.house_url:
try:
parsed = urlparse(lot.house_url)
path = parsed.path if parsed.path else lot.house_url
unique_house_paths.add(path)
except Exception:
continue
if unique_house_paths:
counters.unique_houses += len(unique_house_paths)
# ── Phase 4: enrich houses ─────────────────────────────────
touched_house_ids: set[int] = set()
sweep_consecutive_house_blocks = 0 # #1790: per-anchor счётчик
if enrich_houses and unique_house_paths:
house_paths_list = list(unique_house_paths)
h_idx = 0
while h_idx < len(house_paths_list):
house_path = house_paths_list[h_idx]
try:
enrichment = await fetch_house_catalog(
house_path,
cffi_session=session,
browser_fetcher=shared_bf,
)
hc = save_house_catalog_enrichment(db, enrichment)
counters.houses_enriched += 1
sweep_consecutive_house_blocks = 0
hid = hc.get("house_id")
if hid:
touched_house_ids.add(int(hid))
if h_idx < len(house_paths_list) - 1:
await asyncio.sleep(request_delay_sec)
h_idx += 1
except (AvitoBlockedError, AvitoRateLimitedError) as be:
# #1790: пробуем сменить IP перед re-raise.
sweep_consecutive_house_blocks += 1
counters.houses_failed += 1
counters.errors_count += 1
logger.warning(
"city-sweep run_id=%d: house BLOCKED at %s "
"(consecutive=%d) — %s",
run_id,
house_path,
sweep_consecutive_house_blocks,
be,
)
if (
sweep_consecutive_house_blocks
>= _AVITO_PIPELINE_CONSECUTIVE_BLOCK_ABORT
and sweep_rotations_done < settings.avito_proxy_max_rotations
):
rotated = await _rotate_proxy_ip(
reason=(
f"city-sweep house "
f"consecutive={sweep_consecutive_house_blocks}"
),
rotations_done=sweep_rotations_done,
)
sweep_rotations_done += 1
if rotated:
sweep_consecutive_house_blocks = 0
logger.info(
"city-sweep run_id=%d: house ROTATE — "
"IP changed, reset consecutive "
"(rotation #%d/%d)",
run_id,
sweep_rotations_done,
settings.avito_proxy_max_rotations,
)
# Не инкрементируем h_idx — retry текущего дома
continue
elif sweep_consecutive_house_blocks >= 3:
logger.error(
"city-sweep run_id=%d: house ABORT — "
"%d consecutive blocks, propagating",
run_id,
sweep_consecutive_house_blocks,
)
raise
h_idx += 1
except Exception as he:
logger.warning(
"city-sweep run_id=%d: house_enrich failed %s: %s",
run_id,
house_path,
he,
)
counters.houses_failed += 1
counters.errors_count += 1
try:
db.rollback()
except Exception:
pass
h_idx += 1
# ── Phase 5: enrich detail (top-N) ────────────────────────
if detail_top_n > 0 and anchor_lots:
priority_rows = (
db.execute(
text("""
SELECT source_url
FROM listings
WHERE source = 'avito'
AND source_url IS NOT NULL
AND (
(
detail_enriched_at IS NULL
AND price_rub > 0
AND ST_DWithin(
geom::geography,
ST_MakePoint(:lon, :lat)::geography,
:radius
)
)
OR (
detail_enriched_at IS NULL
AND scraped_at > NOW() - INTERVAL '2 hours'
)
)
ORDER BY scraped_at DESC NULLS LAST
LIMIT :limit
"""),
{
"lat": _lat,
"lon": _lon,
"radius": radius_m * 2,
"limit": detail_top_n,
},
)
.mappings()
.all()
)
consecutive_blocks = 0
consecutive_timeouts = 0
detail_house_paths: set[str] = set()
for d_idx, row in enumerate(priority_rows):
source_url: str = row["source_url"]
counters.detail_attempted += 1
try:
item_url = (
urlparse(source_url).path
if source_url.startswith("http")
else source_url
)
enrichment_detail = await fetch_detail(
item_url,
cffi_session=session,
browser_fetcher=shared_bf,
)
if save_detail_enrichment(db, enrichment_detail):
counters.detail_enriched += 1
if enrichment_detail.house_catalog_url:
hp = urlparse(enrichment_detail.house_catalog_url).path
if hp:
detail_house_paths.add(hp)
consecutive_blocks = 0
consecutive_timeouts = 0
except (AvitoBlockedError, AvitoRateLimitedError) as be:
consecutive_blocks += 1
counters.detail_failed += 1
counters.errors_count += 1
logger.warning(
"city-sweep run_id=%d: detail BLOCKED #%d/%d "
"(consecutive=%d): %s",
run_id,
d_idx + 1,
len(priority_rows),
consecutive_blocks,
be,
)
if consecutive_blocks >= 3:
# #1790: попытка ротации перед abort'ом.
if sweep_rotations_done < settings.avito_proxy_max_rotations:
rotated = await _rotate_proxy_ip(
reason=(
f"city-sweep detail "
f"consecutive={consecutive_blocks}"
),
rotations_done=sweep_rotations_done,
)
sweep_rotations_done += 1
if rotated:
consecutive_blocks = 0
logger.info(
"city-sweep run_id=%d: detail ROTATE — "
"IP changed, reset consecutive "
"(rotation #%d/%d)",
run_id,
sweep_rotations_done,
settings.avito_proxy_max_rotations,
)
continue
logger.error(
"city-sweep run_id=%d: detail ABORT — "
"%d consecutive blocks (IP rate-limited)",
run_id,
consecutive_blocks,
)
raise
except Exception as de:
counters.detail_failed += 1
counters.errors_count += 1
consecutive_timeouts += 1
logger.warning(
"city-sweep run_id=%d: detail failed #%d/%d "
"(consecutive_timeouts=%d): %s",
run_id,
d_idx + 1,
len(priority_rows),
consecutive_timeouts,
de,
)
# (C) Fail-fast: браузер явно деградирует — прерываем
# detail-фазу, не тратим время на оставшиеся запросы.
if consecutive_timeouts >= _AVITO_DETAIL_CONSECUTIVE_TIMEOUT_ABORT:
logger.error(
"city-sweep run_id=%d: detail ABORT — "
"%d consecutive timeouts/errors, browser degraded",
run_id,
consecutive_timeouts,
)
break
try:
db.rollback()
except Exception:
pass
if d_idx < len(priority_rows) - 1:
jitter = random.uniform(0.8, 1.2)
await asyncio.sleep(request_delay_sec * jitter)
# ── Phase 5b: houses найденные через detail-страницы ──
new_house_paths = detail_house_paths - unique_house_paths
if enrich_houses and new_house_paths:
counters.unique_houses += len(new_house_paths)
nh_list = list(new_house_paths)
nh_idx = 0
sweep_consecutive_detail_house_blocks = 0 # #1790
while nh_idx < len(nh_list):
house_path = nh_list[nh_idx]
try:
enrichment = await fetch_house_catalog(
house_path,
cffi_session=session,
browser_fetcher=shared_bf,
)
hc = save_house_catalog_enrichment(db, enrichment)
counters.houses_enriched += 1
sweep_consecutive_detail_house_blocks = 0
hid = hc.get("house_id")
if hid:
touched_house_ids.add(int(hid))
if nh_idx < len(nh_list) - 1:
await asyncio.sleep(request_delay_sec)
nh_idx += 1
except (AvitoBlockedError, AvitoRateLimitedError) as beh:
# #1790: ротация перед re-raise.
sweep_consecutive_detail_house_blocks += 1
counters.houses_failed += 1
counters.errors_count += 1
logger.warning(
"city-sweep run_id=%d: house(detail) BLOCKED %s"
" (consecutive=%d) — %s",
run_id,
house_path,
sweep_consecutive_detail_house_blocks,
beh,
)
if (
sweep_consecutive_detail_house_blocks >= 3
and sweep_rotations_done
< settings.avito_proxy_max_rotations
):
rotated = await _rotate_proxy_ip(
reason=(
"city-sweep house(detail) "
f"consecutive="
f"{sweep_consecutive_detail_house_blocks}"
),
rotations_done=sweep_rotations_done,
)
sweep_rotations_done += 1
if rotated:
sweep_consecutive_detail_house_blocks = 0
logger.info(
"city-sweep run_id=%d: house(detail) "
"ROTATE — IP changed (rotation #%d/%d)",
run_id,
sweep_rotations_done,
settings.avito_proxy_max_rotations,
)
continue
elif sweep_consecutive_detail_house_blocks >= 3:
logger.error(
"city-sweep run_id=%d: house(detail) ABORT"
" — %d consecutive blocks, propagating",
run_id,
sweep_consecutive_detail_house_blocks,
)
raise
nh_idx += 1
except Exception as he:
logger.warning(
"city-sweep run_id=%d: house(detail) failed %s: %s",
run_id,
house_path,
he,
)
counters.houses_failed += 1
counters.errors_count += 1
try:
db.rollback()
except Exception:
pass
nh_idx += 1
all_touched_house_ids.update(touched_house_ids)
logger.info(
"city-sweep run_id=%d anchor %s done: "
"lots=%d houses=%d/%d detail=%d/%d touched=%d errors=%d",
run_id,
_name,
len(anchor_lots),
counters.houses_enriched,
counters.unique_houses,
counters.detail_enriched,
counters.detail_attempted,
len(touched_house_ids),
counters.errors_count,
)
try:
await asyncio.wait_for(_avito_anchor_phases(), timeout=_avito_anchor_timeout)
except TimeoutError:
logger.warning(
"city-sweep run_id=%d: anchor #%d/%d (%s, %.4f, %.4f) "
"timed out after %.0fs — SERP counters preserved, "
"detail phase incomplete",
run_id,
idx,
len(_anchors),
name,
lat,
lon,
_avito_anchor_timeout,
)
counters.errors_count += 1
except (AvitoBlockedError, AvitoRateLimitedError) as e:
logger.error(
"city-sweep run_id=%d ABORT at anchor #%d/%d (%s) — blocked: %s",
run_id,
idx,
len(_anchors),
name,
e,
)
counters.errors_count += 1
counters.anchors_done = idx
scrape_runs.update_heartbeat(db, run_id, counters.to_dict())
scrape_runs.mark_banned(db, run_id, str(e), counters.to_dict())
return counters
except Exception:
logger.exception("city-sweep run_id=%d: anchor %s failed", run_id, name)
counters.errors_count += 1
counters.anchors_done = idx
scrape_runs.update_heartbeat(db, run_id, counters.to_dict())
# ── IMV-фаза: финальный обход тронутых домов ──────────
if enrich_imv and all_touched_house_ids:
if scrape_runs.is_cancelled(db, run_id):
logger.info(
"city-sweep run_id=%d: cancelled before IMV phase (%d houses skipped)",
run_id,
len(all_touched_house_ids),
)
scrape_runs.mark_done(db, run_id, counters.to_dict())
return counters
logger.info(
"city-sweep run_id=%d: IMV phase — %d touched houses",
run_id,
len(all_touched_house_ids),
)
try:
# Периодический heartbeat внутри длинной IMV-фазы (#1363):
# без него десятки/сотни домов × request_delay_sec проходят
# без update_heartbeat, и reap_zombies помечает живой run
# 'zombie' → последующий mark_done становится no-op (дубль-sweep).
def _imv_heartbeat() -> None:
scrape_runs.update_heartbeat(db, run_id, counters.to_dict())
imv_result = await process_houses_imv_batch(
db,
all_touched_house_ids,
request_delay_sec=request_delay_sec,
heartbeat=_imv_heartbeat,
)
counters.imv_attempted += imv_result.checked
counters.imv_enriched += imv_result.saved
counters.imv_failed += imv_result.errors
counters.errors_count += imv_result.errors
scrape_runs.update_heartbeat(db, run_id, counters.to_dict())
logger.info(
"city-sweep run_id=%d: IMV phase done — attempted=%d enriched=%d failed=%d",
run_id,
imv_result.checked,
imv_result.saved,
imv_result.errors,
)
except Exception as exc:
logger.error(
"city-sweep run_id=%d: IMV phase crashed — %r (sweep still done)",
run_id,
exc,
)
counters.errors_count += 1
# #1521: IMV-фаза могла оставить Session в aborted-state
# (необёрнутый eligibility SELECT / _mark_status на отравленной
# сессии). update_heartbeat и mark_done — единственные финализаторы
# БЕЗ defensive rollback, поэтому сбрасываем txn-state сами. Иначе
# PendingRollbackError пролетает мимо внутреннего try во внешний
# except (mark_failed) → успешный sweep ложно помечается 'failed'.
try:
db.rollback()
except Exception:
pass
scrape_runs.update_heartbeat(db, run_id, counters.to_dict())
scrape_runs.mark_done(db, run_id, counters.to_dict())
logger.info(
"city-sweep run_id=%d done: anchors=%d/%d lots=%d (ins=%d/upd=%d) "
"houses=%d/%d detail=%d/%d imv=%d/%d errors=%d",
run_id,
counters.anchors_done,
counters.anchors_total,
counters.lots_fetched,
counters.lots_inserted,
counters.lots_updated,
counters.houses_enriched,
counters.unique_houses,
counters.detail_enriched,
counters.detail_attempted,
counters.imv_enriched,
counters.imv_attempted,
counters.errors_count,
)
return counters
except Exception as exc:
logger.exception("city-sweep run_id=%d: fatal error", run_id)
scrape_runs.mark_failed(db, run_id, str(exc), counters.to_dict())
raise
# ── Avito newbuilding (novostroyka) citywide sweep ────────────────
@dataclass
class NewbuildingSweepCounters:
"""Aggregate counters для Avito novostroyka citywide sweep run."""
lots_fetched: int = 0
lots_inserted: int = 0
lots_updated: int = 0
errors_count: int = 0
def to_dict(self) -> dict[str, int]:
return {f.name: getattr(self, f.name) for f in fields(self)}
async def run_avito_newbuilding_sweep(
db: Session,
*,
run_id: int,
pages: int = 20,
request_delay_sec: float = 7.0,
) -> NewbuildingSweepCounters:
"""Citywide-обход ЕКБ-выборки только новостроек (novostroyka-filter) → save.
В отличие от run_avito_city_sweep (anchor-based fetch_around + houses/detail/IMV),
это простой citywide paginated-обход через AvitoScraper.fetch_newbuildings:
dedicated novostroyka-SERP отдаёт 100% new-build карточки (listing_segment=
"novostroyki", newbuilding_id/newbuilding_url заполнены парсером). Сохраняем через
стандартный save_listings, ведём SERP-счётчики через scrape_runs.
- Single shared AsyncSession/BrowserFetcher на весь sweep (один TLS fingerprint)
- AvitoBlockedError/AvitoRateLimitedError → mark_banned (status='banned')
- is_cancelled-чек до старта SERP-фазы
"""
counters = NewbuildingSweepCounters()
browser_mode = settings.scraper_fetch_mode == "browser"
async with AsyncExitStack() as stack:
session: AsyncSession | None = None
shared_bf: BrowserFetcher | None = None
if browser_mode:
shared_bf = await stack.enter_async_context(BrowserFetcher(source="avito"))
else:
session = await stack.enter_async_context(
AsyncSession(
impersonate="chrome120",
timeout=25,
headers=_CHROME_HEADERS,
proxies=_avito_proxies(),
)
)
try:
if scrape_runs.is_cancelled(db, run_id):
logger.info("nb-sweep run_id=%d: cancelled before SERP phase", run_id)
scrape_runs.mark_done(db, run_id, counters.to_dict())
return counters
scraper = AvitoScraper()
if browser_mode:
scraper._browser = shared_bf
else:
scraper._cffi = session
try:
lots: list[ScrapedLot] = await scraper.fetch_newbuildings(
pages=pages,
delay_override_sec=request_delay_sec,
)
except (AvitoBlockedError, AvitoRateLimitedError) as e:
logger.error("nb-sweep run_id=%d: SERP BLOCKED — %s", run_id, e)
scrape_runs.mark_banned(db, run_id, str(e), counters.to_dict())
return counters
counters.lots_fetched += len(lots)
if lots:
try:
ins, upd = save_listings(db, lots)
counters.lots_inserted += ins
counters.lots_updated += upd
except Exception as save_exc:
logger.exception(
"nb-sweep run_id=%d: save_listings failed: %s", run_id, save_exc
)
counters.errors_count += 1
try:
db.rollback()
except Exception:
pass
scrape_runs.update_heartbeat(db, run_id, counters.to_dict())
scrape_runs.mark_done(db, run_id, counters.to_dict())
logger.info(
"nb-sweep run_id=%d done: lots=%d (ins=%d/upd=%d) errors=%d",
run_id,
counters.lots_fetched,
counters.lots_inserted,
counters.lots_updated,
counters.errors_count,
)
return counters
except Exception as exc:
logger.exception("nb-sweep run_id=%d: fatal error", run_id)
scrape_runs.mark_failed(db, run_id, str(exc), counters.to_dict())
raise
# ── Yandex city sweep ───────────────────────────────────────────
@dataclass
class YandexCitySweepCounters:
"""Aggregate counters для Yandex.Недвижимость city sweep run.
Фазы: SERP (fetch_around_multi_room) → save_listings инкрементально → address-enrich.
address_* — результат T10-обогащения (detail
→ полный адрес с номером дома).
combos_skipped — сколько (room×price) combo пропущено из-за таймаутов gate-API.
"""
anchors_total: int = 0
anchors_done: int = 0
lots_fetched: int = 0
lots_inserted: int = 0
lots_updated: int = 0
price_history_rows: int = 0
combos_skipped: int = 0
address_attempted: int = 0
address_enriched: int = 0
address_failed: int = 0
errors_count: int = 0
def to_dict(self) -> dict[str, int]:
return {f.name: getattr(self, f.name) for f in fields(self)}
# Defensive abort threshold: у YandexRealtyScraper нет block/rate-exception классов
# (fetch_around глотает ошибки и возвращает []). Чтобы не молотить Yandex впустую при
# системном блоке/бане IP, прерываем sweep после N подряд anchor'ов, упавших с
# исключением. NB: пустой результат (0 lots) — это НЕ ошибка, счётчик не растёт.
YANDEX_SWEEP_MAX_CONSECUTIVE_FAILURES = 3
async def run_yandex_city_sweep(
db: Session,
*,
run_id: int,
anchors: list[tuple[float, float, str]] | None = None,
radius_m: int = 25000,
pages_per_anchor: int = 2,
request_delay_sec: float | None = None,
enrich_address: bool = True,
rooms_list: list[str] | None = None,
price_ranges: list[tuple[int | None, int | None]] | None = None,
segments: list[str] | None = None,
) -> YandexCitySweepCounters:
"""Yandex.Недвижимость city sweep: rooms × price combos от центра ЕКБ → save → address-enrich.
Вместо 5 географических anchor'ов итерирует ВСЕ combinations (room × price_range)
из одной центральной точки ЕКБ (56.8400, 60.6050). radius_m=25000 охватывает весь
город; price/room segmentation обходит Yandex SERP ~575-card cap без geo-шрапнели.
Фазы (единственный центральный якорь):
1. SERP: YandexRealtyScraper.fetch_around_multi_room(lat, lon, radius_m,
max_pages=pages_per_anchor, rooms_list=..., price_ranges=...) — combos mode.
2. SAVE: save_listings(db, lots, run_id=run_id).
3. ADDRESS-ENRICH (если enrich_address=True): для yandex-листингов без номера
дома в address (address IS NOT NULL AND NOT (address ~ ',\\s*\\d+'))
запрашиваем detail-страницу и извлекаем полный адрес из .
Параметр `anchors` сохранён для backward-compat тестов (anchors=[...]). В prod
(anchors=None) всегда используется один EKB_CENTER anchor.
Anti-bot / cancel семантика:
- Cooperative cancel: is_cancelled(db, run_id) проверяется перед фазой SERP.
- YandexRealtyScraper не propagate'ит block/rate exceptions (глотает → []).
Ловим broad Exception; при N подряд — mark_banned + abort.
- request_delay_sec: inter-request delay для address-enrich (поверх per-page
sleep внутри scraper'а).
Возвращает YandexCitySweepCounters.
mark_done вызывается ВСЕГДА (кроме cancel / ban / fatal).
"""
from sqlalchemy import text as _text
from app.services.scrapers.yandex_realty import DEFAULT_PRICE_RANGES, ROOM_PATH
from app.services.yandex_address_backfill import (
_RE_HAS_HOUSE_NUMBER,
_extract_address_from_title,
)
# EKB center — единственный anchor для citywide combos sweep.
# anchors-параметр принимается для backward-compat тестов.
if anchors is not None:
_anchors: list[tuple[float, float, str]] = anchors
else:
_anchors = [(56.8400, 60.6050, "Центр (combos)")]
_rooms_list = rooms_list or list(ROOM_PATH.keys())
_price_ranges = price_ranges or DEFAULT_PRICE_RANGES
# newFlat-сегменты: оба прохода (vtorichka + novostroyki) по одним combos.
# Дефолт ["NO", "YES"] — оба; dedup по source_id span'ит оба прохода (один seen
# внутри fetch_around_multi_room). Counters (lots_fetched/inserted/updated)
# агрегируют оба прохода через len(anchor_lots).
_segments = segments or ["NO", "YES"]
counters = YandexCitySweepCounters(anchors_total=len(_anchors))
inter_anchor_delay = request_delay_sec if request_delay_sec is not None else 7.0
enrich_delay = request_delay_sec if request_delay_sec is not None else 3.0
_resolved_delay = request_delay_sec if request_delay_sec is not None else 9.0
consecutive_failures = 0
# Вычисляем watchdog-таймаут для combos-режима (центр, anchors=None).
# В combos-режиме один "anchor" выполняет num_segments × num_combos × max_pages
# fetch'ей — каждый занимает примерно _resolved_delay + _YANDEX_COMBOS_PER_FETCH_S сек.
# Для 2 segments × 30 combos × 3 pages × (9+12) + 300 ≈ 4080s (~68 мин).
# Множитель len(_segments) масштабирует watchdog пропорционально числу проходов
# (vtorichka + novostroyki ≈ ×2 wall-clock) — иначе второй проход гильотинится.
# В explicit-anchor (тест/override) режиме оставляем ANCHOR_TIMEOUT_SEC (240s) —
# там каждый anchor небольшой и watchdog работает как старый защитный барьер.
_num_combos = len(_rooms_list) * len(_price_ranges)
if anchors is None and _num_combos > 0:
_sweep_timeout = max(
ANCHOR_TIMEOUT_SEC,
len(_segments)
* _num_combos
* pages_per_anchor
* (_resolved_delay + _YANDEX_COMBOS_PER_FETCH_S)
+ _YANDEX_ADDRESS_ENRICH_BUDGET_S,
)
else:
# Explicit anchors (backward-compat тесты или ручной override) — legacy timeout.
_sweep_timeout = ANCHOR_TIMEOUT_SEC
# Прокси для address-enrich curl_cffi сессии (зеркало yandex_address_backfill).
_proxy_url = settings.scraper_proxy_url
_proxies = {"http": _proxy_url, "https": _proxy_url} if _proxy_url else None
try:
for idx, (lat, lon, name) in enumerate(_anchors, start=1):
# ── Cooperative cancel ───────────────────────────────────────────
if scrape_runs.is_cancelled(db, run_id):
logger.info(
"yandex-sweep run_id=%d: cancelled at anchor #%d/%d (%s)",
run_id,
idx,
len(_anchors),
name,
)
scrape_runs.update_heartbeat(db, run_id, counters.to_dict())
return counters
logger.info(
"yandex-sweep run_id=%d center-combos #%d/%d (%s, %.4f, %.4f) "
"radius_m=%d rooms=%d price_ranges=%d",
run_id,
idx,
len(_anchors),
name,
lat,
lon,
radius_m,
len(_rooms_list),
len(_price_ranges),
)
# ── Per-anchor phases под watchdog-таймаутом ────────────────────
# Весь SERP + address-enrich одного anchor'а оборачиваем в wait_for.
# TimeoutError → log + errors_count++ + continue (sweep не прерывается).
anchor_lots: list[ScrapedLot] = []
_anchor_timed_out = False
# Capture loop variables in default args to satisfy B023 (flake8/ruff):
# без этого inner async def захватывает изменяемые переменные цикла
# по ссылке → неправильные значения при asyncio scheduling.
_lat, _lon, _name = lat, lon, name
# _combo_lots_accumulator: per-anchor мутируемый контейнер для on_combo.
# Захватывается по ссылке и в _yandex_anchor_phases и в _on_combo —
# избегаем B023 (loop-variable capture) через явный bind в default-arg.
_combo_lots_accumulator: list[ScrapedLot] = anchor_lots
async def _yandex_anchor_phases(
_a_lat: float = _lat,
_a_lon: float = _lon,
_a_name: str = _name,
_al: list[ScrapedLot] = _combo_lots_accumulator,
) -> None:
"""Все фазы одного anchor'а (SERP + address-enrich).
FIX: инкрементальный save per-combo — on_combo вызывается scraper'ом после
каждого (room×price) combo, сохраняет порцию сразу в БД + обновляет heartbeat.
anchor_lots накапливается (через _al) для address-enrich и price-history.
"""
nonlocal consecutive_failures
# ── Phase 1+2: SERP + инкрементальный save per-combo ──────
def _on_combo(
combo_label: str,
new_lots: list[ScrapedLot],
_accumulator: list[ScrapedLot] = _al,
) -> None:
"""Callback: вызывается scraper'ом после каждого combo с новыми лотами.
Сохраняет порцию новых лотов немедленно (анти-зомби: данные в БД
до конца anchor'а) и обновляет heartbeat для мониторинга прогресса.
"""
_accumulator.extend(new_lots)
counters.lots_fetched += len(new_lots)
try:
ins, upd = save_listings(db, new_lots, run_id=run_id)
counters.lots_inserted += ins
counters.lots_updated += upd
except Exception as save_exc:
counters.errors_count += 1
logger.warning(
"yandex-sweep run_id=%d combo %s: save_listings failed: %s",
run_id,
combo_label,
save_exc,
)
try:
db.rollback()
except Exception:
pass
scrape_runs.update_heartbeat(db, run_id, counters.to_dict())
logger.debug(
"yandex-sweep run_id=%d combo %s: saved %d lots "
"(ins=%d upd=%d total_fetched=%d)",
run_id,
combo_label,
len(new_lots),
counters.lots_inserted,
counters.lots_updated,
counters.lots_fetched,
)
async with YandexRealtyScraper() as scraper:
await scraper.fetch_around_multi_room(
_a_lat,
_a_lon,
radius_m,
max_pages=pages_per_anchor,
rooms_list=_rooms_list,
price_ranges=_price_ranges,
segments=_segments,
on_combo=_on_combo,
)
# _al накоплен on_combo; save уже вызван per-combo.
# Price-history из gate price.previous/trend (graceful — провал
# истории НЕ должен ронять сейв листингов).
if _al:
try:
counters.price_history_rows += record_yandex_price_history(db, _al)
except Exception as ph_exc:
logger.warning(
"yandex-sweep run_id=%d: price-history failed: %s",
run_id,
ph_exc,
)
try:
db.rollback()
except Exception:
pass
consecutive_failures = 0
logger.info(
"yandex-sweep run_id=%d center-combos %s: fetched=%d ins=%d upd=%d",
run_id,
_a_name,
counters.lots_fetched,
counters.lots_inserted,
counters.lots_updated,
)
# ── Phase 3: address-enrich ────────────────────────────────
if not (enrich_address and _al):
return
source_ids = [lot.source_id for lot in _al if lot.source_id]
if not source_ids:
return
id_rows = (
db.execute(
_text("""
SELECT id, source_id, address, source_url
FROM listings
WHERE source = 'yandex'
AND source_id = ANY(CAST(:ids AS text[]))
AND address IS NOT NULL
AND NOT (address ~ ',\\s*\\d+')
AND source_url IS NOT NULL
"""),
{"ids": source_ids},
)
.mappings()
.all()
)
items = [dict(r) for r in id_rows]
counters.address_attempted += len(items)
if items:
from curl_cffi.requests import AsyncSession as _AsyncSession
async with _AsyncSession(
impersonate="chrome120",
timeout=30.0,
proxies=_proxies,
headers={"Accept-Language": "ru-RU,ru;q=0.9,en;q=0.8"},
) as enrich_session:
for eidx, item in enumerate(items):
lid: int = item["id"]
old_addr: str = item["address"]
url: str = item["source_url"]
try:
resp = await enrich_session.get(url, allow_redirects=True)
if resp.status_code != 200:
logger.warning(
"yandex-sweep run_id=%d address-enrich:"
" HTTP %d listing_id=%d",
run_id,
resp.status_code,
lid,
)
counters.address_failed += 1
else:
new_addr = _extract_address_from_title(resp.text)
if (
new_addr is None
or new_addr == old_addr
or not _RE_HAS_HOUSE_NUMBER.search(new_addr)
):
pass
else:
try:
db.execute(
_text("""
UPDATE listings
SET address = :addr,
geocode_tried_at = NULL
WHERE id = CAST(:id AS bigint)
"""),
{"addr": new_addr, "id": lid},
)
db.commit()
counters.address_enriched += 1
logger.info(
"yandex-sweep run_id=%d: "
"address updated listing_id=%d "
"%r -> %r",
run_id,
lid,
old_addr,
new_addr,
)
except Exception as upd_exc:
counters.address_failed += 1
counters.errors_count += 1
logger.warning(
"yandex-sweep run_id=%d: "
"address DB update failed "
"listing_id=%d: %s",
run_id,
lid,
upd_exc,
)
try:
db.rollback()
except Exception:
pass
except Exception as fetch_exc:
counters.address_failed += 1
logger.warning(
"yandex-sweep run_id=%d: address-enrich "
"fetch failed listing_id=%d: %s",
run_id,
lid,
fetch_exc,
)
if eidx < len(items) - 1:
await asyncio.sleep(enrich_delay)
logger.info(
"yandex-sweep run_id=%d anchor %s: address enrich=%d/%d failed=%d",
run_id,
_a_name,
counters.address_enriched,
counters.address_attempted,
counters.address_failed,
)
try:
await asyncio.wait_for(_yandex_anchor_phases(), timeout=_sweep_timeout)
except TimeoutError:
logger.warning(
"yandex-sweep run_id=%d: anchor #%d/%d (%s, %.4f, %.4f) "
"timed out after %.0fs — skipping",
run_id,
idx,
len(_anchors),
name,
lat,
lon,
_sweep_timeout,
)
counters.errors_count += 1
consecutive_failures += 1
_anchor_timed_out = True
if consecutive_failures >= YANDEX_SWEEP_MAX_CONSECUTIVE_FAILURES:
logger.error(
"yandex-sweep run_id=%d ABORT at anchor #%d/%d (%s) — "
"%d consecutive failures (last: timeout): "
"IP likely blocked or proxy hung",
run_id,
idx,
len(_anchors),
name,
consecutive_failures,
)
counters.anchors_done = idx
scrape_runs.update_heartbeat(db, run_id, counters.to_dict())
scrape_runs.mark_banned(
db,
run_id,
f"yandex sweep aborted: {consecutive_failures} consecutive "
f"anchor failures (last: anchor timeout {_sweep_timeout:.0f}s)",
counters.to_dict(),
)
return counters
counters.anchors_done = idx
scrape_runs.update_heartbeat(db, run_id, counters.to_dict())
if idx < len(_anchors):
await asyncio.sleep(inter_anchor_delay)
continue
except Exception as e:
# YandexRealtyScraper не propagate'ит block/rate exceptions — ловим broad.
logger.exception("yandex-sweep run_id=%d: anchor %s SERP failed", run_id, name)
counters.errors_count += 1
consecutive_failures += 1
if consecutive_failures >= YANDEX_SWEEP_MAX_CONSECUTIVE_FAILURES:
logger.error(
"yandex-sweep run_id=%d ABORT at anchor #%d/%d (%s) — "
"%d consecutive failures (IP likely blocked): %s",
run_id,
idx,
len(_anchors),
name,
consecutive_failures,
e,
)
counters.anchors_done = idx
scrape_runs.update_heartbeat(db, run_id, counters.to_dict())
scrape_runs.mark_banned(
db,
run_id,
f"yandex sweep aborted: {consecutive_failures} consecutive "
f"anchor failures (last: {e})",
counters.to_dict(),
)
return counters
counters.anchors_done = idx
scrape_runs.update_heartbeat(db, run_id, counters.to_dict())
# Пауза перед следующим anchor'ом даже при ошибке.
if idx < len(_anchors):
await asyncio.sleep(inter_anchor_delay)
continue
counters.anchors_done = idx
scrape_runs.update_heartbeat(db, run_id, counters.to_dict())
# Доп. пауза между anchor'ами (поверх per-page sleep внутри scraper'а).
if idx < len(_anchors):
await asyncio.sleep(inter_anchor_delay)
scrape_runs.mark_done(db, run_id, counters.to_dict())
logger.info(
"yandex-sweep run_id=%d done: anchors=%d/%d lots=%d (ins=%d/upd=%d) "
"address=%d/%d failed=%d errors=%d",
run_id,
counters.anchors_done,
counters.anchors_total,
counters.lots_fetched,
counters.lots_inserted,
counters.lots_updated,
counters.address_enriched,
counters.address_attempted,
counters.address_failed,
counters.errors_count,
)
return counters
except Exception as exc:
logger.exception("yandex-sweep run_id=%d: fatal error", run_id)
scrape_runs.mark_failed(db, run_id, str(exc), counters.to_dict())
raise
# ── Cian city sweep (#860) ──────────────────────────────────────────────────
@dataclass
class CianCitySweepCounters:
"""Aggregate counters для Cian city sweep run.
Поля симметричны CitySweepCounters, но без imv_* (у Cian нет IMV-фазы).
detail_* — обогащение через cian_detail (вкл. price-history).
houses_* — обогащение newbuilding/ЖК через cian_newbuilding.
"""
anchors_total: int = 0
anchors_done: int = 0
lots_fetched: int = 0
lots_dropped_secondary: int = 0 # вторичка, отброшенная при newbuilding_only=True
lots_inserted: int = 0
lots_updated: int = 0
detail_attempted: int = 0
detail_enriched: int = 0
detail_failed: int = 0
houses_attempted: int = 0
houses_enriched: int = 0
houses_failed: int = 0
errors_count: int = 0
def to_dict(self) -> dict[str, int]:
return {f.name: getattr(self, f.name) for f in fields(self)}
# Defensive abort threshold: если N подряд anchor'ов упали с исключением
# (не просто вернули пустой результат) — прерываем sweep.
CIAN_SWEEP_MAX_CONSECUTIVE_FAILURES = 3
async def run_cian_city_sweep(
db: Session,
*,
run_id: int,
anchors: list[tuple[float, float, str]] | None = None,
radius_m: int = 1500,
pages_per_anchor: int = 3,
request_delay_sec: float = 5.0,
detail_top_n: int = 10,
enrich_houses: bool = True,
newbuilding_only: bool = True,
) -> CianCitySweepCounters:
"""Cian newbuilding city sweep: SERP → detail(+price-history) → newbuilding/houses.
Симметрично run_avito_city_sweep. Использует EKB_ANCHORS (общий список).
newbuilding_only (default True): в SERP-фазе сохраняются только лоты с
listing_segment == "novostroyki" — вторичку отбрасываем, ею авторитетно владеет
run_cian_full_load (exhaustive региональный сбор). Зеркало secondary_only в
fetch_all_secondary, которое фильтрует обратное (!= "novostroyki"). DETAIL-фаза
(обогащение cian-листингов без detail_enriched_at) и HOUSES-фаза (newbuilding-enrich)
не затронуты фильтром.
Фазы per anchor:
1. SERP: CianScraper.fetch_around_multi_room(lat, lon, radius_m, pages=pages_per_anchor)
2. FILTER (newbuilding_only): отбросить вторичку (listing_segment != "novostroyki")
3. SAVE: save_listings(db, lots, run_id=run_id)
4. DETAIL: top-N свежих cian-листингов без detail_enriched_at →
fetch_detail(source_url) + save_detail_enrichment(db, listing_id, enrichment)
(save_detail_enrichment уже пишет price_changes → offer_price_history)
5. HOUSES (если enrich_houses=True): для домов из этого anchor'а с cian_zhk_url →
fetch_newbuilding(zhk_url) + save_newbuilding_enrichment(db, house_id, enr)
Anti-bot:
- CianScraper управляет своей curl_cffi-сессией (прокси уже внутри).
- asyncio.sleep(request_delay_sec) между detail-запросами (jitter ±20%).
- Defensive abort при CIAN_SWEEP_MAX_CONSECUTIVE_FAILURES SERP-ошибок подряд.
Cooperative cancel: scrape_runs.is_cancelled(db, run_id) перед каждым anchor.
mark_done вызывается ВСЕГДА (finally outer).
"""
from app.services.scrapers.cian import CianScraper
from app.services.scrapers.cian_detail import fetch_detail, save_detail_enrichment
from app.services.scrapers.cian_newbuilding import (
fetch_newbuilding,
save_newbuilding_enrichment,
)
_anchors = anchors if anchors is not None else EKB_ANCHORS
counters = CianCitySweepCounters(anchors_total=len(_anchors))
consecutive_failures = 0
try:
for idx, (lat, lon, name) in enumerate(_anchors, start=1):
# Cooperative cancel перед каждым anchor
if scrape_runs.is_cancelled(db, run_id):
logger.info(
"cian-sweep run_id=%d: cancelled at anchor #%d/%d (%s)",
run_id,
idx,
len(_anchors),
name,
)
scrape_runs.update_heartbeat(db, run_id, counters.to_dict())
return counters
logger.info(
"cian-sweep run_id=%d anchor #%d/%d (%s, %.4f, %.4f)",
run_id,
idx,
len(_anchors),
name,
lat,
lon,
)
# ── Per-anchor phases под watchdog-таймаутом ────────────────────
# Весь SERP + detail + houses одного anchor'а оборачиваем в wait_for.
# TimeoutError → log + errors_count++ + continue (sweep не прерывается).
anchor_lots: list[ScrapedLot] = []
# Capture loop variables in default args (B023): prevents stale binding
# if the coroutine is scheduled after the loop variable changes.
_c_lat, _c_lon, _c_name = lat, lon, name
async def _cian_anchor_phases(
_a_lat: float = _c_lat,
_a_lon: float = _c_lon,
_a_name: str = _c_name,
) -> None:
"""Все фазы одного cian anchor'а (SERP + detail + houses)."""
nonlocal anchor_lots, consecutive_failures
from sqlalchemy import text as _text
# ── Phase 1+2: SERP + save ─────────────────────────────────
async with CianScraper() as scraper:
anchor_lots = await scraper.fetch_around_multi_room(
_a_lat,
_a_lon,
radius_m,
pages=pages_per_anchor,
)
counters.lots_fetched += len(anchor_lots)
# ── Фильтр вторички (newbuilding_only) ─────────────────────
# Зеркало secondary_only в fetch_all_secondary: тут оставляем
# только новостройки — вторичкой авторитетно владеет run_cian_full_load.
if newbuilding_only:
_before = len(anchor_lots)
anchor_lots = [
lot for lot in anchor_lots if lot.listing_segment == "novostroyki"
]
counters.lots_dropped_secondary += _before - len(anchor_lots)
if anchor_lots:
inserted, updated = save_listings(db, anchor_lots, run_id=run_id)
counters.lots_inserted += inserted
counters.lots_updated += updated
consecutive_failures = 0
logger.info(
"cian-sweep run_id=%d anchor %s: SERP fetched=%d nb_kept=%d "
"dropped_secondary=%d ins=%d upd=%d",
run_id,
_a_name,
counters.lots_fetched,
len(anchor_lots),
counters.lots_dropped_secondary,
counters.lots_inserted,
counters.lots_updated,
)
# ── Phase 3: detail enrichment (top-N) ────────────────────
if detail_top_n > 0:
priority_rows = (
db.execute(
_text("""
SELECT id, source_url
FROM listings
WHERE source = 'cian'
AND source_url IS NOT NULL
AND detail_enriched_at IS NULL
AND scraped_at > NOW() - INTERVAL '2 hours'
ORDER BY scraped_at DESC NULLS LAST
LIMIT :lim
"""),
{"lim": detail_top_n},
)
.mappings()
.all()
)
for didx, row in enumerate(priority_rows):
listing_id: int = row["id"]
source_url: str = row["source_url"]
counters.detail_attempted += 1
try:
enrichment = await fetch_detail(source_url)
if enrichment is not None:
save_detail_enrichment(db, listing_id, enrichment)
counters.detail_enriched += 1
else:
counters.detail_failed += 1
logger.warning(
"cian-sweep run_id=%d: detail returned None for "
"listing_id=%d url=%s",
run_id,
listing_id,
source_url,
)
except Exception as exc:
counters.detail_failed += 1
counters.errors_count += 1
logger.warning(
"cian-sweep run_id=%d: detail failed listing_id=%d: %s",
run_id,
listing_id,
exc,
)
if didx < len(priority_rows) - 1:
jitter = random.uniform(0.8, 1.2)
await asyncio.sleep(request_delay_sec * jitter)
logger.info(
"cian-sweep run_id=%d anchor %s: detail=%d/%d failed=%d",
run_id,
_a_name,
counters.detail_enriched,
counters.detail_attempted,
counters.detail_failed,
)
# ── Phase 4: newbuilding/houses enrichment ─────────────────
if not (enrich_houses and anchor_lots):
return
cian_nb_ext_ids = {
lot.house_ext_id
for lot in anchor_lots
if lot.house_source == "cian_newbuilding" and lot.house_ext_id
}
if not cian_nb_ext_ids:
return
nb_id_list = list(cian_nb_ext_ids)
try:
house_rows = (
db.execute(
_text("""
SELECT h.id, h.cian_zhk_url
FROM houses h
JOIN house_sources hs ON hs.house_id = h.id
WHERE hs.ext_source = 'cian_newbuilding'
AND hs.ext_id = ANY(CAST(:ids AS text[]))
AND h.cian_zhk_url IS NOT NULL
"""),
{"ids": nb_id_list},
)
.mappings()
.all()
)
except Exception as exc:
logger.exception(
"cian-sweep run_id=%d anchor %s: houses DB query failed — "
"skipping houses phase (%d ids): %s",
run_id,
_a_name,
len(nb_id_list),
exc,
)
counters.houses_failed += len(nb_id_list)
try:
db.rollback()
except Exception:
pass
return
for hidx, hrow in enumerate(house_rows):
house_id_val: int = hrow["id"]
zhk_url: str = hrow["cian_zhk_url"]
counters.houses_attempted += 1
try:
nb_enrichment = await fetch_newbuilding(zhk_url)
if nb_enrichment is not None:
save_newbuilding_enrichment(db, house_id_val, nb_enrichment)
counters.houses_enriched += 1
else:
counters.houses_failed += 1
logger.warning(
"cian-sweep run_id=%d: newbuilding returned None "
"house_id=%d url=%s",
run_id,
house_id_val,
zhk_url,
)
except Exception as exc:
counters.houses_failed += 1
counters.errors_count += 1
logger.warning(
"cian-sweep run_id=%d: houses failed house_id=%d: %s",
run_id,
house_id_val,
exc,
)
try:
db.rollback()
except Exception:
pass
if hidx < len(house_rows) - 1:
await asyncio.sleep(request_delay_sec)
logger.info(
"cian-sweep run_id=%d anchor %s: houses=%d/%d failed=%d",
run_id,
_a_name,
counters.houses_enriched,
counters.houses_attempted,
counters.houses_failed,
)
try:
await asyncio.wait_for(_cian_anchor_phases(), timeout=ANCHOR_TIMEOUT_SEC)
except TimeoutError:
logger.warning(
"cian-sweep run_id=%d: anchor #%d/%d (%s, %.4f, %.4f) "
"timed out after %ds — skipping",
run_id,
idx,
len(_anchors),
name,
lat,
lon,
ANCHOR_TIMEOUT_SEC,
)
counters.errors_count += 1
consecutive_failures += 1
if consecutive_failures >= CIAN_SWEEP_MAX_CONSECUTIVE_FAILURES:
logger.error(
"cian-sweep run_id=%d ABORT at anchor #%d/%d (%s) — "
"%d consecutive failures (last: timeout): "
"IP likely blocked or proxy hung",
run_id,
idx,
len(_anchors),
name,
consecutive_failures,
)
counters.anchors_done = idx
scrape_runs.update_heartbeat(db, run_id, counters.to_dict())
scrape_runs.mark_banned(
db,
run_id,
f"cian sweep aborted: {consecutive_failures} consecutive "
f"anchor failures (last: anchor timeout {ANCHOR_TIMEOUT_SEC}s)",
counters.to_dict(),
)
return counters
counters.anchors_done = idx
scrape_runs.update_heartbeat(db, run_id, counters.to_dict())
continue
except Exception as e:
logger.exception("cian-sweep run_id=%d: anchor %s SERP failed", run_id, name)
counters.errors_count += 1
consecutive_failures += 1
if consecutive_failures >= CIAN_SWEEP_MAX_CONSECUTIVE_FAILURES:
logger.error(
"cian-sweep run_id=%d ABORT at anchor #%d/%d (%s) — "
"%d consecutive SERP failures (IP likely blocked): %s",
run_id,
idx,
len(_anchors),
name,
consecutive_failures,
e,
)
counters.anchors_done = idx
scrape_runs.update_heartbeat(db, run_id, counters.to_dict())
scrape_runs.mark_banned(
db,
run_id,
f"cian sweep aborted: {consecutive_failures} consecutive "
f"anchor SERP failures (last: {e})",
counters.to_dict(),
)
return counters
counters.anchors_done = idx
scrape_runs.update_heartbeat(db, run_id, counters.to_dict())
continue
counters.anchors_done = idx
scrape_runs.update_heartbeat(db, run_id, counters.to_dict())
# Пауза между anchor'ами (поверх per-request sleep внутри scraper'а)
if idx < len(_anchors):
await asyncio.sleep(request_delay_sec)
scrape_runs.mark_done(db, run_id, counters.to_dict())
logger.info(
"cian-sweep run_id=%d done: anchors=%d/%d lots=%d (ins=%d/upd=%d) "
"detail=%d/%d houses=%d/%d errors=%d",
run_id,
counters.anchors_done,
counters.anchors_total,
counters.lots_fetched,
counters.lots_inserted,
counters.lots_updated,
counters.detail_enriched,
counters.detail_attempted,
counters.houses_enriched,
counters.houses_attempted,
counters.errors_count,
)
return counters
except Exception as exc:
logger.exception("cian-sweep run_id=%d: fatal error", run_id)
scrape_runs.mark_failed(db, run_id, str(exc), counters.to_dict())
raise
# ── Cian exhaustive full load (region-wide, без anchor'ов) ─────────────────
@dataclass
class CianFullLoadCounters:
"""Aggregate counters для run_cian_full_load.
Без anchor-полей — единственный региональный проход.
detail_* — опциональное обогащение detail-страниц (если enrich_detail=True).
"""
unique_fetched: int = 0
saved_inserted: int = 0
saved_updated: int = 0
detail_attempted: int = 0
detail_enriched: int = 0
detail_failed: int = 0
errors_count: int = 0
def to_dict(self) -> dict[str, int]:
return {f.name: getattr(self, f.name) for f in fields(self)}
async def run_cian_full_load(
db: Session,
*,
run_id: int,
price_cap_per_bucket: int = 1400,
detail_top_n: int = 0,
request_delay_sec: float = 4.0,
enrich_detail: bool = False,
concurrency: int = 5,
resume_run_id: int | None = None,
) -> CianFullLoadCounters:
"""Exhaustive региональный сбор Cian ЕКБ вторички (БЕЗ anchor'ов).
Cian SERP — region-wide, не geo-bbox. Anchor-loop избыточен (5 anchor'ов
возвращают одну и ту же выдачу ×5). Один проход с партиционированием по
КОМНАТНОСТИ × ЦЕНЕ (адаптивное бинарное деление) охватывает весь ЕКБ.
Параметры:
price_cap_per_bucket: максимум офферов на price-бакет (< Cian SERP-cap ~1500).
detail_top_n: сколько свежих cian-листингов без detail_enriched_at обогащать
(актуально только при enrich_detail=True).
request_delay_sec: задержка между SERP-запросами (перезаписывает scraper default).
enrich_detail: включить detail-обогащение (по умолчанию отключено — тяжело).
concurrency: параллельных page-фетчей в leaf-бакете (default=5).
Инкрементальный save: on_bucket коммитит каждый leaf-бакет в БД сразу после сбора.
Краш/бан больше не приводит к потере всего накопленного.
Cooperative cancel: scrape_runs.is_cancelled проверяется per-bucket.
"""
from app.services.scrapers.cian import CianScraper
counters = CianFullLoadCounters()
# ── Checkpoint/resume: читаем done_buckets из прошлого run ───────────────
skip_set: set[str] = set()
if resume_run_id is not None:
_prev_row = db.execute(
text("SELECT counters FROM scrape_runs WHERE id = CAST(:rid AS bigint)"),
{"rid": resume_run_id},
).fetchone()
if _prev_row is not None and _prev_row.counters:
_prev_counters: dict = (
_prev_row.counters if isinstance(_prev_row.counters, dict) else {}
)
skip_set = set(_prev_counters.get("done_buckets", []))
logger.info(
"cian-full-load run_id=%d: resuming from run %d — %d buckets already done",
run_id,
resume_run_id,
len(skip_set),
)
done: set[str] = set(skip_set) # накапливаем завершённые бакеты этого прогона
def _on_bucket(bucket_key: str, lots: list) -> None: # type: ignore[type-arg]
"""Инкрементальный save после каждого leaf-бакета. Дописывает bucket_key в done."""
nonlocal done
if scrape_runs.is_cancelled(db, run_id):
logger.info("cian-full-load run_id=%d: cancel detected in on_bucket", run_id)
raise RuntimeError("cancelled")
if not lots:
done.add(bucket_key)
scrape_runs.update_heartbeat(
db, run_id, {**counters.to_dict(), "done_buckets": sorted(done)}
)
return
inserted, updated = save_listings(
db, lots, run_id=run_id, skip_seen_today=settings.scraper_skip_seen_today
)
# save_listings вызывает db.commit() внутри — данные в БД сразу
counters.saved_inserted += inserted
counters.saved_updated += updated
counters.unique_fetched += len(lots)
done.add(bucket_key)
scrape_runs.update_heartbeat(
db, run_id, {**counters.to_dict(), "done_buckets": sorted(done)}
)
logger.info(
"cian-full-load run_id=%d: bucket %s saved ins=%d upd=%d total_unique=%d",
run_id,
bucket_key,
inserted,
updated,
counters.unique_fetched,
)
def _on_progress(unique_count: int) -> None:
"""Heartbeat per room-bucket (на случай пустых бакетов без on_bucket вызовов)."""
scrape_runs.update_heartbeat(db, run_id, counters.to_dict())
try:
async with CianScraper() as scraper:
scraper.request_delay_sec = request_delay_sec
await scraper.fetch_all_secondary(
price_cap_per_bucket=price_cap_per_bucket,
concurrency=concurrency,
secondary_only=True,
on_bucket=_on_bucket,
on_progress=_on_progress,
skip_buckets=skip_set if skip_set else None,
)
logger.info(
"cian-full-load run_id=%d: fetch done — unique=%d ins=%d upd=%d",
run_id,
counters.unique_fetched,
counters.saved_inserted,
counters.saved_updated,
)
scrape_runs.update_heartbeat(db, run_id, counters.to_dict())
# ── Опциональный detail-энричмент ─────────────────────────────────────
if enrich_detail and detail_top_n > 0:
from sqlalchemy import text as _text
from app.services.scrapers.cian_detail import (
fetch_detail,
save_detail_enrichment,
)
priority_rows = (
db.execute(
_text("""
SELECT id, source_url
FROM listings
WHERE source = 'cian'
AND source_url IS NOT NULL
AND detail_enriched_at IS NULL
AND scraped_at > NOW() - INTERVAL '4 hours'
ORDER BY scraped_at DESC NULLS LAST
LIMIT :lim
"""),
{"lim": detail_top_n},
)
.mappings()
.all()
)
for didx, row in enumerate(priority_rows):
if scrape_runs.is_cancelled(db, run_id):
logger.info("cian-full-load run_id=%d: cancelled during detail enrich", run_id)
break
listing_id: int = row["id"]
source_url: str = row["source_url"]
counters.detail_attempted += 1
try:
enrichment = await fetch_detail(source_url)
if enrichment is not None:
save_detail_enrichment(db, listing_id, enrichment)
counters.detail_enriched += 1
else:
counters.detail_failed += 1
except Exception as exc:
counters.detail_failed += 1
counters.errors_count += 1
logger.warning(
"cian-full-load run_id=%d: detail failed listing_id=%d: %s",
run_id,
listing_id,
exc,
)
if didx < len(priority_rows) - 1:
await asyncio.sleep(request_delay_sec * random.uniform(0.8, 1.2))
logger.info(
"cian-full-load run_id=%d: detail done — enriched=%d/%d failed=%d",
run_id,
counters.detail_enriched,
counters.detail_attempted,
counters.detail_failed,
)
scrape_runs.update_heartbeat(db, run_id, counters.to_dict())
scrape_runs.mark_done(db, run_id, {**counters.to_dict(), "done_buckets": sorted(done)})
logger.info(
"cian-full-load run_id=%d done: unique=%d ins=%d upd=%d detail=%d/%d errors=%d",
run_id,
counters.unique_fetched,
counters.saved_inserted,
counters.saved_updated,
counters.detail_enriched,
counters.detail_attempted,
counters.errors_count,
)
return counters
except RuntimeError as exc:
# on_bucket кидает RuntimeError("cancelled") при cooperative cancel
if str(exc) == "cancelled":
logger.info(
"cian-full-load run_id=%d: cancelled — partial results unique=%d ins=%d upd=%d",
run_id,
counters.unique_fetched,
counters.saved_inserted,
counters.saved_updated,
)
scrape_runs.mark_done(db, run_id, {**counters.to_dict(), "done_buckets": sorted(done)})
return counters
logger.exception("cian-full-load run_id=%d: fatal error", run_id)
scrape_runs.mark_failed(db, run_id, str(exc), counters.to_dict())
raise
except Exception as exc:
logger.exception("cian-full-load run_id=%d: fatal error", run_id)
scrape_runs.mark_failed(db, run_id, str(exc), counters.to_dict())
raise
# ── Yandex exhaustive full load (region-wide, без anchor'ов) ──────────────────
@dataclass
class YandexFullLoadCounters:
"""Aggregate counters для run_yandex_full_load.
Без detail_* полей — Yandex full load не делает detail-обогащение.
"""
unique_fetched: int = 0
saved_inserted: int = 0
saved_updated: int = 0
price_history_rows: int = 0
errors_count: int = 0
def to_dict(self) -> dict[str, int]:
return {f.name: getattr(self, f.name) for f in fields(self)}
async def run_yandex_full_load(
db: Session,
*,
run_id: int,
price_cap_per_bucket: int = 500,
request_delay_sec: float = 2.0,
concurrency: int = 4,
resume_run_id: int | None = None,
) -> YandexFullLoadCounters:
"""Exhaustive региональный сбор Yandex ЕКБ вторички (БЕЗ anchor'ов).
Yandex SERP — path-based по городу, не geo-bbox. Anchor-loop избыточен
(все anchor'ы дают одну и ту же городскую выдачу). Один проход с
партиционированием по КОМНАТНОСТИ × ЦЕНЕ (адаптивное бинарное деление)
охватывает весь ЕКБ.
Параметры:
price_cap_per_bucket: максимум офферов на price-бакет (< Yandex SERP-cap ~575).
request_delay_sec: задержка между SERP-запросами (перезаписывает scraper default).
concurrency: параллельных page-фетчей в leaf-бакете.
resume_run_id: если задан — читает done_buckets из counters прошлого run и
пропускает уже завершённые бакеты. Передать run_id предыдущего (прерванного)
yandex_full_load прогона для resume. Без этого поля — full walk с нуля.
Инкрементальный save: on_bucket коммитит каждый leaf-бакет в БД сразу после сбора.
Cooperative cancel: scrape_runs.is_cancelled проверяется per-bucket.
"""
from app.services.scrapers.yandex_realty import YandexRealtyScraper
counters = YandexFullLoadCounters()
skip_set: set[str] = set()
if resume_run_id is not None:
_prev_row = db.execute(
text("SELECT counters FROM scrape_runs WHERE id = CAST(:rid AS bigint)"),
{"rid": resume_run_id},
).fetchone()
if _prev_row is not None and _prev_row.counters:
_prev_counters: dict = (
_prev_row.counters if isinstance(_prev_row.counters, dict) else {}
)
skip_set = set(_prev_counters.get("done_buckets", []))
logger.info(
"yandex-full-load run_id=%d: resuming from run %d — %d buckets already done",
run_id,
resume_run_id,
len(skip_set),
)
done: set[str] = set(skip_set)
def _on_bucket(bucket_key: str, lots: list) -> None: # type: ignore[type-arg]
"""Инкрементальный save после каждого leaf-бакета."""
nonlocal done
if scrape_runs.is_cancelled(db, run_id):
logger.info("yandex-full-load run_id=%d: cancel detected in on_bucket", run_id)
raise RuntimeError("cancelled")
if not lots:
done.add(bucket_key)
scrape_runs.update_heartbeat(
db, run_id, {**counters.to_dict(), "done_buckets": sorted(done)}
)
return
inserted, updated = save_listings(
db, lots, run_id=run_id, skip_seen_today=settings.scraper_skip_seen_today
)
# save_listings вызывает db.commit() внутри — данные в БД сразу
counters.saved_inserted += inserted
counters.saved_updated += updated
counters.unique_fetched += len(lots)
# Price-history из gate price.previous/trend (graceful — провал истории
# НЕ должен ронять сейв листингов/бакета).
try:
counters.price_history_rows += record_yandex_price_history(db, lots)
except Exception as ph_exc:
logger.warning(
"yandex-full-load run_id=%d: price-history failed bucket=%s: %s",
run_id,
bucket_key,
ph_exc,
)
try:
db.rollback()
except Exception:
pass
done.add(bucket_key)
scrape_runs.update_heartbeat(
db, run_id, {**counters.to_dict(), "done_buckets": sorted(done)}
)
logger.info(
"yandex-full-load run_id=%d: bucket=%s saved ins=%d upd=%d total_unique=%d",
run_id,
bucket_key,
inserted,
updated,
counters.unique_fetched,
)
def _on_progress(unique_count: int) -> None:
"""Heartbeat per room-bucket (на случай пустых бакетов без on_bucket вызовов)."""
scrape_runs.update_heartbeat(db, run_id, counters.to_dict())
try:
async with YandexRealtyScraper() as scraper:
scraper.request_delay_sec = request_delay_sec
await scraper.fetch_all_secondary(
price_cap_per_bucket=price_cap_per_bucket,
concurrency=concurrency,
on_bucket=_on_bucket,
on_progress=_on_progress,
skip_buckets=skip_set if skip_set else None,
)
logger.info(
"yandex-full-load run_id=%d: fetch done — unique=%d ins=%d upd=%d",
run_id,
counters.unique_fetched,
counters.saved_inserted,
counters.saved_updated,
)
scrape_runs.update_heartbeat(db, run_id, counters.to_dict())
scrape_runs.mark_done(db, run_id, {**counters.to_dict(), "done_buckets": sorted(done)})
logger.info(
"yandex-full-load run_id=%d done: unique=%d ins=%d upd=%d errors=%d",
run_id,
counters.unique_fetched,
counters.saved_inserted,
counters.saved_updated,
counters.errors_count,
)
return counters
except RuntimeError as exc:
# on_bucket кидает RuntimeError("cancelled") при cooperative cancel
if str(exc) == "cancelled":
logger.info(
"yandex-full-load run_id=%d: cancelled — partial results unique=%d ins=%d upd=%d",
run_id,
counters.unique_fetched,
counters.saved_inserted,
counters.saved_updated,
)
scrape_runs.mark_done(db, run_id, {**counters.to_dict(), "done_buckets": sorted(done)})
return counters
logger.exception("yandex-full-load run_id=%d: fatal error", run_id)
scrape_runs.mark_failed(db, run_id, str(exc), counters.to_dict())
raise
except Exception as exc:
logger.exception("yandex-full-load run_id=%d: fatal error", run_id)
scrape_runs.mark_failed(db, run_id, str(exc), counters.to_dict())
raise
# ── Avito exhaustive full load (citywide, без anchor'ов) ──────────────────────
@dataclass
class AvitoFullLoadCounters:
"""Aggregate counters для run_avito_full_load.
Зеркало YandexFullLoadCounters — citywide room×price bisection без anchor/detail.
"""
unique_fetched: int = 0
saved_inserted: int = 0
saved_updated: int = 0
errors_count: int = 0
def to_dict(self) -> dict[str, int]:
return {f.name: getattr(self, f.name) for f in fields(self)}
async def run_avito_full_load(
db: Session,
*,
run_id: int,
price_cap_per_bucket: int = 1400,
request_delay_sec: float = 7.0,
concurrency: int = 5,
secondary_only: bool = True,
resume_run_id: int | None = None,
incremental_days: int | None = None,
) -> AvitoFullLoadCounters:
"""Exhaustive ИЛИ инкрементальный региональный сбор Avito ЕКБ вторички (БЕЗ anchor'ов).
Avito anchor/citywide упирается в SERP-cap ~5000 результатов на запрос →
вторичка катастрофически недобрана. Решение зеркалит cian/yandex full_load:
один проход с партиционированием по КОМНАТНОСТИ × ЦЕНЕ (адаптивное бинарное
деление ценового диапазона через `pmin`/`pmax`) охватывает весь ЕКБ.
Параметры:
price_cap_per_bucket: максимум результатов на price-бакет (< Avito SERP-cap ~5000).
request_delay_sec: задержка между SERP-запросами (перезаписывает scraper default).
concurrency: параллельных page-фетчей в leaf-бакете.
secondary_only: отбрасывать новостройки (listing_segment=="novostroyki") до save.
resume_run_id: если задан — читает done_buckets из counters прошлого run и
пропускает уже завершённые бакеты. Передать run_id предыдущего (прерванного)
avito_full_load прогона для resume. Без поля — full walk с нуля.
incremental_days: если задан — shallow инкрементальный обход. Считается
since = date.today() - timedelta(days=incremental_days) и передаётся в
fetch_all_secondary(since=since): date early-stop останавливает пагинацию
до глубоких страниц (избегает page-36 429-банов). None → since=None →
exhaustive bisection (полный обход, текущее поведение).
Инкрементальный save: on_bucket коммитит каждый leaf-бакет в БД сразу после сбора.
Cooperative cancel: scrape_runs.is_cancelled проверяется per-bucket.
AvitoBlockedError/AvitoRateLimitedError → mark_banned (status='banned'), как в
run_avito_city_sweep — partial results уже сохранены через on_bucket.
"""
counters = AvitoFullLoadCounters()
# ── Инкрементальный режим: since-cutoff для shallow обхода ────────────────
since: date | None = None
if incremental_days is not None:
since = date.today() - timedelta(days=incremental_days)
logger.info(
"avito-full-load run_id=%d: INCREMENTAL mode — since=%s (last %d days)",
run_id,
since.isoformat(),
incremental_days,
)
# ── Checkpoint/resume: читаем done_buckets из прошлого run ───────────────
skip_set: set[str] = set()
if resume_run_id is not None:
_prev_row = db.execute(
text("SELECT counters FROM scrape_runs WHERE id = CAST(:rid AS bigint)"),
{"rid": resume_run_id},
).fetchone()
if _prev_row is not None and _prev_row.counters:
_prev_counters: dict = (
_prev_row.counters if isinstance(_prev_row.counters, dict) else {}
)
skip_set = set(_prev_counters.get("done_buckets", []))
logger.info(
"avito-full-load run_id=%d: resuming from run %d — %d buckets already done",
run_id,
resume_run_id,
len(skip_set),
)
done: set[str] = set(skip_set)
def _on_bucket(bucket_key: str, lots: list) -> None: # type: ignore[type-arg]
"""Инкрементальный save после каждого leaf-бакета."""
nonlocal done
if scrape_runs.is_cancelled(db, run_id):
logger.info("avito-full-load run_id=%d: cancel detected in on_bucket", run_id)
raise RuntimeError("cancelled")
if not lots:
done.add(bucket_key)
scrape_runs.update_heartbeat(
db, run_id, {**counters.to_dict(), "done_buckets": sorted(done)}
)
return
inserted, updated = save_listings(
db, lots, run_id=run_id, skip_seen_today=settings.scraper_skip_seen_today
)
# save_listings вызывает db.commit() внутри — данные в БД сразу
counters.saved_inserted += inserted
counters.saved_updated += updated
counters.unique_fetched += len(lots)
done.add(bucket_key)
scrape_runs.update_heartbeat(
db, run_id, {**counters.to_dict(), "done_buckets": sorted(done)}
)
logger.info(
"avito-full-load run_id=%d: bucket=%s saved ins=%d upd=%d total_unique=%d",
run_id,
bucket_key,
inserted,
updated,
counters.unique_fetched,
)
def _on_progress(unique_count: int) -> None:
"""Heartbeat per room-bucket (на случай пустых бакетов без on_bucket вызовов)."""
scrape_runs.update_heartbeat(db, run_id, counters.to_dict())
try:
async with AvitoScraper() as scraper:
scraper.request_delay_sec = request_delay_sec
await scraper.fetch_all_secondary(
price_cap_per_bucket=price_cap_per_bucket,
concurrency=concurrency,
secondary_only=secondary_only,
on_bucket=_on_bucket,
on_progress=_on_progress,
skip_buckets=skip_set if skip_set else None,
since=since,
)
logger.info(
"avito-full-load run_id=%d: fetch done — unique=%d ins=%d upd=%d",
run_id,
counters.unique_fetched,
counters.saved_inserted,
counters.saved_updated,
)
scrape_runs.update_heartbeat(db, run_id, counters.to_dict())
scrape_runs.mark_done(db, run_id, {**counters.to_dict(), "done_buckets": sorted(done)})
logger.info(
"avito-full-load run_id=%d done: unique=%d ins=%d upd=%d errors=%d",
run_id,
counters.unique_fetched,
counters.saved_inserted,
counters.saved_updated,
counters.errors_count,
)
return counters
except (AvitoBlockedError, AvitoRateLimitedError) as exc:
# IP-бан/rate-limit: partial results уже сохранены через on_bucket.
logger.error("avito-full-load run_id=%d: blocked/rate-limited — %s", run_id, exc)
counters.errors_count += 1
scrape_runs.mark_banned(
db,
run_id,
f"avito full load aborted: {exc}",
{**counters.to_dict(), "done_buckets": sorted(done)},
)
return counters
except RuntimeError as exc:
# on_bucket кидает RuntimeError("cancelled") при cooperative cancel
if str(exc) == "cancelled":
logger.info(
"avito-full-load run_id=%d: cancelled — partial results unique=%d ins=%d upd=%d",
run_id,
counters.unique_fetched,
counters.saved_inserted,
counters.saved_updated,
)
scrape_runs.mark_done(db, run_id, {**counters.to_dict(), "done_buckets": sorted(done)})
return counters
logger.exception("avito-full-load run_id=%d: fatal error", run_id)
scrape_runs.mark_failed(db, run_id, str(exc), counters.to_dict())
raise
except Exception as exc:
logger.exception("avito-full-load run_id=%d: fatal error", run_id)
scrape_runs.mark_failed(db, run_id, str(exc), counters.to_dict())
raise
# ── DomClick citywide sweep (без anchor'ов) ────────────────────────────────
@dataclass
class DomClickCitySweepCounters:
"""Aggregate counters для DomClick citywide sweep run.
DomClick SERP — citywide (city_id), не geo-bbox: anchor-loop отсутствует.
Один проход fetch_city перебирает rooms × pages внутри scraper'а.
Поля симметричны sibling-counters, но без detail_*/houses_* (DomClick SERP-only)
и без anchors_* (нет anchor-фазы). pages_fetched = rooms × pages worst-case.
"""
lots_fetched: int = 0
lots_inserted: int = 0
lots_updated: int = 0
pages_fetched: int = 0
errors_count: int = 0
def to_dict(self) -> dict[str, int]:
return {f.name: getattr(self, f.name) for f in fields(self)}
# Дефолтные параметры sweep'а (EKB city_id=4, все комнатности).
DOMCLICK_DEFAULT_CITY_ID: int = 4
DOMCLICK_DEFAULT_ROOMS: list[int] = [0, 1, 2, 3, 4]
# Оценка времени одного fetch'а (network + camoufox render + parse) для watchdog.
_DOMCLICK_PER_FETCH_S: float = 12.0
# Буфер сверху расчётного бюджета (cold browser start, save-фаза).
_DOMCLICK_SWEEP_BUDGET_S: float = 300.0
async def run_domclick_city_sweep(
db: Session,
*,
run_id: int,
city_id: int = DOMCLICK_DEFAULT_CITY_ID,
rooms: list[int] | None = None,
pages: int = 5,
request_delay_sec: float | None = None,
) -> DomClickCitySweepCounters:
"""DomClick citywide sweep: SERP по city_id × rooms × pages → save → house-match.
Структурно зеркалит run_cian_city_sweep / run_yandex_city_sweep, но CITYWIDE:
DomClick SERP не поддерживает geo-radius (fetch_around → NotImplementedError),
поэтому anchor-loop отсутствует. DomClickScraper.fetch_city уже перебирает
rooms × pages через внутренний BrowserFetcher (camoufox headless) и
дедуплицирует по source_id.
Фазы (единственный citywide проход):
1. SERP: DomClickScraper().fetch_city(city_id, rooms, pages).
2. SAVE: save_listings(db, lots, run_id=run_id) — BaseScraper save path,
триггерит существующий house-match hook (адрес-fingerprint tier, т.к.
DomClick SERP не отдаёт lat/lon).
Watchdog: весь fetch_city оборачивается в asyncio.wait_for. Таймаут считается
из rooms × pages × per-fetch + buffer (минимум ANCHOR_TIMEOUT_SEC).
Cooperative cancel: is_cancelled(db, run_id) проверяется перед SERP-фазой.
mark_done вызывается ВСЕГДА (кроме cancel / fatal). lat=lon=None у всех lots →
house-match использует только address_fingerprint tier.
Возвращает DomClickCitySweepCounters.
"""
from app.services.scrapers.domclick import DomClickScraper
_rooms = rooms if rooms is not None else list(DOMCLICK_DEFAULT_ROOMS)
_resolved_delay = request_delay_sec if request_delay_sec is not None else 8.0
counters = DomClickCitySweepCounters()
# Watchdog-таймаут: число fetch'ей = len(rooms) × pages (worst-case, без
# early-break на пустой странице). Каждый fetch ≈ delay + render/parse overhead.
_num_fetches = max(1, len(_rooms)) * max(1, pages)
_sweep_timeout = max(
ANCHOR_TIMEOUT_SEC,
int(_num_fetches * (_resolved_delay + _DOMCLICK_PER_FETCH_S) + _DOMCLICK_SWEEP_BUDGET_S),
)
try:
# ── Cooperative cancel перед SERP-фазой ──────────────────────────────
if scrape_runs.is_cancelled(db, run_id):
logger.info("domclick-sweep run_id=%d: cancelled before SERP", run_id)
scrape_runs.update_heartbeat(db, run_id, counters.to_dict())
return counters
logger.info(
"domclick-sweep run_id=%d: citywide SERP city_id=%d rooms=%s pages=%d (watchdog %ds)",
run_id,
city_id,
_rooms,
pages,
_sweep_timeout,
)
lots: list[ScrapedLot] = []
async def _domclick_phase() -> None:
"""Единственная citywide-фаза: fetch_city + save."""
nonlocal lots
async with DomClickScraper() as scraper:
if request_delay_sec is not None:
scraper.request_delay_sec = _resolved_delay
lots = await scraper.fetch_city(city_id=city_id, rooms=_rooms, pages=pages)
counters.lots_fetched += len(lots)
if lots:
inserted, updated = save_listings(db, lots, run_id=run_id)
counters.lots_inserted += inserted
counters.lots_updated += updated
try:
await asyncio.wait_for(_domclick_phase(), timeout=_sweep_timeout)
except TimeoutError:
logger.warning(
"domclick-sweep run_id=%d: SERP phase timed out after %ds — partial results",
run_id,
_sweep_timeout,
)
counters.errors_count += 1
except Exception:
logger.exception("domclick-sweep run_id=%d: SERP phase failed", run_id)
counters.errors_count += 1
# pages_fetched: worst-case число запрошенных страниц (rooms × pages).
counters.pages_fetched = _num_fetches
scrape_runs.update_heartbeat(db, run_id, counters.to_dict())
scrape_runs.mark_done(db, run_id, counters.to_dict())
logger.info(
"domclick-sweep run_id=%d done: lots=%d (ins=%d/upd=%d) pages=%d errors=%d",
run_id,
counters.lots_fetched,
counters.lots_inserted,
counters.lots_updated,
counters.pages_fetched,
counters.errors_count,
)
return counters
except Exception as exc:
logger.exception("domclick-sweep run_id=%d: fatal error", run_id)
scrape_runs.mark_failed(db, run_id, str(exc), counters.to_dict())
raise
# ── Public re-exports ───────────────────────────────────────────
__all__ = [
"ANCHOR_TIMEOUT_SEC",
"CIAN_SWEEP_MAX_CONSECUTIVE_FAILURES",
"EKB_ANCHORS",
"CianCitySweepCounters",
"CianFullLoadCounters",
"CitySweepCounters",
"DomClickCitySweepCounters",
"NewbuildingSweepCounters",
"PipelineCounters",
"PipelineResult",
"YandexCitySweepCounters",
"YandexFullLoadCounters",
"run_avito_city_sweep",
"run_avito_newbuilding_sweep",
"run_avito_pipeline",
"run_cian_city_sweep",
"run_cian_full_load",
"run_domclick_city_sweep",
"run_yandex_city_sweep",
"run_yandex_full_load",
]