perf(tradein): cian full-load incremental save + concurrent pagination #929
4 changed files with 276 additions and 128 deletions
|
|
@ -958,10 +958,19 @@ class CianFullLoadRequest(BaseModel):
|
||||||
),
|
),
|
||||||
)
|
)
|
||||||
request_delay_sec: float = Field(
|
request_delay_sec: float = Field(
|
||||||
default=4.0,
|
default=1.5,
|
||||||
ge=3.0,
|
ge=0.5,
|
||||||
le=15.0,
|
le=15.0,
|
||||||
description="Задержка между SERP-запросами (сек). Увеличить при бане.",
|
description=(
|
||||||
|
"Задержка между SERP-запросами (сек). Выделенный прокси + IP-ротация "
|
||||||
|
"позволяют снизить до 1.5s. Увеличить при бане."
|
||||||
|
),
|
||||||
|
)
|
||||||
|
concurrency: int = Field(
|
||||||
|
default=5,
|
||||||
|
ge=1,
|
||||||
|
le=10,
|
||||||
|
description="Параллельных page-фетчей внутри одного leaf-бакета.",
|
||||||
)
|
)
|
||||||
enrich_detail: bool = Field(
|
enrich_detail: bool = Field(
|
||||||
default=False,
|
default=False,
|
||||||
|
|
@ -1004,6 +1013,7 @@ async def start_cian_full_load(
|
||||||
run_id=run_id,
|
run_id=run_id,
|
||||||
price_cap_per_bucket=payload.price_cap_per_bucket,
|
price_cap_per_bucket=payload.price_cap_per_bucket,
|
||||||
request_delay_sec=payload.request_delay_sec,
|
request_delay_sec=payload.request_delay_sec,
|
||||||
|
concurrency=payload.concurrency,
|
||||||
enrich_detail=payload.enrich_detail,
|
enrich_detail=payload.enrich_detail,
|
||||||
detail_top_n=payload.detail_top_n,
|
detail_top_n=payload.detail_top_n,
|
||||||
)
|
)
|
||||||
|
|
|
||||||
|
|
@ -1557,6 +1557,7 @@ async def run_cian_full_load(
|
||||||
detail_top_n: int = 0,
|
detail_top_n: int = 0,
|
||||||
request_delay_sec: float = 4.0,
|
request_delay_sec: float = 4.0,
|
||||||
enrich_detail: bool = False,
|
enrich_detail: bool = False,
|
||||||
|
concurrency: int = 5,
|
||||||
) -> CianFullLoadCounters:
|
) -> CianFullLoadCounters:
|
||||||
"""Exhaustive региональный сбор Cian ЕКБ вторички (БЕЗ anchor'ов).
|
"""Exhaustive региональный сбор Cian ЕКБ вторички (БЕЗ anchor'ов).
|
||||||
|
|
||||||
|
|
@ -1570,54 +1571,63 @@ async def run_cian_full_load(
|
||||||
(актуально только при enrich_detail=True).
|
(актуально только при enrich_detail=True).
|
||||||
request_delay_sec: задержка между SERP-запросами (перезаписывает scraper default).
|
request_delay_sec: задержка между SERP-запросами (перезаписывает scraper default).
|
||||||
enrich_detail: включить detail-обогащение (по умолчанию отключено — тяжело).
|
enrich_detail: включить detail-обогащение (по умолчанию отключено — тяжело).
|
||||||
|
concurrency: параллельных page-фетчей в leaf-бакете (default=5).
|
||||||
|
|
||||||
Cooperative cancel: scrape_runs.is_cancelled проверяется в on_progress callback
|
Инкрементальный save: on_bucket коммитит каждый leaf-бакет в БД сразу после сбора.
|
||||||
(раз в room-bucket). При cancel — возвращаем частичный результат, mark_done.
|
Краш/бан больше не приводит к потере всего накопленного.
|
||||||
|
Cooperative cancel: scrape_runs.is_cancelled проверяется per-bucket.
|
||||||
"""
|
"""
|
||||||
from app.services.scrapers.cian import CianScraper
|
from app.services.scrapers.cian import CianScraper
|
||||||
|
|
||||||
counters = CianFullLoadCounters()
|
counters = CianFullLoadCounters()
|
||||||
_cancelled = False
|
|
||||||
|
def _on_bucket(lots: list) -> None: # type: ignore[type-arg]
|
||||||
|
"""Инкрементальный save после каждого leaf-бакета."""
|
||||||
|
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:
|
||||||
|
return
|
||||||
|
inserted, updated = save_listings(db, lots, run_id=run_id)
|
||||||
|
# save_listings вызывает db.commit() внутри — данные в БД сразу
|
||||||
|
counters.saved_inserted += inserted
|
||||||
|
counters.saved_updated += updated
|
||||||
|
counters.unique_fetched += len(lots)
|
||||||
|
scrape_runs.update_heartbeat(db, run_id, counters.to_dict())
|
||||||
|
logger.info(
|
||||||
|
"cian-full-load run_id=%d: bucket saved ins=%d upd=%d total_unique=%d",
|
||||||
|
run_id,
|
||||||
|
inserted,
|
||||||
|
updated,
|
||||||
|
counters.unique_fetched,
|
||||||
|
)
|
||||||
|
|
||||||
def _on_progress(unique_count: int) -> None:
|
def _on_progress(unique_count: int) -> None:
|
||||||
nonlocal _cancelled
|
"""Heartbeat per room-bucket (на случай пустых бакетов без on_bucket вызовов)."""
|
||||||
counters.unique_fetched = unique_count
|
|
||||||
scrape_runs.update_heartbeat(db, run_id, counters.to_dict())
|
scrape_runs.update_heartbeat(db, run_id, counters.to_dict())
|
||||||
if scrape_runs.is_cancelled(db, run_id):
|
|
||||||
_cancelled = True
|
|
||||||
logger.info("cian-full-load run_id=%d: cancel detected in on_progress", run_id)
|
|
||||||
|
|
||||||
try:
|
try:
|
||||||
async with CianScraper() as scraper:
|
async with CianScraper() as scraper:
|
||||||
scraper.request_delay_sec = request_delay_sec
|
scraper.request_delay_sec = request_delay_sec
|
||||||
|
|
||||||
lots = await scraper.fetch_all_secondary(
|
await scraper.fetch_all_secondary(
|
||||||
price_cap_per_bucket=price_cap_per_bucket,
|
price_cap_per_bucket=price_cap_per_bucket,
|
||||||
|
concurrency=concurrency,
|
||||||
|
on_bucket=_on_bucket,
|
||||||
on_progress=_on_progress,
|
on_progress=_on_progress,
|
||||||
)
|
)
|
||||||
|
|
||||||
counters.unique_fetched = len(lots)
|
|
||||||
logger.info(
|
logger.info(
|
||||||
"cian-full-load run_id=%d: fetch done — unique=%d cancelled=%s",
|
"cian-full-load run_id=%d: fetch done — unique=%d ins=%d upd=%d",
|
||||||
run_id,
|
run_id,
|
||||||
len(lots),
|
counters.unique_fetched,
|
||||||
_cancelled,
|
counters.saved_inserted,
|
||||||
|
counters.saved_updated,
|
||||||
)
|
)
|
||||||
|
|
||||||
if lots:
|
|
||||||
inserted, updated = save_listings(db, lots, run_id=run_id)
|
|
||||||
counters.saved_inserted = inserted
|
|
||||||
counters.saved_updated = updated
|
|
||||||
logger.info(
|
|
||||||
"cian-full-load run_id=%d: save done — ins=%d upd=%d",
|
|
||||||
run_id,
|
|
||||||
inserted,
|
|
||||||
updated,
|
|
||||||
)
|
|
||||||
scrape_runs.update_heartbeat(db, run_id, counters.to_dict())
|
scrape_runs.update_heartbeat(db, run_id, counters.to_dict())
|
||||||
|
|
||||||
# ── Опциональный detail-энричмент ─────────────────────────────────────
|
# ── Опциональный detail-энричмент ─────────────────────────────────────
|
||||||
if enrich_detail and detail_top_n > 0 and not _cancelled:
|
if enrich_detail and detail_top_n > 0:
|
||||||
from sqlalchemy import text as _text
|
from sqlalchemy import text as _text
|
||||||
|
|
||||||
from app.services.scrapers.cian_detail import (
|
from app.services.scrapers.cian_detail import (
|
||||||
|
|
@ -1645,7 +1655,6 @@ async def run_cian_full_load(
|
||||||
for didx, row in enumerate(priority_rows):
|
for didx, row in enumerate(priority_rows):
|
||||||
if scrape_runs.is_cancelled(db, run_id):
|
if scrape_runs.is_cancelled(db, run_id):
|
||||||
logger.info("cian-full-load run_id=%d: cancelled during detail enrich", run_id)
|
logger.info("cian-full-load run_id=%d: cancelled during detail enrich", run_id)
|
||||||
_cancelled = True
|
|
||||||
break
|
break
|
||||||
listing_id: int = row["id"]
|
listing_id: int = row["id"]
|
||||||
source_url: str = row["source_url"]
|
source_url: str = row["source_url"]
|
||||||
|
|
@ -1680,8 +1689,7 @@ async def run_cian_full_load(
|
||||||
|
|
||||||
scrape_runs.mark_done(db, run_id, counters.to_dict())
|
scrape_runs.mark_done(db, run_id, counters.to_dict())
|
||||||
logger.info(
|
logger.info(
|
||||||
"cian-full-load run_id=%d done: unique=%d ins=%d upd=%d "
|
"cian-full-load run_id=%d done: unique=%d ins=%d upd=%d detail=%d/%d errors=%d",
|
||||||
"detail=%d/%d errors=%d cancelled=%s",
|
|
||||||
run_id,
|
run_id,
|
||||||
counters.unique_fetched,
|
counters.unique_fetched,
|
||||||
counters.saved_inserted,
|
counters.saved_inserted,
|
||||||
|
|
@ -1689,10 +1697,25 @@ async def run_cian_full_load(
|
||||||
counters.detail_enriched,
|
counters.detail_enriched,
|
||||||
counters.detail_attempted,
|
counters.detail_attempted,
|
||||||
counters.errors_count,
|
counters.errors_count,
|
||||||
_cancelled,
|
|
||||||
)
|
)
|
||||||
return counters
|
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())
|
||||||
|
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:
|
except Exception as exc:
|
||||||
logger.exception("cian-full-load run_id=%d: fatal error", run_id)
|
logger.exception("cian-full-load run_id=%d: fatal error", run_id)
|
||||||
scrape_runs.mark_failed(db, run_id, str(exc), counters.to_dict())
|
scrape_runs.mark_failed(db, run_id, str(exc), counters.to_dict())
|
||||||
|
|
|
||||||
|
|
@ -23,6 +23,7 @@ from __future__ import annotations
|
||||||
|
|
||||||
import asyncio
|
import asyncio
|
||||||
import hashlib
|
import hashlib
|
||||||
|
import inspect
|
||||||
import logging
|
import logging
|
||||||
import math
|
import math
|
||||||
from collections.abc import Callable
|
from collections.abc import Callable
|
||||||
|
|
@ -316,6 +317,8 @@ class CianScraper(BaseScraper):
|
||||||
rooms_buckets: list[tuple[int, ...]] | None = None,
|
rooms_buckets: list[tuple[int, ...]] | None = None,
|
||||||
price_cap_per_bucket: int = 1400,
|
price_cap_per_bucket: int = 1400,
|
||||||
max_pages_per_bucket: int = 54,
|
max_pages_per_bucket: int = 54,
|
||||||
|
concurrency: int = 5,
|
||||||
|
on_bucket: Callable[..., Any] | None = None,
|
||||||
on_progress: Callable[[int], None] | None = None,
|
on_progress: Callable[[int], None] | None = None,
|
||||||
) -> list[ScrapedLot]:
|
) -> list[ScrapedLot]:
|
||||||
"""Exhaustive-загрузка Cian ЕКБ вторички через партиционирование КОМНАТНОСТЬ × ЦЕНА.
|
"""Exhaustive-загрузка Cian ЕКБ вторички через партиционирование КОМНАТНОСТЬ × ЦЕНА.
|
||||||
|
|
@ -323,13 +326,17 @@ class CianScraper(BaseScraper):
|
||||||
Обходит Cian SERP-cap (~54 стр/запрос ≈ 1500 результатов на запрос):
|
Обходит Cian SERP-cap (~54 стр/запрос ≈ 1500 результатов на запрос):
|
||||||
внутри каждой комнатности адаптивно бьёт диапазон цены на бакеты так, чтобы
|
внутри каждой комнатности адаптивно бьёт диапазон цены на бакеты так, чтобы
|
||||||
в каждом totalOffers < price_cap_per_bucket → пагинирует бакет полностью.
|
в каждом totalOffers < price_cap_per_bucket → пагинирует бакет полностью.
|
||||||
|
Страницы внутри leaf-бакета запрашиваются параллельно (asyncio.gather + Semaphore).
|
||||||
Дедуп по source_id (dict seen).
|
Дедуп по source_id (dict seen).
|
||||||
|
|
||||||
Параметры:
|
Параметры:
|
||||||
rooms_buckets: список room-кодов Cian для перебора (default: 1-6).
|
rooms_buckets: список room-кодов Cian для перебора (default: 1-6).
|
||||||
price_cap_per_bucket: максимум офферов в бакете перед делением (< 1500).
|
price_cap_per_bucket: максимум офферов в бакете перед делением (< 1500).
|
||||||
max_pages_per_bucket: Cian hard cap ~54; не превышать.
|
max_pages_per_bucket: Cian hard cap ~54; не превышать.
|
||||||
on_progress: опциональный callback(unique_count) для heartbeat.
|
concurrency: максимум параллельных page-фетчей в leaf-бакете (default=5).
|
||||||
|
on_bucket: опциональный callback(list[ScrapedLot]) после каждого leaf-бакета.
|
||||||
|
Может быть async или sync. Если кидает исключение — прерывает прогон.
|
||||||
|
on_progress: опциональный callback(unique_count) для heartbeat (per room-bucket).
|
||||||
|
|
||||||
Возвращает list[ScrapedLot] уникальных лотов (дедуп по source_id/source_url).
|
Возвращает list[ScrapedLot] уникальных лотов (дедуп по source_id/source_url).
|
||||||
"""
|
"""
|
||||||
|
|
@ -339,16 +346,21 @@ class CianScraper(BaseScraper):
|
||||||
for rooms in _buckets:
|
for rooms in _buckets:
|
||||||
room_label = f"room{'_'.join(str(r) for r in rooms)}"
|
room_label = f"room{'_'.join(str(r) for r in rooms)}"
|
||||||
logger.info(
|
logger.info(
|
||||||
"cian exhaustive: starting %s (price_cap=%d)", room_label, price_cap_per_bucket
|
"cian exhaustive: starting %s (price_cap=%d concurrency=%d)",
|
||||||
|
room_label,
|
||||||
|
price_cap_per_bucket,
|
||||||
|
concurrency,
|
||||||
)
|
)
|
||||||
before = len(seen)
|
before = len(seen)
|
||||||
await self._walk_price_range(
|
await self._walk_price_range(
|
||||||
rooms=rooms,
|
rooms=rooms,
|
||||||
lo=0,
|
lo=0,
|
||||||
hi=None,
|
hi=_MAX_PRICE,
|
||||||
seen=seen,
|
seen=seen,
|
||||||
price_cap_per_bucket=price_cap_per_bucket,
|
price_cap_per_bucket=price_cap_per_bucket,
|
||||||
max_pages_per_bucket=max_pages_per_bucket,
|
max_pages_per_bucket=max_pages_per_bucket,
|
||||||
|
concurrency=concurrency,
|
||||||
|
on_bucket=on_bucket,
|
||||||
)
|
)
|
||||||
room_collected = len(seen) - before
|
room_collected = len(seen) - before
|
||||||
logger.info(
|
logger.info(
|
||||||
|
|
@ -368,27 +380,32 @@ class CianScraper(BaseScraper):
|
||||||
*,
|
*,
|
||||||
rooms: tuple[int, ...],
|
rooms: tuple[int, ...],
|
||||||
lo: int,
|
lo: int,
|
||||||
hi: int | None,
|
hi: int,
|
||||||
seen: dict[str, ScrapedLot],
|
seen: dict[str, ScrapedLot],
|
||||||
price_cap_per_bucket: int,
|
price_cap_per_bucket: int,
|
||||||
max_pages_per_bucket: int,
|
max_pages_per_bucket: int,
|
||||||
|
concurrency: int = 5,
|
||||||
|
on_bucket: Callable[..., Any] | None = None,
|
||||||
_depth: int = 0,
|
_depth: int = 0,
|
||||||
) -> None:
|
) -> None:
|
||||||
"""Рекурсивное адаптивное бинарное партиционирование ценового диапазона [lo, hi].
|
"""Рекурсивное адаптивное бинарное партиционирование ценового диапазона [lo, hi].
|
||||||
|
|
||||||
Алгоритм:
|
Алгоритм:
|
||||||
1. Запросить page=1 с min_price=lo, max_price=hi → totalOffers.
|
1. Запросить page=1 с min_price=lo, max_price=hi → totalOffers.
|
||||||
2. Если totalOffers <= cap → пагинировать бакет полностью.
|
2. Если totalOffers <= cap → пагинировать leaf-бакет параллельно (concurrency).
|
||||||
3. Если totalOffers > cap → разбить бакет пополам (рекурсия).
|
3. Если totalOffers > cap → разбить бакет пополам (рекурсия).
|
||||||
Guard: hi - lo < _MIN_BRACKET → пагинировать как есть (логируем WARNING).
|
Guard: hi - lo < _MIN_BRACKET → пагинировать как есть (логируем WARNING).
|
||||||
|
|
||||||
|
После пагинации leaf-бакета вызывает on_bucket(bucket_lots) если задан.
|
||||||
|
on_bucket может быть async или sync. Исключение в on_bucket прерывает прогон.
|
||||||
"""
|
"""
|
||||||
# Нормализация: None-hi → _MAX_PRICE для первого уровня
|
effective_hi = hi
|
||||||
effective_hi = hi if hi is not None else _MAX_PRICE
|
|
||||||
|
|
||||||
room_label = f"room{'_'.join(str(r) for r in rooms)}"
|
room_label = f"room{'_'.join(str(r) for r in rooms)}"
|
||||||
|
_lo_param = lo if lo > 0 else None
|
||||||
|
|
||||||
# ── Шаг 1: probe page 1 ────────────────────────────────────────────────
|
# ── Шаг 1: probe page 1 ────────────────────────────────────────────────
|
||||||
html = await self._fetch_page_html(rooms, 1, lo if lo > 0 else None, hi)
|
html = await self._fetch_page_html(rooms, 1, _lo_param, hi)
|
||||||
await self.sleep_between_requests()
|
await self.sleep_between_requests()
|
||||||
|
|
||||||
total: int | None = None
|
total: int | None = None
|
||||||
|
|
@ -398,7 +415,7 @@ class CianScraper(BaseScraper):
|
||||||
# Ретрай на captcha/ошибку: rotate IP + 1 retry
|
# Ретрай на captcha/ошибку: rotate IP + 1 retry
|
||||||
if total is None:
|
if total is None:
|
||||||
logger.warning(
|
logger.warning(
|
||||||
"cian: totalOffers=None for %s [%d, %s] depth=%d — rotating IP + retry",
|
"cian: totalOffers=None for %s [%d, %d] depth=%d — rotating IP + retry",
|
||||||
room_label,
|
room_label,
|
||||||
lo,
|
lo,
|
||||||
hi,
|
hi,
|
||||||
|
|
@ -406,14 +423,14 @@ class CianScraper(BaseScraper):
|
||||||
)
|
)
|
||||||
rotated = await self._rotate_ip()
|
rotated = await self._rotate_ip()
|
||||||
if rotated:
|
if rotated:
|
||||||
html = await self._fetch_page_html(rooms, 1, lo if lo > 0 else None, hi)
|
html = await self._fetch_page_html(rooms, 1, _lo_param, hi)
|
||||||
await self.sleep_between_requests()
|
await self.sleep_between_requests()
|
||||||
if html is not None:
|
if html is not None:
|
||||||
total = self._extract_total_offers(html)
|
total = self._extract_total_offers(html)
|
||||||
|
|
||||||
if total is None:
|
if total is None:
|
||||||
logger.error(
|
logger.error(
|
||||||
"cian: skipping bucket %s [%d, %s] — totalOffers unavailable after retry",
|
"cian: skipping bucket %s [%d, %d] — totalOffers unavailable after retry",
|
||||||
room_label,
|
room_label,
|
||||||
lo,
|
lo,
|
||||||
hi,
|
hi,
|
||||||
|
|
@ -421,7 +438,7 @@ class CianScraper(BaseScraper):
|
||||||
return
|
return
|
||||||
|
|
||||||
logger.info(
|
logger.info(
|
||||||
"cian: %s [%d, %s] totalOffers=%d depth=%d",
|
"cian: %s [%d, %d] totalOffers=%d depth=%d",
|
||||||
room_label,
|
room_label,
|
||||||
lo,
|
lo,
|
||||||
hi,
|
hi,
|
||||||
|
|
@ -439,7 +456,7 @@ class CianScraper(BaseScraper):
|
||||||
|
|
||||||
if need_split and too_narrow:
|
if need_split and too_narrow:
|
||||||
logger.warning(
|
logger.warning(
|
||||||
"cian: %s [%d, %s] totalOffers=%d > cap=%d but bracket=%d < MIN_BRACKET=%d "
|
"cian: %s [%d, %d] totalOffers=%d > cap=%d but bracket=%d < MIN_BRACKET=%d "
|
||||||
"— paginating as-is (tail loss ~%d)",
|
"— paginating as-is (tail loss ~%d)",
|
||||||
room_label,
|
room_label,
|
||||||
lo,
|
lo,
|
||||||
|
|
@ -453,18 +470,6 @@ class CianScraper(BaseScraper):
|
||||||
need_split = False # принудительно пагинируем
|
need_split = False # принудительно пагинируем
|
||||||
|
|
||||||
if need_split:
|
if need_split:
|
||||||
# При hi=None: сначала устанавливаем hi = _MAX_PRICE и рекурсируем
|
|
||||||
if hi is None:
|
|
||||||
await self._walk_price_range(
|
|
||||||
rooms=rooms,
|
|
||||||
lo=lo,
|
|
||||||
hi=_MAX_PRICE,
|
|
||||||
seen=seen,
|
|
||||||
price_cap_per_bucket=price_cap_per_bucket,
|
|
||||||
max_pages_per_bucket=max_pages_per_bucket,
|
|
||||||
_depth=_depth + 1,
|
|
||||||
)
|
|
||||||
return
|
|
||||||
mid = (lo + effective_hi) // 2
|
mid = (lo + effective_hi) // 2
|
||||||
# [lo, mid]
|
# [lo, mid]
|
||||||
await self._walk_price_range(
|
await self._walk_price_range(
|
||||||
|
|
@ -474,6 +479,8 @@ class CianScraper(BaseScraper):
|
||||||
seen=seen,
|
seen=seen,
|
||||||
price_cap_per_bucket=price_cap_per_bucket,
|
price_cap_per_bucket=price_cap_per_bucket,
|
||||||
max_pages_per_bucket=max_pages_per_bucket,
|
max_pages_per_bucket=max_pages_per_bucket,
|
||||||
|
concurrency=concurrency,
|
||||||
|
on_bucket=on_bucket,
|
||||||
_depth=_depth + 1,
|
_depth=_depth + 1,
|
||||||
)
|
)
|
||||||
# [mid+1, hi]
|
# [mid+1, hi]
|
||||||
|
|
@ -484,101 +491,82 @@ class CianScraper(BaseScraper):
|
||||||
seen=seen,
|
seen=seen,
|
||||||
price_cap_per_bucket=price_cap_per_bucket,
|
price_cap_per_bucket=price_cap_per_bucket,
|
||||||
max_pages_per_bucket=max_pages_per_bucket,
|
max_pages_per_bucket=max_pages_per_bucket,
|
||||||
|
concurrency=concurrency,
|
||||||
|
on_bucket=on_bucket,
|
||||||
_depth=_depth + 1,
|
_depth=_depth + 1,
|
||||||
)
|
)
|
||||||
return
|
return
|
||||||
|
|
||||||
# ── Пагинация бакета ───────────────────────────────────────────────────
|
# ── Параллельная пагинация leaf-бакета ────────────────────────────────
|
||||||
max_pages = min(
|
max_pages = min(
|
||||||
math.ceil(total / _CIAN_OFFERS_PER_PAGE),
|
math.ceil(total / _CIAN_OFFERS_PER_PAGE),
|
||||||
max_pages_per_bucket,
|
max_pages_per_bucket,
|
||||||
)
|
)
|
||||||
|
|
||||||
# Страница 1 уже есть (html из probe выше)
|
# Страница 1 уже есть (html из probe выше); остальные — параллельно.
|
||||||
collected_this_bucket = 0
|
sem = asyncio.Semaphore(concurrency)
|
||||||
if html:
|
|
||||||
page1_lots = self._parse_serp_html(html)
|
|
||||||
for lot in page1_lots:
|
|
||||||
key = lot.source_id or lot.source_url
|
|
||||||
if key:
|
|
||||||
seen[key] = lot
|
|
||||||
collected_this_bucket += len(page1_lots)
|
|
||||||
if not page1_lots:
|
|
||||||
logger.info("cian: %s [%d, %s] page=1 empty — early stop", room_label, lo, hi)
|
|
||||||
return
|
|
||||||
|
|
||||||
for page in range(2, max_pages + 1):
|
async def _one_page(p: int) -> list[ScrapedLot]:
|
||||||
page_html: str | None = None
|
if p == 1:
|
||||||
try:
|
# Используем уже полученный HTML от probe
|
||||||
page_html = await self._fetch_page_html(rooms, page, lo if lo > 0 else None, hi)
|
return self._parse_serp_html(html) if html else []
|
||||||
await self.sleep_between_requests()
|
async with sem:
|
||||||
except Exception:
|
page_html = await self._fetch_page_html(rooms, p, _lo_param, hi)
|
||||||
logger.exception(
|
await asyncio.sleep(self.request_delay_sec)
|
||||||
"cian: fetch failed %s [%d, %s] page=%d — rotating IP",
|
if page_html is None:
|
||||||
|
logger.warning(
|
||||||
|
"cian: page_html=None %s [%d, %d] page=%d — skipping page",
|
||||||
|
room_label,
|
||||||
|
lo,
|
||||||
|
hi,
|
||||||
|
p,
|
||||||
|
)
|
||||||
|
return []
|
||||||
|
return self._parse_serp_html(page_html)
|
||||||
|
|
||||||
|
page_results = await asyncio.gather(
|
||||||
|
*[_one_page(p) for p in range(1, max_pages + 1)],
|
||||||
|
return_exceptions=True,
|
||||||
|
)
|
||||||
|
|
||||||
|
bucket_lots: list[ScrapedLot] = []
|
||||||
|
for p_idx, res in enumerate(page_results, start=1):
|
||||||
|
if isinstance(res, BaseException):
|
||||||
|
logger.warning(
|
||||||
|
"cian: page exception %s [%d, %d] page=%d — %r",
|
||||||
room_label,
|
room_label,
|
||||||
lo,
|
lo,
|
||||||
hi,
|
hi,
|
||||||
page,
|
p_idx,
|
||||||
|
res,
|
||||||
)
|
)
|
||||||
await self._rotate_ip()
|
else:
|
||||||
try:
|
bucket_lots.extend(res)
|
||||||
page_html = await self._fetch_page_html(rooms, page, lo if lo > 0 else None, hi)
|
|
||||||
await self.sleep_between_requests()
|
|
||||||
except Exception:
|
|
||||||
logger.error(
|
|
||||||
"cian: retry failed %s [%d, %s] page=%d — breaking bucket",
|
|
||||||
room_label,
|
|
||||||
lo,
|
|
||||||
hi,
|
|
||||||
page,
|
|
||||||
)
|
|
||||||
break
|
|
||||||
|
|
||||||
if page_html is None:
|
collected_this_bucket = len(bucket_lots)
|
||||||
# Повторная попытка после ротации
|
|
||||||
rotated = await self._rotate_ip()
|
|
||||||
if rotated:
|
|
||||||
try:
|
|
||||||
page_html = await self._fetch_page_html(
|
|
||||||
rooms, page, lo if lo > 0 else None, hi
|
|
||||||
)
|
|
||||||
await self.sleep_between_requests()
|
|
||||||
except Exception:
|
|
||||||
pass
|
|
||||||
if page_html is None:
|
|
||||||
logger.error(
|
|
||||||
"cian: page_html=None %s [%d, %s] page=%d — breaking bucket",
|
|
||||||
room_label,
|
|
||||||
lo,
|
|
||||||
hi,
|
|
||||||
page,
|
|
||||||
)
|
|
||||||
break
|
|
||||||
|
|
||||||
page_lots = self._parse_serp_html(page_html)
|
# Дедуп в общий seen
|
||||||
if not page_lots:
|
for lot in bucket_lots:
|
||||||
logger.info(
|
key = lot.source_id or lot.source_url
|
||||||
"cian: %s [%d, %s] page=%d empty — early stop", room_label, lo, hi, page
|
if key:
|
||||||
)
|
seen[key] = lot
|
||||||
break
|
|
||||||
|
|
||||||
for lot in page_lots:
|
|
||||||
key = lot.source_id or lot.source_url
|
|
||||||
if key:
|
|
||||||
seen[key] = lot
|
|
||||||
collected_this_bucket += len(page_lots)
|
|
||||||
|
|
||||||
logger.info(
|
logger.info(
|
||||||
"cian: %s [%d, %s] paginated=%d/%d pages collected=%d unique_total=%d",
|
"cian: %s [%d, %d] paginated=%d pages collected=%d unique_total=%d",
|
||||||
room_label,
|
room_label,
|
||||||
lo,
|
lo,
|
||||||
hi,
|
hi,
|
||||||
min(max_pages, page if "page" in locals() else 1),
|
|
||||||
max_pages,
|
max_pages,
|
||||||
collected_this_bucket,
|
collected_this_bucket,
|
||||||
len(seen),
|
len(seen),
|
||||||
)
|
)
|
||||||
|
|
||||||
|
# ── on_bucket callback: инкрементальный save ──────────────────────────
|
||||||
|
if on_bucket is not None and bucket_lots:
|
||||||
|
res_cb = on_bucket(bucket_lots)
|
||||||
|
if inspect.isawaitable(res_cb):
|
||||||
|
await res_cb
|
||||||
|
|
||||||
def _parse_serp_html(self, html: str) -> list[ScrapedLot]:
|
def _parse_serp_html(self, html: str) -> list[ScrapedLot]:
|
||||||
"""Извлечь offers из Cian Redux state.
|
"""Извлечь offers из Cian Redux state.
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -202,5 +202,132 @@ async def test_fetch_all_secondary_min_bracket_guard(scraper: CianScraper) -> No
|
||||||
max_pages_per_bucket=54,
|
max_pages_per_bucket=54,
|
||||||
)
|
)
|
||||||
|
|
||||||
# Тест прошёл если нет RuntimeError (guard сработал)
|
# Guard сработал: нет RuntimeError (бесконечная рекурсия не возникла).
|
||||||
assert call_count <= 10, f"Слишком много вызовов ({call_count}) — guard не сработал"
|
# Параллельная пагинация запускает до max_pages_per_bucket=54 страниц единовременно
|
||||||
|
# (все страницы бакета параллельны), поэтому call_count может быть до 54.
|
||||||
|
# Важно что нет рекурсии (bracket < MIN_BRACKET → один проход пагинации, не деление).
|
||||||
|
assert call_count <= 55, f"Слишком много вызовов ({call_count}) — guard не сработал"
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.asyncio
|
||||||
|
async def test_fetch_all_secondary_on_bucket_called_per_leaf(scraper: CianScraper) -> None:
|
||||||
|
"""on_bucket вызывается после каждого leaf-бакета с лотами бакета."""
|
||||||
|
# totalOffers=56 ≤ cap → leaf-бакет, пагинируется параллельно 2 страницы
|
||||||
|
bucket_calls: list[list[ScrapedLot]] = []
|
||||||
|
|
||||||
|
def fake_on_bucket(lots: list[ScrapedLot]) -> None:
|
||||||
|
bucket_calls.append(list(lots))
|
||||||
|
|
||||||
|
async def fake_fetch_page_html(
|
||||||
|
rooms: tuple, page: int, min_price: int | None, max_price: int | None
|
||||||
|
) -> str:
|
||||||
|
return f"<html>page={page}</html>"
|
||||||
|
|
||||||
|
def fake_extract_total_offers(html: str) -> int | None:
|
||||||
|
return 56 # ceil(56/28) = 2 страницы
|
||||||
|
|
||||||
|
def fake_parse_serp_html(html: str) -> list[ScrapedLot]:
|
||||||
|
import re
|
||||||
|
|
||||||
|
page_m = re.search(r"page=(\d+)", html)
|
||||||
|
page = int(page_m.group(1)) if page_m else 1
|
||||||
|
return [_make_lot(f"lot_p{page}_{i}") for i in range(10)]
|
||||||
|
|
||||||
|
with (
|
||||||
|
patch.object(scraper, "_fetch_page_html", side_effect=fake_fetch_page_html),
|
||||||
|
patch.object(scraper, "_extract_total_offers", side_effect=fake_extract_total_offers),
|
||||||
|
patch.object(scraper, "_parse_serp_html", side_effect=fake_parse_serp_html),
|
||||||
|
patch.object(scraper, "sleep_between_requests", new_callable=AsyncMock),
|
||||||
|
patch.object(scraper, "_rotate_ip", return_value=False),
|
||||||
|
):
|
||||||
|
await scraper.fetch_all_secondary(
|
||||||
|
rooms_buckets=[(1,)],
|
||||||
|
price_cap_per_bucket=1400,
|
||||||
|
concurrency=3,
|
||||||
|
on_bucket=fake_on_bucket,
|
||||||
|
)
|
||||||
|
|
||||||
|
# on_bucket вызван ровно 1 раз (один leaf-бакет для одной комнатности)
|
||||||
|
assert len(bucket_calls) == 1, f"Ожидался 1 вызов on_bucket, получено {len(bucket_calls)}"
|
||||||
|
# Лоты из обеих страниц переданы в on_bucket
|
||||||
|
assert (
|
||||||
|
len(bucket_calls[0]) == 20
|
||||||
|
), f"Ожидалось 20 лотов в on_bucket (2 стр × 10), получено {len(bucket_calls[0])}"
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.asyncio
|
||||||
|
async def test_fetch_all_secondary_on_bucket_cancel_stops_run(scraper: CianScraper) -> None:
|
||||||
|
"""Если on_bucket кидает RuntimeError('cancelled') — прогон прерывается."""
|
||||||
|
|
||||||
|
def cancel_on_bucket(lots: list[ScrapedLot]) -> None:
|
||||||
|
raise RuntimeError("cancelled")
|
||||||
|
|
||||||
|
async def fake_fetch_page_html(
|
||||||
|
rooms: tuple, page: int, min_price: int | None, max_price: int | None
|
||||||
|
) -> str:
|
||||||
|
return f"<html>page={page}</html>"
|
||||||
|
|
||||||
|
def fake_extract_total_offers(html: str) -> int | None:
|
||||||
|
return 28 # один лист → 1 страница
|
||||||
|
|
||||||
|
def fake_parse_serp_html(html: str) -> list[ScrapedLot]:
|
||||||
|
return [_make_lot("lot_1")]
|
||||||
|
|
||||||
|
with (
|
||||||
|
patch.object(scraper, "_fetch_page_html", side_effect=fake_fetch_page_html),
|
||||||
|
patch.object(scraper, "_extract_total_offers", side_effect=fake_extract_total_offers),
|
||||||
|
patch.object(scraper, "_parse_serp_html", side_effect=fake_parse_serp_html),
|
||||||
|
patch.object(scraper, "sleep_between_requests", new_callable=AsyncMock),
|
||||||
|
patch.object(scraper, "_rotate_ip", return_value=False),
|
||||||
|
):
|
||||||
|
with pytest.raises(RuntimeError, match="cancelled"):
|
||||||
|
await scraper.fetch_all_secondary(
|
||||||
|
rooms_buckets=[(1,), (2,)],
|
||||||
|
price_cap_per_bucket=1400,
|
||||||
|
on_bucket=cancel_on_bucket,
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.asyncio
|
||||||
|
async def test_fetch_all_secondary_concurrent_pages_deduped(scraper: CianScraper) -> None:
|
||||||
|
"""Параллельная пагинация не ломает дедупликацию по source_id."""
|
||||||
|
fetch_pages: list[int] = []
|
||||||
|
|
||||||
|
async def fake_fetch_page_html(
|
||||||
|
rooms: tuple, page: int, min_price: int | None, max_price: int | None
|
||||||
|
) -> str:
|
||||||
|
fetch_pages.append(page)
|
||||||
|
return f"<html>page={page}</html>"
|
||||||
|
|
||||||
|
def fake_extract_total_offers(html: str) -> int | None:
|
||||||
|
return 84 # ceil(84/28) = 3 страницы
|
||||||
|
|
||||||
|
def fake_parse_serp_html(html: str) -> list[ScrapedLot]:
|
||||||
|
import re
|
||||||
|
|
||||||
|
page_m = re.search(r"page=(\d+)", html)
|
||||||
|
page = int(page_m.group(1)) if page_m else 1
|
||||||
|
# Страница 1 и 2 дают уникальные лоты; страница 3 дублирует страницу 1
|
||||||
|
if page == 3:
|
||||||
|
return [_make_lot(f"lot_p1_{i}") for i in range(10)] # дубликаты
|
||||||
|
return [_make_lot(f"lot_p{page}_{i}") for i in range(10)]
|
||||||
|
|
||||||
|
with (
|
||||||
|
patch.object(scraper, "_fetch_page_html", side_effect=fake_fetch_page_html),
|
||||||
|
patch.object(scraper, "_extract_total_offers", side_effect=fake_extract_total_offers),
|
||||||
|
patch.object(scraper, "_parse_serp_html", side_effect=fake_parse_serp_html),
|
||||||
|
patch.object(scraper, "sleep_between_requests", new_callable=AsyncMock),
|
||||||
|
patch.object(scraper, "_rotate_ip", return_value=False),
|
||||||
|
):
|
||||||
|
lots = await scraper.fetch_all_secondary(
|
||||||
|
rooms_buckets=[(1,)],
|
||||||
|
price_cap_per_bucket=1400,
|
||||||
|
concurrency=5,
|
||||||
|
)
|
||||||
|
|
||||||
|
# Страницы 1+2 = 20 уникальных, страница 3 дублирует → итого 20
|
||||||
|
assert len(lots) == 20, f"Ожидалось 20 уникальных лотов (дедуп), получено {len(lots)}"
|
||||||
|
# Все 3 страницы запрошены (параллельно или нет — нам важен факт)
|
||||||
|
assert 3 in fetch_pages or 3 in [
|
||||||
|
p for p in fetch_pages
|
||||||
|
], f"Страница 3 не была запрошена, fetch_pages={fetch_pages}"
|
||||||
|
|
|
||||||
Loading…
Add table
Reference in a new issue