Some checks failed
CI Trade-In / changes (pull_request) Successful in 14s
CI / changes (pull_request) Successful in 18s
CI Trade-In / browser-tests (pull_request) Has been skipped
CI Trade-In / frontend-checks (pull_request) Has been skipped
CI / backend-tests (pull_request) Has been skipped
CI / frontend-tests (pull_request) Has been skipped
CI / openapi-codegen-check (pull_request) Has been skipped
CI Trade-In / backend-tests (pull_request) Failing after 6m17s
Скрапер один давал 3006 error-строк в сутки из ~3700 по всему Trade-In — ленту перестали читать, и настоящая поломка терялась в ней. Принцип: ожидаемый исход сбора (площадка забанила, пул прокси пуст, капча/недогруз, серия блоков перевалила порог circuit breaker) — это состояние работы против недружелюбного источника, а не инцидент. В error остаётся только неожиданное: изменившаяся вёрстка/схема (Cian markup change), просроченный токен ротации прокси (ASocks 401), неразобранное исключение. Переведено error -> warning в 9 файлах, 12 мест: "СТОП — пул прокси пуст" (avito/domclick/cian_history/yandex_newbuilding_sweep x2), "пул прокси исчерпан" (cian_session, cian_price_history, yandex_address_backfill), FAIL-CLOSED без здорового узла для source (proxy_egress), ABORT по счётчику подтверждённых блоков площадки (avito, domclick, yandex_detail_backfill). Оставлено error намеренно: cookie-алерты Cian/DomClick (#2658, #2674) — они рассчитаны именно на LoggingIntegration(event_level=ERROR) в scheduler_main.py и без него молчат по 37 дней; ABORT по смешанным/soft причинам без единого подтверждённого блока площадки (#3272, #2674/#3196) — это может быть наш баг, а не бан, сигнал сознательно не приглушали. GlitchTip: сентри-интеграция скрапера уже настроена как LoggingIntegration(level=INFO, event_level=ERROR) в scheduler_main.py — отдельной правки sentry_scrub.py не требуется, понижение уровня logger само убирает эти записи из GlitchTip. Итоговая FINISHED-строка со счётчиками (attempted/enriched/blocked/failed) уже существует в каждом detail_backfill — новую не добавлял. Refs #3471
214 lines
8.5 KiB
Python
214 lines
8.5 KiB
Python
"""T8a: Cian offer price-history backfill service.
|
||
|
||
Fetches Cian detail pages (curl_cffi, без Playwright) for listings that have
|
||
no rows in offer_price_history, extracts priceChanges from
|
||
_cianConfig['frontend-offer-card'] defaultState, writes to offer_price_history.
|
||
|
||
Deliberately NOT wired into scheduler — rate-limit risk with datacenter IPs.
|
||
Triggered manually via POST /api/v1/admin/scrape/cian-price-history.
|
||
"""
|
||
|
||
from __future__ import annotations
|
||
|
||
import asyncio
|
||
import logging
|
||
import time
|
||
from dataclasses import dataclass, field
|
||
|
||
# #2306: fetch_detail/save_detail_enrichment migrated to scraper_kit (byte-identical
|
||
# golden-parity была доказана против legacy cian_detail-модуля до его удаления,
|
||
# #2397 Part E2; extract_state/ScrapedLot parity-тесты убраны вместе с остальным
|
||
# legacy scrapers-каталогом, #2397 финальный шаг E — kit единственный живой путь).
|
||
from scraper_kit.providers.cian.detail import fetch_detail, save_detail_enrichment
|
||
from scraper_kit.proxy_errors import NoProxyAvailableError
|
||
from sqlalchemy import text
|
||
from sqlalchemy.orm import Session
|
||
|
||
from app.services.scraper_adapters import (
|
||
RealMatcherAdapter,
|
||
RealProxyProvider,
|
||
RealScraperConfig,
|
||
)
|
||
from app.services.scraper_settings import get_scraper_delay
|
||
|
||
logger = logging.getLogger(__name__)
|
||
|
||
|
||
@dataclass
|
||
class CianPriceHistoryResult:
|
||
checked: int = 0
|
||
saved: int = 0
|
||
skipped: int = 0
|
||
errors: int = 0
|
||
duration_sec: float = field(default=0.0)
|
||
|
||
|
||
async def backfill_cian_price_history(
|
||
db: Session,
|
||
*,
|
||
batch_size: int = 50,
|
||
listing_id: int | None = None,
|
||
) -> CianPriceHistoryResult:
|
||
"""Fetch Cian detail pages and write missing price-history rows.
|
||
|
||
Args:
|
||
db: SQLAlchemy session (caller-owned; commits internally per listing).
|
||
batch_size: max listings to process when listing_id is None.
|
||
listing_id: process a single specific listing (ignores batch_size).
|
||
|
||
Selection query picks cian listings with no existing offer_price_history rows.
|
||
Idempotent: re-run safe via ON CONFLICT DO NOTHING in save_detail_enrichment.
|
||
"""
|
||
result = CianPriceHistoryResult()
|
||
t0 = time.time()
|
||
delay = get_scraper_delay("cian") # default 5.0s
|
||
# Egress через пул с учётом `scrape_proxy_source_bans` (#2830): узел выбирает
|
||
# `curl_proxy_url` внутри `fetch_detail`, он же на выходе возвращает вердикт
|
||
# (mark_banned на CianBlockedError / mark_health / release).
|
||
#
|
||
# Флаг читается из окружения как у всех (#3386 хвост): до #3387 у контейнера
|
||
# `backend` не было `USE_PROXY_POOL_CURL`, и здесь стоял подкласс с зашитым
|
||
# `use_proxy_pool_curl = True` — иначе `curl_proxy_url` игнорировал бы
|
||
# `proxy_provider`. Теперь переменная задана и сервису `backend` (compose), а
|
||
# зашитая константа делала рубильник неотключаемым ровно на этом пути.
|
||
scraper_config = RealScraperConfig()
|
||
proxy_provider = RealProxyProvider()
|
||
|
||
if listing_id is not None:
|
||
rows = (
|
||
db.execute(
|
||
text("""
|
||
SELECT l.id, l.source_url
|
||
FROM listings l
|
||
WHERE l.id = CAST(:lid AS bigint)
|
||
AND l.source = 'cian'
|
||
AND l.source_url IS NOT NULL
|
||
"""),
|
||
{"lid": listing_id},
|
||
)
|
||
.mappings()
|
||
.all()
|
||
)
|
||
else:
|
||
rows = (
|
||
db.execute(
|
||
text("""
|
||
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
|
||
ORDER BY l.id
|
||
LIMIT :lim
|
||
"""),
|
||
{"lim": batch_size},
|
||
)
|
||
.mappings()
|
||
.all()
|
||
)
|
||
|
||
result.checked = len(rows)
|
||
logger.info(
|
||
"cian_price_history backfill: checked=%d delay=%.1fs",
|
||
result.checked,
|
||
delay,
|
||
)
|
||
|
||
for i, row in enumerate(rows):
|
||
lid: int = row["id"]
|
||
url: str = row["source_url"]
|
||
|
||
try:
|
||
# config= обязателен — без него kit fetch_detail идёт напрямую, а без прокси
|
||
# datacenter-IP блокируется Cian (#806). proxy_provider= — узел из пула
|
||
# (#2830): раньше здесь был статичный SCRAPER_PROXY_URL, не знающий про
|
||
# `scrape_proxy_source_bans`, и 403 от отбитого узла никому не сообщался.
|
||
enrichment = await fetch_detail(
|
||
url, config=scraper_config, proxy_provider=proxy_provider
|
||
)
|
||
except NoProxyAvailableError as exc:
|
||
# Fail-closed (#2616): пул пуст/недоступен в проде. Остальные листинги
|
||
# упрутся в то же самое — рвём батч сразу, а не 50 раз по 5 секунд с
|
||
# логом, который читается как «Циан нас блокирует».
|
||
logger.warning(
|
||
"cian_price_history: нет доступного прокси в пуле (%s) — батч прерван "
|
||
"на listing_id=%s (обработано %d из %d)",
|
||
exc,
|
||
lid,
|
||
i,
|
||
len(rows),
|
||
)
|
||
result.errors += 1
|
||
break
|
||
except Exception as exc:
|
||
logger.warning(
|
||
"cian_price_history: fetch failed listing_id=%s url=%s: %s",
|
||
lid,
|
||
url,
|
||
exc,
|
||
)
|
||
result.errors += 1
|
||
await asyncio.sleep(delay)
|
||
continue
|
||
|
||
if enrichment is None:
|
||
logger.warning(
|
||
"cian_price_history: fetch returned None listing_id=%s url=%s",
|
||
lid,
|
||
url,
|
||
)
|
||
result.errors += 1
|
||
await asyncio.sleep(delay)
|
||
continue
|
||
|
||
if not enrichment.price_changes:
|
||
logger.debug("cian_price_history: no price_changes listing_id=%s", lid)
|
||
result.skipped += 1
|
||
else:
|
||
try:
|
||
# Count rows actually inserted: save_detail_enrichment skips
|
||
# changes without change_time/price_rub and uses ON CONFLICT
|
||
# DO NOTHING, so len(price_changes) overcounts on invalid
|
||
# elements or idempotent re-runs. Diff the row count instead.
|
||
before = db.execute(
|
||
text(
|
||
"SELECT COUNT(*) FROM offer_price_history "
|
||
"WHERE listing_id = CAST(:lid AS bigint)"
|
||
),
|
||
{"lid": lid},
|
||
).scalar_one()
|
||
# matcher (#2435): резолвит канонический дом из bti_data (address/geo
|
||
# листинга) — иначе BTI-снапшот распарсен, но выбрасывается.
|
||
save_detail_enrichment(db, lid, enrichment, matcher=RealMatcherAdapter())
|
||
after = db.execute(
|
||
text(
|
||
"SELECT COUNT(*) FROM offer_price_history "
|
||
"WHERE listing_id = CAST(:lid AS bigint)"
|
||
),
|
||
{"lid": lid},
|
||
).scalar_one()
|
||
result.saved += max(0, int(after) - int(before))
|
||
except Exception as exc:
|
||
logger.warning("cian_price_history: save failed listing_id=%s: %s", lid, exc)
|
||
result.errors += 1
|
||
try:
|
||
db.rollback()
|
||
except Exception:
|
||
pass
|
||
await asyncio.sleep(delay)
|
||
continue
|
||
|
||
if i < len(rows) - 1:
|
||
await asyncio.sleep(delay)
|
||
|
||
result.duration_sec = time.time() - t0
|
||
logger.info(
|
||
"cian_price_history done: checked=%d saved=%d skipped=%d errors=%d %.1fs",
|
||
result.checked,
|
||
result.saved,
|
||
result.skipped,
|
||
result.errors,
|
||
result.duration_sec,
|
||
)
|
||
return result
|