improve scraper runtime

This commit is contained in:
qananasikq
2026-04-08 13:30:45 +03:00
parent 21bdf17f35
commit f96e5ca5f8
9 changed files with 207 additions and 21 deletions

View File

@@ -1,5 +1,6 @@
import logging
import time
import uuid
from pathlib import Path
from typing import Any
@@ -9,7 +10,7 @@ from playwright.sync_api import TimeoutError as PlaywrightTimeoutError
from .browser import BrowserFactory, HumanPacer, NetworkCapture
from .core.config import Settings, settings
from .core.logs import setup_logging
from .core.logs import set_trace_id, setup_logging
from .core.retry import retryable
from .core.utils import save_to_json
from .parsing.mapper import CarMapper
@@ -24,8 +25,11 @@ logger = logging.getLogger("iaai_scraper.scraper")
class IAAIScraper:
def __init__(self, runtime_settings: Settings | None = None) -> None:
# Базовые зависимости и сервисы скрапера.
self.settings = runtime_settings or settings
setup_logging(self.settings.log_level, self.settings.log_file)
self.trace_id = uuid.uuid4().hex[:12]
set_trace_id(self.trace_id)
self.playwright = None
self.browser = None
self.context: BrowserContext | None = None
@@ -38,20 +42,51 @@ class IAAIScraper:
# browser lifecycle
def _new_trace_id(self, prefix: str) -> str:
# Новый trace id для отдельной операции.
trace_id = f"{prefix}-{uuid.uuid4().hex[:8]}"
self.trace_id = trace_id
set_trace_id(trace_id)
return trace_id
def __enter__(self) -> "IAAIScraper":
self.playwright = sync_playwright().start()
self.browser = self.browser_factory.create_browser(self.playwright)
if self.playwright is None:
self.playwright = sync_playwright().start()
if self.browser is None:
self.browser = self.browser_factory.create_browser(self.playwright)
return self
def __exit__(self, exc_type, exc, tb) -> None:
if self.context:
self.context.close()
if self.browser:
self.browser.close()
if self.playwright:
self.playwright.stop()
self.close()
def close(self) -> None:
# Закрываем ресурсы аккуратно, даже если Playwright уже в ошибке.
if self.context is not None:
try:
self.context.close()
except PlaywrightError:
pass
finally:
self.context = None
if self.browser is not None:
try:
self.browser.close()
except PlaywrightError:
pass
finally:
self.browser = None
if self.playwright is not None:
try:
self.playwright.stop()
except PlaywrightError:
pass
finally:
self.playwright = None
def _new_context(self, storage_state: str | None = None) -> BrowserContext:
# Контекст пересоздаётся, чтобы не тянуть старое состояние страницы.
if self.browser is None:
self.__enter__()
if self.context:
self.context.close()
self.context = self.browser_factory.create_context(self.browser, storage_state=storage_state)
@@ -68,6 +103,7 @@ class IAAIScraper:
# listing
def collect_listing(self, make: str | None = None, model: str | None = None):
# Собираем ссылки карточек с публичного листинга.
page = self._get_unauthenticated_page()
try:
@@ -87,8 +123,11 @@ class IAAIScraper:
# scrape
@retryable(max_attempts=3)
@retryable(max_attempts=3, jitter_seconds=0.25)
def open_vehicle_page(self, vehicle_url: str):
# Лёгкое открытие страницы без сетевого дампа и сохранения в БД.
trace_id = self._new_trace_id("open")
started_at = time.perf_counter()
page = self._get_page()
try:
self.pacer.before_vehicle_open()
@@ -107,15 +146,17 @@ class IAAIScraper:
dom_text = ""
parsed = self.vehicle_parser.normalize(vehicle_url, html, dom_text, {"json_responses": []})
return {
"trace_id": trace_id,
"source_url": vehicle_url,
"opened_in_gentle_mode": True,
"fetched_at_epoch": int(time.time()),
"elapsed_seconds": round(time.perf_counter() - started_at, 3),
**parsed,
}
finally:
page.close()
@retryable(max_attempts=3)
@retryable(max_attempts=3, jitter_seconds=0.25)
def scrape_vehicle_detail(self, vehicle_url: str):
"""Открыть страницу, перехватить JSON, вернуть данные."""
page = self._get_page()
@@ -125,6 +166,9 @@ class IAAIScraper:
page.close()
def _scrape_on_page(self, page: Page, vehicle_url: str):
# Полный проход по карточке с network capture и маппингом в DB-модель.
trace_id = self._new_trace_id("scrape")
started_at = time.perf_counter()
capture = NetworkCapture(self.settings)
capture.attach(page)
@@ -150,8 +194,10 @@ class IAAIScraper:
)
result = {
"trace_id": trace_id,
"source_url": vehicle_url,
"fetched_at_epoch": int(time.time()),
"elapsed_seconds": round(time.perf_counter() - started_at, 3),
"network": network_dump,
**parsed,
"db_record": db_record.model_dump(mode="json"),
@@ -165,6 +211,9 @@ class IAAIScraper:
def sync_vehicle(self, vehicle_url: str, lane: str = "iaai"):
"""Scrape + upsert одного авто."""
# Один URL -> один scrape -> один upsert.
trace_id = self._new_trace_id("sync-vehicle")
started_at = time.perf_counter()
self.persistence.create_tables()
run_id = self.persistence.start_sync_run(lane=lane)
ids_fetched = 1
@@ -184,11 +233,13 @@ class IAAIScraper:
images_upserted = int(upsert.get("images_upserted", 0))
status = "success"
return {
"trace_id": trace_id,
"status": status,
"run_id": run_id,
"vehicle_url": vehicle_url,
"db_action": upsert.get("action"),
"images_upserted": images_upserted,
"elapsed_seconds": round(time.perf_counter() - started_at, 3),
"db_record": db_record,
}
except Exception as exc:
@@ -207,11 +258,16 @@ class IAAIScraper:
error_summary=error_summary,
)
def sync_listing(self, make: str | None = None, model: str | None = None, lane: str = "iaai_cars"):
def sync_listing(self, make: str | None = None, model: str | None = None, lane: str = "iaai_cars", limit: int | None = None):
"""Листинг + sync всех найденных машин."""
# Массовая синхронизация с общим run_id и сбором ошибок.
trace_id = self._new_trace_id("sync-listing")
started_at = time.perf_counter()
self.persistence.create_tables()
listing = self.collect_listing(make=make, model=model)
vehicle_urls = list(listing.get("vehicle_urls", []))
if limit is not None:
vehicle_urls = vehicle_urls[:max(0, limit)]
run_id = self.persistence.start_sync_run(lane=lane)
cars_upserted = 0
cars_failed = 0
@@ -225,6 +281,7 @@ class IAAIScraper:
try:
for index, vehicle_url in enumerate(vehicle_urls, start=1):
try:
# Повторно используем страницу, чтобы не создавать лишний overhead.
logger.info("[%d/%d] Scraping %s", index, total, vehicle_url)
scrape_result = self._scrape_on_page(page, vehicle_url)
db_record = scrape_result.get("db_record")
@@ -242,6 +299,7 @@ class IAAIScraper:
record.year or "?", record.price or "N/A", img_count,
)
except Exception as exc:
# После сбоя пересоздаём page, чтобы не остаться в битом состоянии.
cars_failed += 1
failures.append({"vehicle_url": vehicle_url, "error": str(exc)})
logger.error("[%d/%d] Failed %s: %s", index, total, vehicle_url, exc)
@@ -252,7 +310,7 @@ class IAAIScraper:
page = self._get_page()
if index < total:
time.sleep(1.0)
self.pacer.between_vehicles()
finally:
try:
page.close()
@@ -275,18 +333,21 @@ class IAAIScraper:
run_id, cars_upserted, total, cars_failed, images_upserted,
)
return {
"trace_id": trace_id,
"status": status,
"run_id": run_id,
"listing": listing,
"cars_upserted": cars_upserted,
"cars_failed": cars_failed,
"images_upserted": images_upserted,
"elapsed_seconds": round(time.perf_counter() - started_at, 3),
"failures": failures,
}
# scheduler
def run_scheduled(self) -> None:
# Простой бесконечный цикл без внешнего планировщика.
interval = self.settings.scheduler_interval_minutes * 60
logger.info(
"Scheduler started: syncing every %d minutes",
@@ -332,6 +393,7 @@ class IAAIScraper:
def _warm_page(self, page: Page) -> None:
"""Scroll для lazy-load."""
# Небольшой прогрев страницы для ленивых блоков и картинок.
rounds = max(0, self.settings.gentle.warm_scroll_rounds)
pause_seconds = max(0.1, self.settings.gentle.scroll_pause_ms / 1000)
for _ in range(rounds):