feat(avito): dedicated newbuilding (novostroyka) citywide sweep
This commit is contained in:
parent
bf91eb0be0
commit
1a658a65a6
6 changed files with 552 additions and 8 deletions
|
|
@ -9,6 +9,8 @@
|
||||||
|
|
||||||
Sources:
|
Sources:
|
||||||
- avito_city_sweep → run_avito_city_sweep (scrape_pipeline.py)
|
- avito_city_sweep → run_avito_city_sweep (scrape_pipeline.py)
|
||||||
|
- avito_newbuilding_sweep → run_avito_newbuilding_sweep (scrape_pipeline.py; dedicated
|
||||||
|
novostroyka-filtered citywide SERP → save, 100% new-build cards)
|
||||||
- yandex_city_sweep → run_yandex_city_sweep (scrape_pipeline.py, #561; shipped DORMANT)
|
- yandex_city_sweep → run_yandex_city_sweep (scrape_pipeline.py, #561; shipped DORMANT)
|
||||||
- cian_history_backfill → run_cian_backfill (this module, #560)
|
- cian_history_backfill → run_cian_backfill (this module, #560)
|
||||||
- rosreestr_dkp_import → import_rosreestr_dkp (this module, #563)
|
- rosreestr_dkp_import → import_rosreestr_dkp (this module, #563)
|
||||||
|
|
@ -85,6 +87,7 @@ from app.core.db import SessionLocal
|
||||||
from app.services import scrape_runs as runs_mod
|
from app.services import scrape_runs as runs_mod
|
||||||
from app.services.scrape_pipeline import (
|
from app.services.scrape_pipeline import (
|
||||||
run_avito_city_sweep,
|
run_avito_city_sweep,
|
||||||
|
run_avito_newbuilding_sweep,
|
||||||
run_cian_city_sweep,
|
run_cian_city_sweep,
|
||||||
run_domclick_city_sweep,
|
run_domclick_city_sweep,
|
||||||
run_yandex_city_sweep,
|
run_yandex_city_sweep,
|
||||||
|
|
@ -299,6 +302,44 @@ async def trigger_avito_city_sweep_run(db: Session, schedule_row: dict[str, Any]
|
||||||
return run_id
|
return run_id
|
||||||
|
|
||||||
|
|
||||||
|
async def trigger_avito_newbuilding_sweep_run(
|
||||||
|
db: Session, schedule_row: dict[str, Any]
|
||||||
|
) -> int | None:
|
||||||
|
"""Создать новый scrape_runs + launch run_avito_newbuilding_sweep в asyncio.create_task.
|
||||||
|
|
||||||
|
Зеркало trigger_avito_city_sweep_run, но citywide novostroyka-обход (без anchor/
|
||||||
|
houses/detail-параметров): dedicated novostroyka-SERP → save.
|
||||||
|
|
||||||
|
Returns run_id (или None если skip — есть running run).
|
||||||
|
"""
|
||||||
|
run_id = _claim_run(db, schedule_row)
|
||||||
|
if run_id is None:
|
||||||
|
return None
|
||||||
|
|
||||||
|
params = schedule_row.get("default_params") or {}
|
||||||
|
|
||||||
|
# Spawn asyncio task — pass NEW session (avoid sharing with this scheduler tick)
|
||||||
|
async def _run() -> None:
|
||||||
|
run_db = SessionLocal()
|
||||||
|
try:
|
||||||
|
await run_avito_newbuilding_sweep(
|
||||||
|
run_db,
|
||||||
|
run_id=run_id,
|
||||||
|
pages=int(params.get("pages", 20)),
|
||||||
|
request_delay_sec=float(params.get("request_delay_sec", 7.0)),
|
||||||
|
)
|
||||||
|
except Exception:
|
||||||
|
logger.exception("scheduler: run_avito_newbuilding_sweep crashed run_id=%d", run_id)
|
||||||
|
finally:
|
||||||
|
run_db.close()
|
||||||
|
|
||||||
|
task = asyncio.create_task(_run())
|
||||||
|
# Keep reference to avoid GC before task completes (RUF006)
|
||||||
|
task.add_done_callback(lambda t: t.exception() if not t.cancelled() else None)
|
||||||
|
logger.info("scheduler: triggered newbuilding_sweep run_id=%d", run_id)
|
||||||
|
return run_id
|
||||||
|
|
||||||
|
|
||||||
async def trigger_yandex_city_sweep_run(db: Session, schedule_row: dict[str, Any]) -> int | None:
|
async def trigger_yandex_city_sweep_run(db: Session, schedule_row: dict[str, Any]) -> int | None:
|
||||||
"""Создать новый scrape_runs + launch run_yandex_city_sweep в asyncio.create_task.
|
"""Создать новый scrape_runs + launch run_yandex_city_sweep в asyncio.create_task.
|
||||||
|
|
||||||
|
|
@ -1350,6 +1391,8 @@ async def scheduler_loop() -> None:
|
||||||
source = sch["source"]
|
source = sch["source"]
|
||||||
if source == "avito_city_sweep":
|
if source == "avito_city_sweep":
|
||||||
await trigger_avito_city_sweep_run(db, sch)
|
await trigger_avito_city_sweep_run(db, sch)
|
||||||
|
elif source == "avito_newbuilding_sweep":
|
||||||
|
await trigger_avito_newbuilding_sweep_run(db, sch)
|
||||||
elif source == "yandex_city_sweep":
|
elif source == "yandex_city_sweep":
|
||||||
await trigger_yandex_city_sweep_run(db, sch)
|
await trigger_yandex_city_sweep_run(db, sch)
|
||||||
elif source == "cian_history_backfill":
|
elif source == "cian_history_backfill":
|
||||||
|
|
|
||||||
|
|
@ -990,6 +990,114 @@ async def run_avito_city_sweep(
|
||||||
raise
|
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 ───────────────────────────────────────────
|
# ── Yandex city sweep ───────────────────────────────────────────
|
||||||
|
|
||||||
|
|
||||||
|
|
@ -2266,8 +2374,7 @@ async def run_domclick_city_sweep(
|
||||||
return counters
|
return counters
|
||||||
|
|
||||||
logger.info(
|
logger.info(
|
||||||
"domclick-sweep run_id=%d: citywide SERP city_id=%d rooms=%s pages=%d "
|
"domclick-sweep run_id=%d: citywide SERP city_id=%d rooms=%s pages=%d (watchdog %ds)",
|
||||||
"(watchdog %ds)",
|
|
||||||
run_id,
|
run_id,
|
||||||
city_id,
|
city_id,
|
||||||
_rooms,
|
_rooms,
|
||||||
|
|
@ -2333,11 +2440,13 @@ __all__ = [
|
||||||
"CianFullLoadCounters",
|
"CianFullLoadCounters",
|
||||||
"CitySweepCounters",
|
"CitySweepCounters",
|
||||||
"DomClickCitySweepCounters",
|
"DomClickCitySweepCounters",
|
||||||
|
"NewbuildingSweepCounters",
|
||||||
"PipelineCounters",
|
"PipelineCounters",
|
||||||
"PipelineResult",
|
"PipelineResult",
|
||||||
"YandexCitySweepCounters",
|
"YandexCitySweepCounters",
|
||||||
"YandexFullLoadCounters",
|
"YandexFullLoadCounters",
|
||||||
"run_avito_city_sweep",
|
"run_avito_city_sweep",
|
||||||
|
"run_avito_newbuilding_sweep",
|
||||||
"run_avito_pipeline",
|
"run_avito_pipeline",
|
||||||
"run_cian_city_sweep",
|
"run_cian_city_sweep",
|
||||||
"run_cian_full_load",
|
"run_cian_full_load",
|
||||||
|
|
|
||||||
|
|
@ -26,6 +26,7 @@ from __future__ import annotations
|
||||||
import asyncio
|
import asyncio
|
||||||
import logging
|
import logging
|
||||||
import re
|
import re
|
||||||
|
from collections.abc import Callable
|
||||||
from datetime import date, datetime, timedelta, timezone
|
from datetime import date, datetime, timedelta, timezone
|
||||||
from typing import Any
|
from typing import Any
|
||||||
from urllib.parse import urlencode, urljoin, urlparse, urlunparse
|
from urllib.parse import urlencode, urljoin, urlparse, urlunparse
|
||||||
|
|
@ -216,6 +217,12 @@ def _build_sort_timestamp_map(html: str) -> dict[str, date]:
|
||||||
# маркерам в начале документа (firewall-страница объёмна, не сканируем целиком).
|
# маркерам в начале документа (firewall-страница объёмна, не сканируем целиком).
|
||||||
_FIREWALL_MARKERS = ("доступ ограничен", "проблема с ip", "firewall-container")
|
_FIREWALL_MARKERS = ("доступ ограничен", "проблема с ip", "firewall-container")
|
||||||
|
|
||||||
|
# EKB-фильтр «Новостройка» (slug категории kvartiry/prodam). Кодирует
|
||||||
|
# ASgB-параметры выборки только новостроек — SERP возвращает 100% new-build
|
||||||
|
# карточки (с data-marker="item-development-name"). Извлечён из реального
|
||||||
|
# search-URL: /ekaterinburg/kvartiry/prodam/novostroyka-ASgBAgICAkSSA8YQ5geOUg
|
||||||
|
NOVOSTROYKA_SLUG = "novostroyka-ASgBAgICAkSSA8YQ5geOUg"
|
||||||
|
|
||||||
|
|
||||||
def _is_firewall_page(html: str) -> bool:
|
def _is_firewall_page(html: str) -> bool:
|
||||||
"""True если Avito вернул firewall-страницу IP-блока (на HTTP 200)."""
|
"""True если Avito вернул firewall-страницу IP-блока (на HTTP 200)."""
|
||||||
|
|
@ -500,6 +507,16 @@ class AvitoScraper(BaseScraper):
|
||||||
params = {"s": 104, "p": page}
|
params = {"s": 104, "p": page}
|
||||||
return f"{self.base_url}/ekaterinburg/kvartiry/prodam-ASgBAgICAUSSA8YQ?{urlencode(params)}"
|
return f"{self.base_url}/ekaterinburg/kvartiry/prodam-ASgBAgICAUSSA8YQ?{urlencode(params)}"
|
||||||
|
|
||||||
|
def _build_newbuilding_url(self, page: int = 1) -> str:
|
||||||
|
"""URL ЕКБ-выборки только новостроек (novostroyka-filter), сортировка по дате.
|
||||||
|
|
||||||
|
Тот же shape параметров, что _build_citywide_url (s=104, p=page), но slug
|
||||||
|
категории NOVOSTROYKA_SLUG отдаёт исключительно new-build карточки.
|
||||||
|
"""
|
||||||
|
params = {"s": 104, "p": page}
|
||||||
|
path = f"/ekaterinburg/kvartiry/prodam/{NOVOSTROYKA_SLUG}"
|
||||||
|
return f"{self.base_url}{path}?{urlencode(params)}"
|
||||||
|
|
||||||
def _build_rooms_url(self, room_slug: str, page: int = 1) -> str:
|
def _build_rooms_url(self, room_slug: str, page: int = 1) -> str:
|
||||||
"""T6: URL с фильтром по комнатности для всего ЕКБ (no geo).
|
"""T6: URL с фильтром по комнатности для всего ЕКБ (no geo).
|
||||||
|
|
||||||
|
|
@ -532,6 +549,30 @@ class AvitoScraper(BaseScraper):
|
||||||
Returns:
|
Returns:
|
||||||
Список ScrapedLot (все страницы, дедуп по source_id).
|
Список ScrapedLot (все страницы, дедуп по source_id).
|
||||||
"""
|
"""
|
||||||
|
all_lots = await self._paginate_sweep(
|
||||||
|
pages,
|
||||||
|
self._build_citywide_url,
|
||||||
|
label="citywide",
|
||||||
|
delay_override_sec=delay_override_sec,
|
||||||
|
)
|
||||||
|
logger.info("avito fetch_city_wide pages=%d total_lots=%d", pages, len(all_lots))
|
||||||
|
return all_lots
|
||||||
|
|
||||||
|
async def _paginate_sweep(
|
||||||
|
self,
|
||||||
|
pages: int,
|
||||||
|
url_builder: Callable[[int], str],
|
||||||
|
*,
|
||||||
|
label: str,
|
||||||
|
delay_override_sec: float | None = None,
|
||||||
|
) -> list[ScrapedLot]:
|
||||||
|
"""Общий paginated-обход ЕКБ (citywide / novostroyka) с break-on-empty.
|
||||||
|
|
||||||
|
url_builder(page) — функция построения URL страницы (citywide или
|
||||||
|
novostroyka). Сохраняет anti-block pipeline (_fetch_serp_html: firewall,
|
||||||
|
IP rotation, ретраи), дедуп по source_id, break-on-empty. label — только
|
||||||
|
для логов. Не меняет наблюдаемое поведение fetch_city_wide.
|
||||||
|
"""
|
||||||
if delay_override_sec is not None:
|
if delay_override_sec is not None:
|
||||||
self.request_delay_sec = delay_override_sec
|
self.request_delay_sec = delay_override_sec
|
||||||
|
|
||||||
|
|
@ -539,21 +580,21 @@ class AvitoScraper(BaseScraper):
|
||||||
seen_ids: set[str] = set()
|
seen_ids: set[str] = set()
|
||||||
|
|
||||||
for page in range(1, pages + 1):
|
for page in range(1, pages + 1):
|
||||||
url = self._build_citywide_url(page)
|
url = url_builder(page)
|
||||||
try:
|
try:
|
||||||
html = await self._fetch_serp_html(url, page)
|
html = await self._fetch_serp_html(url, page)
|
||||||
except (AvitoBlockedError, AvitoRateLimitedError):
|
except (AvitoBlockedError, AvitoRateLimitedError):
|
||||||
raise
|
raise
|
||||||
except Exception:
|
except Exception:
|
||||||
logger.exception("avito citywide page=%d fetch failed url=%s", page, url)
|
logger.exception("avito %s page=%d fetch failed url=%s", label, page, url)
|
||||||
break
|
break
|
||||||
if html is None:
|
if html is None:
|
||||||
logger.info("avito citywide page=%d: non-200 — end of pagination", page)
|
logger.info("avito %s page=%d: non-200 — end of pagination", label, page)
|
||||||
break
|
break
|
||||||
|
|
||||||
lots = self._parse_html(html, source_url_base=url)
|
lots = self._parse_html(html, source_url_base=url)
|
||||||
if not lots:
|
if not lots:
|
||||||
logger.info("avito citywide page=%d: 0 lots — end of pagination", page)
|
logger.info("avito %s page=%d: 0 lots — end of pagination", label, page)
|
||||||
break
|
break
|
||||||
|
|
||||||
new_lots = [lot for lot in lots if lot.source_id not in seen_ids]
|
new_lots = [lot for lot in lots if lot.source_id not in seen_ids]
|
||||||
|
|
@ -562,7 +603,8 @@ class AvitoScraper(BaseScraper):
|
||||||
seen_ids.add(lot.source_id)
|
seen_ids.add(lot.source_id)
|
||||||
all_lots.extend(new_lots)
|
all_lots.extend(new_lots)
|
||||||
logger.info(
|
logger.info(
|
||||||
"avito citywide page=%d: %d lots (%d new, %d dedup'd, total=%d)",
|
"avito %s page=%d: %d lots (%d new, %d dedup'd, total=%d)",
|
||||||
|
label,
|
||||||
page,
|
page,
|
||||||
len(lots),
|
len(lots),
|
||||||
len(new_lots),
|
len(new_lots),
|
||||||
|
|
@ -571,8 +613,37 @@ class AvitoScraper(BaseScraper):
|
||||||
)
|
)
|
||||||
if page < pages:
|
if page < pages:
|
||||||
await self.sleep_between_requests()
|
await self.sleep_between_requests()
|
||||||
|
return all_lots
|
||||||
|
|
||||||
logger.info("avito fetch_city_wide pages=%d total_lots=%d", pages, len(all_lots))
|
async def fetch_newbuildings(
|
||||||
|
self,
|
||||||
|
pages: int = 30,
|
||||||
|
*,
|
||||||
|
delay_override_sec: float | None = None,
|
||||||
|
) -> list[ScrapedLot]:
|
||||||
|
"""Обход ЕКБ-выборки только новостроек (novostroyka-filter), paginated.
|
||||||
|
|
||||||
|
Avito отдаёт dedicated SERP с фильтром «Новостройка» — 100% new-build
|
||||||
|
карточек (с маркером застройщика → listing_segment="novostroyki",
|
||||||
|
newbuilding_id/newbuilding_url заполнены). Citywide (без geo/anchor).
|
||||||
|
Возвращает дедуплицированный список (по source_id), break-on-empty.
|
||||||
|
Сохраняет весь anti-block pipeline. Не трогает fetch_city_wide/fetch_around.
|
||||||
|
|
||||||
|
Args:
|
||||||
|
pages: максимальное число страниц (default 30).
|
||||||
|
delay_override_sec: если задан — переопределяет request_delay_sec для
|
||||||
|
этого вызова.
|
||||||
|
|
||||||
|
Returns:
|
||||||
|
Список ScrapedLot новостроек (все страницы, дедуп по source_id).
|
||||||
|
"""
|
||||||
|
all_lots = await self._paginate_sweep(
|
||||||
|
pages,
|
||||||
|
self._build_newbuilding_url,
|
||||||
|
label="newbuilding",
|
||||||
|
delay_override_sec=delay_override_sec,
|
||||||
|
)
|
||||||
|
logger.info("avito fetch_newbuildings pages=%d total_lots=%d", pages, len(all_lots))
|
||||||
return all_lots
|
return all_lots
|
||||||
|
|
||||||
async def fetch_by_rooms(
|
async def fetch_by_rooms(
|
||||||
|
|
|
||||||
|
|
@ -0,0 +1,41 @@
|
||||||
|
-- 123_avito_newbuilding_sweep_schedule.sql
|
||||||
|
-- Seed row для Avito newbuilding (novostroyka) citywide sweep.
|
||||||
|
--
|
||||||
|
-- Avito отдаёт dedicated novostroyka-filtered SERP (slug
|
||||||
|
-- novostroyka-ASgBAgICAkSSA8YQ5geOUg) — 100% new-build карточки с маркером
|
||||||
|
-- застройщика (listing_segment='novostroyki', newbuilding_id заполнен парсером).
|
||||||
|
-- Текущий general sweep ловит новостройки лишь инцидентно (~20 строк); этот
|
||||||
|
-- dedicated обход даёт сотни. Citywide (без anchor/geo), save через стандартный путь.
|
||||||
|
--
|
||||||
|
-- default_params:
|
||||||
|
-- pages — макс. число страниц SERP за прогон.
|
||||||
|
-- request_delay_sec — пауза между страницами (anti-block).
|
||||||
|
--
|
||||||
|
-- next_run_at = завтра 02:00 UTC — чтобы НЕ выстрелить мгновенно на деплое (enabled
|
||||||
|
-- row с NULL next_run_at трактуется due немедленно); первый прогон в окне 02:00-05:00.
|
||||||
|
--
|
||||||
|
-- ЗАВИСИМОСТИ: 052_scrape_schedules.sql (таблица + UNIQUE(source)).
|
||||||
|
-- Idempotent: ON CONFLICT (source) DO NOTHING — безопасно запускать повторно.
|
||||||
|
|
||||||
|
BEGIN;
|
||||||
|
|
||||||
|
INSERT INTO scrape_schedules (
|
||||||
|
source,
|
||||||
|
enabled,
|
||||||
|
window_start_hour,
|
||||||
|
window_end_hour,
|
||||||
|
next_run_at,
|
||||||
|
default_params
|
||||||
|
)
|
||||||
|
VALUES
|
||||||
|
(
|
||||||
|
'avito_newbuilding_sweep',
|
||||||
|
true,
|
||||||
|
2,
|
||||||
|
5,
|
||||||
|
((CURRENT_DATE + INTERVAL '1 day') + make_interval(hours => 2)) AT TIME ZONE 'UTC',
|
||||||
|
'{"pages": 20, "request_delay_sec": 8}'::jsonb
|
||||||
|
)
|
||||||
|
ON CONFLICT (source) DO NOTHING;
|
||||||
|
|
||||||
|
COMMIT;
|
||||||
35
tradein-mvp/backend/tests/fixtures/avito_serp_novostroyka.html
vendored
Normal file
35
tradein-mvp/backend/tests/fixtures/avito_serp_novostroyka.html
vendored
Normal file
|
|
@ -0,0 +1,35 @@
|
||||||
|
<!doctype html>
|
||||||
|
<!--
|
||||||
|
Minimal Avito novostroyka (new-build) SERP fixture, distilled from the dedicated
|
||||||
|
novostroyka-filtered SERP (/ekaterinburg/kvartiry/prodam/novostroyka-...). Каждая
|
||||||
|
карточка несёт data-marker="item-development-name" (название ЖК/застройщика) +
|
||||||
|
якорь на ЖК-сабдомен https://zhk-<slug>-ekaterinburg.avito.ru → avito.py::_parse_html
|
||||||
|
классифицирует их как listing_segment='novostroyki' и заполняет newbuilding_id/url.
|
||||||
|
-->
|
||||||
|
<html lang="ru"><head><meta charset="utf-8"><title>Avito novostroyka SERP fixture</title></head>
|
||||||
|
<body>
|
||||||
|
<div class="items-items">
|
||||||
|
<div data-marker="item" data-item-id="9100000001" id="i9100000001">
|
||||||
|
<a data-marker="item-title" href="/ekaterinburg/kvartiry/1-k._kvartira_40m_510et._9100000001">1-к. квартира, 40 м², 5/10 эт.</a>
|
||||||
|
<meta itemprop="price" content="6800000">
|
||||||
|
<div data-marker="item-line"><div data-marker="item-address"><p>ул. Меридиан, 1</p></div></div>
|
||||||
|
<div data-marker="item-development-name">ЖК «Меридиан»</div>
|
||||||
|
<a href="https://zhk-meridian-ekaterinburg.avito.ru?context=abc">ЖК «Меридиан»</a>
|
||||||
|
<p data-marker="item-date">2 часа назад</p>
|
||||||
|
</div>
|
||||||
|
<div data-marker="item" data-item-id="9100000002" id="i9100000002">
|
||||||
|
<a data-marker="item-title" href="/ekaterinburg/kvartiry/2-k._kvartira_62m_812et._9100000002">2-к. квартира, 62 м², 8/12 эт.</a>
|
||||||
|
<meta itemprop="price" content="9100000">
|
||||||
|
<div data-marker="item-line"><div data-marker="item-address"><p>ул. Северная, 7</p></div></div>
|
||||||
|
<div data-marker="item-development-name">ЖК «Северный квартал»</div>
|
||||||
|
<a href="https://zhk-severnyy-kvartal-ekaterinburg.avito.ru?context=def">ЖК «Северный квартал»</a>
|
||||||
|
<p data-marker="item-date">3 часа назад</p>
|
||||||
|
</div>
|
||||||
|
</div>
|
||||||
|
<script>
|
||||||
|
window.__avito_state_fragment = {"catalog":{"items":[
|
||||||
|
{"id":9100000001,"category":{"id":24,"slug":"kvartiry"},"addressDetailed":{"locationName":"Екатеринбург"},"sortTimeStamp":1700000000000,"allowTimeStamp":1700000000000},
|
||||||
|
{"id":9100000002,"category":{"id":24,"slug":"kvartiry"},"addressDetailed":{"locationName":"Екатеринбург"},"sortTimeStamp":1700086400000,"allowTimeStamp":1700086400000}
|
||||||
|
]}};
|
||||||
|
</script>
|
||||||
|
</body></html>
|
||||||
245
tradein-mvp/backend/tests/test_avito_newbuilding_sweep.py
Normal file
245
tradein-mvp/backend/tests/test_avito_newbuilding_sweep.py
Normal file
|
|
@ -0,0 +1,245 @@
|
||||||
|
"""Avito newbuilding (novostroyka) citywide sweep.
|
||||||
|
|
||||||
|
Тесты (без сети, без реального DB):
|
||||||
|
1. _build_newbuilding_url — novostroyka-slug в пути, s=104, page param, без geo.
|
||||||
|
2. _parse_html на novostroyka-SERP fixture → listing_segment='novostroyki' +
|
||||||
|
newbuilding_id заполнен (DOM-маркер застройщика + zhk-якорь).
|
||||||
|
3. fetch_newbuildings break-on-empty + dedup (mirror fetch_city_wide).
|
||||||
|
4. NewbuildingSweepCounters — defaults + to_dict.
|
||||||
|
5. run_avito_newbuilding_sweep wiring — мокаем fetch_newbuildings + save_listings,
|
||||||
|
assert save вызван и SERP-счётчики записаны (mark_done).
|
||||||
|
"""
|
||||||
|
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import os
|
||||||
|
from pathlib import Path
|
||||||
|
from typing import Any
|
||||||
|
from unittest.mock import AsyncMock, MagicMock, patch
|
||||||
|
|
||||||
|
import pytest
|
||||||
|
|
||||||
|
os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost:5432/test")
|
||||||
|
|
||||||
|
from app.services import scrape_pipeline
|
||||||
|
from app.services.scrapers.avito import NOVOSTROYKA_SLUG, AvitoScraper
|
||||||
|
from app.services.scrapers.base import ScrapedLot
|
||||||
|
|
||||||
|
FIXTURE = Path(__file__).parent / "fixtures" / "avito_serp_novostroyka.html"
|
||||||
|
|
||||||
|
|
||||||
|
def _make_lot(source_id: str) -> ScrapedLot:
|
||||||
|
return ScrapedLot(
|
||||||
|
source="avito",
|
||||||
|
source_url=f"https://www.avito.ru/ekaterinburg/kvartiry/{source_id}",
|
||||||
|
source_id=source_id,
|
||||||
|
price_rub=6_000_000,
|
||||||
|
listing_segment="novostroyki",
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
# ── 1. _build_newbuilding_url ───────────────────────────────────────────────
|
||||||
|
|
||||||
|
|
||||||
|
def test_build_newbuilding_url_contains_slug_and_sort() -> None:
|
||||||
|
s = AvitoScraper()
|
||||||
|
url = s._build_newbuilding_url(page=1)
|
||||||
|
assert NOVOSTROYKA_SLUG in url
|
||||||
|
assert "novostroyka-" in url
|
||||||
|
assert "geoCoords" not in url
|
||||||
|
assert "radius" not in url
|
||||||
|
assert "s=104" in url
|
||||||
|
assert "p=1" in url
|
||||||
|
assert "ekaterinburg/kvartiry/prodam/" in url
|
||||||
|
|
||||||
|
|
||||||
|
def test_build_newbuilding_url_page_param() -> None:
|
||||||
|
s = AvitoScraper()
|
||||||
|
url = s._build_newbuilding_url(page=2)
|
||||||
|
assert "novostroyka-" in url
|
||||||
|
assert "p=2" in url
|
||||||
|
|
||||||
|
|
||||||
|
def test_build_newbuilding_url_differs_from_citywide() -> None:
|
||||||
|
s = AvitoScraper()
|
||||||
|
assert s._build_newbuilding_url(page=1) != s._build_citywide_url(page=1)
|
||||||
|
|
||||||
|
|
||||||
|
# ── 2. _parse_html classifies novostroyka cards ─────────────────────────────
|
||||||
|
|
||||||
|
|
||||||
|
def test_parse_novostroyka_fixture_segments_and_newbuilding_id() -> None:
|
||||||
|
html = FIXTURE.read_text(encoding="utf-8")
|
||||||
|
s = AvitoScraper()
|
||||||
|
lots = s._parse_html(html, source_url_base=s._build_newbuilding_url(page=1))
|
||||||
|
|
||||||
|
assert len(lots) == 2
|
||||||
|
for lot in lots:
|
||||||
|
assert lot.listing_segment == "novostroyki"
|
||||||
|
assert lot.newbuilding_id is not None
|
||||||
|
assert lot.newbuilding_url is not None
|
||||||
|
|
||||||
|
ids = {lot.newbuilding_id for lot in lots}
|
||||||
|
assert "meridian-ekaterinburg" in ids
|
||||||
|
assert "severnyy-kvartal-ekaterinburg" in ids
|
||||||
|
|
||||||
|
|
||||||
|
# ── 3. fetch_newbuildings break-on-empty + dedup ────────────────────────────
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.asyncio
|
||||||
|
async def test_fetch_newbuildings_break_on_empty() -> None:
|
||||||
|
s = AvitoScraper()
|
||||||
|
pages_data: list[list[ScrapedLot]] = [[_make_lot("A"), _make_lot("B")], []]
|
||||||
|
call_n = 0
|
||||||
|
|
||||||
|
async def mock_html(url: str, page: int) -> str:
|
||||||
|
assert "novostroyka-" in url
|
||||||
|
return f"<html>{page}</html>"
|
||||||
|
|
||||||
|
def mock_parse(html: str, source_url_base: str) -> list[ScrapedLot]:
|
||||||
|
nonlocal call_n
|
||||||
|
idx = call_n
|
||||||
|
call_n += 1
|
||||||
|
return pages_data[idx] if idx < len(pages_data) else []
|
||||||
|
|
||||||
|
with patch.object(s, "_fetch_serp_html", side_effect=mock_html):
|
||||||
|
with patch.object(s, "_parse_html", side_effect=mock_parse):
|
||||||
|
with patch.object(s, "sleep_between_requests", new=AsyncMock()):
|
||||||
|
result = await s.fetch_newbuildings(pages=10, delay_override_sec=0)
|
||||||
|
|
||||||
|
assert len(result) == 2
|
||||||
|
assert call_n == 2 # page1 + page2(empty) → stop
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.asyncio
|
||||||
|
async def test_fetch_newbuildings_dedup_by_source_id() -> None:
|
||||||
|
s = AvitoScraper()
|
||||||
|
dup = _make_lot("SAME")
|
||||||
|
pages_data: list[list[ScrapedLot]] = [[dup], [dup], []]
|
||||||
|
call_n = 0
|
||||||
|
|
||||||
|
async def mock_html(url: str, page: int) -> str:
|
||||||
|
return f"<html>{page}</html>"
|
||||||
|
|
||||||
|
def mock_parse(html: str, source_url_base: str) -> list[ScrapedLot]:
|
||||||
|
nonlocal call_n
|
||||||
|
idx = call_n
|
||||||
|
call_n += 1
|
||||||
|
return pages_data[idx] if idx < len(pages_data) else []
|
||||||
|
|
||||||
|
with patch.object(s, "_fetch_serp_html", side_effect=mock_html):
|
||||||
|
with patch.object(s, "_parse_html", side_effect=mock_parse):
|
||||||
|
with patch.object(s, "sleep_between_requests", new=AsyncMock()):
|
||||||
|
result = await s.fetch_newbuildings(pages=10, delay_override_sec=0)
|
||||||
|
|
||||||
|
assert len(result) == 1
|
||||||
|
assert result[0].source_id == "SAME"
|
||||||
|
|
||||||
|
|
||||||
|
def test_fetch_city_wide_url_unchanged() -> None:
|
||||||
|
"""fetch_city_wide остаётся citywide-URL — рефактор не сменил его поведение."""
|
||||||
|
s = AvitoScraper()
|
||||||
|
assert "novostroyka-" not in s._build_citywide_url(page=1)
|
||||||
|
assert "prodam-ASgBAgICAUSSA8YQ" in s._build_citywide_url(page=1)
|
||||||
|
|
||||||
|
|
||||||
|
# ── 4. NewbuildingSweepCounters ─────────────────────────────────────────────
|
||||||
|
|
||||||
|
|
||||||
|
def test_newbuilding_counters_defaults_and_to_dict() -> None:
|
||||||
|
from app.services.scrape_pipeline import NewbuildingSweepCounters
|
||||||
|
|
||||||
|
c = NewbuildingSweepCounters()
|
||||||
|
assert c.lots_fetched == 0
|
||||||
|
assert c.lots_inserted == 0
|
||||||
|
assert c.lots_updated == 0
|
||||||
|
assert c.errors_count == 0
|
||||||
|
d = c.to_dict()
|
||||||
|
assert set(d.keys()) == {"lots_fetched", "lots_inserted", "lots_updated", "errors_count"}
|
||||||
|
|
||||||
|
|
||||||
|
# ── 5. run_avito_newbuilding_sweep wiring ───────────────────────────────────
|
||||||
|
|
||||||
|
|
||||||
|
class _FakeRuns:
|
||||||
|
def __init__(self) -> None:
|
||||||
|
self.heartbeats: list[dict[str, int]] = []
|
||||||
|
self.done: dict[str, int] | None = None
|
||||||
|
self.failed: tuple[str, dict[str, int]] | None = None
|
||||||
|
self.banned: tuple[str, dict[str, int]] | None = None
|
||||||
|
|
||||||
|
def is_cancelled(self, _db: Any, _run_id: int) -> bool:
|
||||||
|
return False
|
||||||
|
|
||||||
|
def update_heartbeat(self, _db: Any, _run_id: int, counters: dict[str, int]) -> None:
|
||||||
|
self.heartbeats.append(dict(counters))
|
||||||
|
|
||||||
|
def mark_done(self, _db: Any, _run_id: int, counters: dict[str, int]) -> None:
|
||||||
|
self.done = dict(counters)
|
||||||
|
|
||||||
|
def mark_failed(self, _db: Any, _run_id: int, error: str, counters: dict[str, int]) -> None:
|
||||||
|
self.failed = (error, dict(counters))
|
||||||
|
|
||||||
|
def mark_banned(self, _db: Any, _run_id: int, error: str, counters: dict[str, int]) -> None:
|
||||||
|
self.banned = (error, dict(counters))
|
||||||
|
|
||||||
|
|
||||||
|
def _install_runs(monkeypatch: pytest.MonkeyPatch, fake: _FakeRuns) -> None:
|
||||||
|
monkeypatch.setattr(scrape_pipeline.scrape_runs, "is_cancelled", fake.is_cancelled)
|
||||||
|
monkeypatch.setattr(scrape_pipeline.scrape_runs, "update_heartbeat", fake.update_heartbeat)
|
||||||
|
monkeypatch.setattr(scrape_pipeline.scrape_runs, "mark_done", fake.mark_done)
|
||||||
|
monkeypatch.setattr(scrape_pipeline.scrape_runs, "mark_failed", fake.mark_failed)
|
||||||
|
monkeypatch.setattr(scrape_pipeline.scrape_runs, "mark_banned", fake.mark_banned)
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.asyncio
|
||||||
|
async def test_run_avito_newbuilding_sweep_saves_and_counts(
|
||||||
|
monkeypatch: pytest.MonkeyPatch,
|
||||||
|
) -> None:
|
||||||
|
fake_runs = _FakeRuns()
|
||||||
|
_install_runs(monkeypatch, fake_runs)
|
||||||
|
|
||||||
|
fake_lots = [_make_lot("N1"), _make_lot("N2"), _make_lot("N3")]
|
||||||
|
|
||||||
|
async def fake_fetch_newbuildings(
|
||||||
|
self: AvitoScraper,
|
||||||
|
pages: int = 30,
|
||||||
|
*,
|
||||||
|
delay_override_sec: float | None = None,
|
||||||
|
) -> list[ScrapedLot]:
|
||||||
|
return list(fake_lots)
|
||||||
|
|
||||||
|
monkeypatch.setattr(AvitoScraper, "fetch_newbuildings", fake_fetch_newbuildings)
|
||||||
|
|
||||||
|
save_calls: list[int] = []
|
||||||
|
|
||||||
|
def fake_save_listings(_db: Any, lots: Any, **_kw: Any) -> tuple[int, int]:
|
||||||
|
save_calls.append(len(lots))
|
||||||
|
return (len(lots), 0)
|
||||||
|
|
||||||
|
monkeypatch.setattr(scrape_pipeline, "save_listings", fake_save_listings)
|
||||||
|
# Force cffi-mode (no browser) so we don't spin up BrowserFetcher.
|
||||||
|
monkeypatch.setattr(scrape_pipeline.settings, "scraper_fetch_mode", "cffi")
|
||||||
|
|
||||||
|
# Stub shared AsyncSession (used by the sweep in cffi-mode) — no real session.
|
||||||
|
fake_session = AsyncMock()
|
||||||
|
fake_session.__aenter__ = AsyncMock(return_value=fake_session)
|
||||||
|
fake_session.__aexit__ = AsyncMock(return_value=None)
|
||||||
|
monkeypatch.setattr(scrape_pipeline, "AsyncSession", lambda **_kw: fake_session)
|
||||||
|
|
||||||
|
counters = await scrape_pipeline.run_avito_newbuilding_sweep(
|
||||||
|
MagicMock(),
|
||||||
|
run_id=42,
|
||||||
|
pages=2,
|
||||||
|
request_delay_sec=0.0,
|
||||||
|
)
|
||||||
|
|
||||||
|
assert save_calls == [3]
|
||||||
|
assert counters.lots_fetched == 3
|
||||||
|
assert counters.lots_inserted == 3
|
||||||
|
assert counters.lots_updated == 0
|
||||||
|
assert fake_runs.done is not None
|
||||||
|
assert fake_runs.done["lots_fetched"] == 3
|
||||||
|
assert fake_runs.failed is None
|
||||||
|
assert fake_runs.banned is None
|
||||||
Loading…
Add table
Reference in a new issue