IAAI scraper: Playwright + SQLAlchemy, парсинг авто с аукциона, Docker-ready
This commit is contained in:
342
iaai_scraper/scraper.py
Normal file
342
iaai_scraper/scraper.py
Normal file
@@ -0,0 +1,342 @@
|
||||
import logging
|
||||
import time
|
||||
from pathlib import Path
|
||||
from typing import Any
|
||||
|
||||
from playwright.sync_api import Error as PlaywrightError
|
||||
from playwright.sync_api import BrowserContext, Page, sync_playwright
|
||||
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.retry import retryable
|
||||
from .core.utils import save_to_json
|
||||
from .parsing.mapper import CarMapper
|
||||
from .parsing.parser import VehicleParser
|
||||
from .storage.db import PersistenceService
|
||||
from .storage.listing import ListingCollector
|
||||
from .storage.schemas import CarRecord
|
||||
|
||||
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.playwright = None
|
||||
self.browser = None
|
||||
self.context: BrowserContext | None = None
|
||||
self.browser_factory = BrowserFactory(self.settings)
|
||||
self.pacer = HumanPacer(self.settings)
|
||||
self.listing_collector = ListingCollector(self.settings, self.pacer)
|
||||
self.vehicle_parser = VehicleParser()
|
||||
self.car_mapper = CarMapper()
|
||||
self.persistence = PersistenceService(self.settings)
|
||||
|
||||
# browser lifecycle
|
||||
|
||||
def __enter__(self) -> "IAAIScraper":
|
||||
self.playwright = sync_playwright().start()
|
||||
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()
|
||||
|
||||
def _new_context(self, storage_state: str | None = None) -> BrowserContext:
|
||||
if self.context:
|
||||
self.context.close()
|
||||
self.context = self.browser_factory.create_context(self.browser, storage_state=storage_state)
|
||||
return self.context
|
||||
|
||||
def init_db(self):
|
||||
self.persistence.create_tables()
|
||||
return {"status": "ok", "database_url": self.settings.database.url}
|
||||
|
||||
def _get_unauthenticated_page(self) -> Page:
|
||||
context = self._new_context()
|
||||
return context.new_page()
|
||||
|
||||
# listing
|
||||
|
||||
def collect_listing(self, make: str | None = None, model: str | None = None):
|
||||
page = self._get_unauthenticated_page()
|
||||
|
||||
try:
|
||||
listing = self.listing_collector.collect_listing_links(page, make=make, model=model)
|
||||
finally:
|
||||
page.close()
|
||||
|
||||
return {
|
||||
"status": "ok",
|
||||
**listing,
|
||||
}
|
||||
|
||||
def _get_page(self) -> Page:
|
||||
if not self.context:
|
||||
self._new_context()
|
||||
return self.context.new_page()
|
||||
|
||||
# scrape
|
||||
|
||||
@retryable(max_attempts=3)
|
||||
def open_vehicle_page(self, vehicle_url: str):
|
||||
page = self._get_page()
|
||||
try:
|
||||
self.pacer.before_vehicle_open()
|
||||
page.goto(vehicle_url, wait_until="domcontentloaded", timeout=60_000)
|
||||
try:
|
||||
page.wait_for_load_state("networkidle", timeout=15000)
|
||||
except PlaywrightTimeoutError:
|
||||
pass
|
||||
self._warm_page(page)
|
||||
self.pacer.after_vehicle_open()
|
||||
|
||||
html = page.content()
|
||||
try:
|
||||
dom_text = page.locator("body").inner_text(timeout=10_000)
|
||||
except Exception:
|
||||
dom_text = ""
|
||||
parsed = self.vehicle_parser.normalize(vehicle_url, html, dom_text, {"json_responses": []})
|
||||
return {
|
||||
"source_url": vehicle_url,
|
||||
"opened_in_gentle_mode": True,
|
||||
"fetched_at_epoch": int(time.time()),
|
||||
**parsed,
|
||||
}
|
||||
finally:
|
||||
page.close()
|
||||
|
||||
@retryable(max_attempts=3)
|
||||
def scrape_vehicle_detail(self, vehicle_url: str):
|
||||
"""Открыть страницу, перехватить JSON, вернуть данные."""
|
||||
page = self._get_page()
|
||||
try:
|
||||
return self._scrape_on_page(page, vehicle_url)
|
||||
finally:
|
||||
page.close()
|
||||
|
||||
def _scrape_on_page(self, page: Page, vehicle_url: str):
|
||||
capture = NetworkCapture(self.settings)
|
||||
capture.attach(page)
|
||||
|
||||
page.goto(vehicle_url, wait_until="domcontentloaded", timeout=60_000)
|
||||
try:
|
||||
page.wait_for_load_state("networkidle", timeout=10000)
|
||||
except PlaywrightTimeoutError:
|
||||
pass
|
||||
time.sleep(self.settings.network_settle_ms / 1000)
|
||||
|
||||
html = page.content()
|
||||
try:
|
||||
dom_text = page.locator("body").inner_text(timeout=10_000)
|
||||
except Exception:
|
||||
dom_text = ""
|
||||
network_dump = capture.export()
|
||||
parsed = self.vehicle_parser.normalize(vehicle_url, html, dom_text, network_dump)
|
||||
|
||||
db_record = self.car_mapper.map_to_car_record(
|
||||
vehicle_url=vehicle_url,
|
||||
vehicle_summary=parsed.get("vehicle_summary", {}),
|
||||
payload_insights=parsed.get("payload_insights", {}),
|
||||
)
|
||||
|
||||
result = {
|
||||
"source_url": vehicle_url,
|
||||
"fetched_at_epoch": int(time.time()),
|
||||
"network": network_dump,
|
||||
**parsed,
|
||||
"db_record": db_record.model_dump(mode="json"),
|
||||
}
|
||||
|
||||
if self.settings.raw_output_json:
|
||||
save_to_json(network_dump, self.settings.raw_output_json)
|
||||
return result
|
||||
|
||||
# db sync
|
||||
|
||||
def sync_vehicle(self, vehicle_url: str, lane: str = "iaai"):
|
||||
"""Scrape + upsert одного авто."""
|
||||
self.persistence.create_tables()
|
||||
run_id = self.persistence.start_sync_run(lane=lane)
|
||||
ids_fetched = 1
|
||||
cars_upserted = 0
|
||||
cars_failed = 0
|
||||
images_upserted = 0
|
||||
status = "failed"
|
||||
error_summary = None
|
||||
try:
|
||||
scrape_result = self.scrape_vehicle_detail(vehicle_url)
|
||||
db_record = scrape_result.get("db_record")
|
||||
if not db_record:
|
||||
raise RuntimeError("Scrape result does not contain db_record")
|
||||
record = CarRecord.model_validate(db_record)
|
||||
upsert = self.persistence.upsert_car(record)
|
||||
cars_upserted = 1
|
||||
images_upserted = int(upsert.get("images_upserted", 0))
|
||||
status = "success"
|
||||
return {
|
||||
"status": status,
|
||||
"run_id": run_id,
|
||||
"vehicle_url": vehicle_url,
|
||||
"db_action": upsert.get("action"),
|
||||
"images_upserted": images_upserted,
|
||||
"db_record": db_record,
|
||||
}
|
||||
except Exception as exc:
|
||||
cars_failed = 1
|
||||
status = "failed"
|
||||
error_summary = str(exc)
|
||||
raise
|
||||
finally:
|
||||
self.persistence.finish_sync_run(
|
||||
run_id,
|
||||
status=status,
|
||||
ids_fetched=ids_fetched,
|
||||
cars_upserted=cars_upserted,
|
||||
cars_failed=cars_failed,
|
||||
images_upserted=images_upserted,
|
||||
error_summary=error_summary,
|
||||
)
|
||||
|
||||
def sync_listing(self, make: str | None = None, model: str | None = None, lane: str = "iaai_cars"):
|
||||
"""Листинг + sync всех найденных машин."""
|
||||
self.persistence.create_tables()
|
||||
listing = self.collect_listing(make=make, model=model)
|
||||
vehicle_urls = list(listing.get("vehicle_urls", []))
|
||||
run_id = self.persistence.start_sync_run(lane=lane)
|
||||
cars_upserted = 0
|
||||
cars_failed = 0
|
||||
images_upserted = 0
|
||||
failures: list[dict[str, str]] = []
|
||||
|
||||
total = len(vehicle_urls)
|
||||
logger.info("Starting sync: %d vehicles to process", total)
|
||||
|
||||
page = self._get_page()
|
||||
try:
|
||||
for index, vehicle_url in enumerate(vehicle_urls, start=1):
|
||||
try:
|
||||
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")
|
||||
if not db_record:
|
||||
raise RuntimeError("Scrape result does not contain db_record")
|
||||
record = CarRecord.model_validate(db_record)
|
||||
upsert = self.persistence.upsert_car(record)
|
||||
cars_upserted += 1
|
||||
img_count = int(upsert.get("images_upserted", 0))
|
||||
images_upserted += img_count
|
||||
logger.info(
|
||||
"[%d/%d] %s %s: %s %s %s — %s, %d images",
|
||||
index, total, upsert.get("action", "?"),
|
||||
record.origin_id, record.brand, record.model,
|
||||
record.year or "?", record.price or "N/A", img_count,
|
||||
)
|
||||
except Exception as exc:
|
||||
cars_failed += 1
|
||||
failures.append({"vehicle_url": vehicle_url, "error": str(exc)})
|
||||
logger.error("[%d/%d] Failed %s: %s", index, total, vehicle_url, exc)
|
||||
try:
|
||||
page.close()
|
||||
except Exception:
|
||||
pass
|
||||
page = self._get_page()
|
||||
|
||||
if index < total:
|
||||
time.sleep(1.0)
|
||||
finally:
|
||||
try:
|
||||
page.close()
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
status = "success" if not failures else ("partial_success" if cars_upserted else "failed")
|
||||
error_summary = "; ".join(item["error"] for item in failures[:10]) if failures else None
|
||||
self.persistence.finish_sync_run(
|
||||
run_id,
|
||||
status=status,
|
||||
ids_fetched=total,
|
||||
cars_upserted=cars_upserted,
|
||||
cars_failed=cars_failed,
|
||||
images_upserted=images_upserted,
|
||||
error_summary=error_summary,
|
||||
)
|
||||
logger.info(
|
||||
"Sync run #%d finished: %d/%d upserted, %d failed, %d images",
|
||||
run_id, cars_upserted, total, cars_failed, images_upserted,
|
||||
)
|
||||
return {
|
||||
"status": status,
|
||||
"run_id": run_id,
|
||||
"listing": listing,
|
||||
"cars_upserted": cars_upserted,
|
||||
"cars_failed": cars_failed,
|
||||
"images_upserted": images_upserted,
|
||||
"failures": failures,
|
||||
}
|
||||
|
||||
# scheduler
|
||||
|
||||
def run_scheduled(self) -> None:
|
||||
interval = self.settings.scheduler_interval_minutes * 60
|
||||
logger.info(
|
||||
"Scheduler started: syncing every %d minutes",
|
||||
self.settings.scheduler_interval_minutes,
|
||||
)
|
||||
cycle = 0
|
||||
while True:
|
||||
cycle += 1
|
||||
logger.info("=== Scheduler cycle #%d starting ===", cycle)
|
||||
start = time.time()
|
||||
try:
|
||||
if self.context:
|
||||
try:
|
||||
self.context.close()
|
||||
except PlaywrightError:
|
||||
pass
|
||||
self.context = None
|
||||
|
||||
result = self.sync_listing()
|
||||
elapsed = time.time() - start
|
||||
logger.info(
|
||||
"Cycle #%d done in %.1fs: %d upserted, %d failed",
|
||||
cycle, elapsed,
|
||||
result.get("cars_upserted", 0),
|
||||
result.get("cars_failed", 0),
|
||||
)
|
||||
except Exception as exc:
|
||||
elapsed = time.time() - start
|
||||
logger.error("Cycle #%d failed after %.1fs: %s", cycle, elapsed, exc)
|
||||
if self.context:
|
||||
try:
|
||||
self.context.close()
|
||||
except PlaywrightError:
|
||||
pass
|
||||
self.context = None
|
||||
|
||||
sleep_time = max(0, interval - (time.time() - start))
|
||||
if sleep_time > 0:
|
||||
logger.info("Sleeping %.0f seconds until next cycle...", sleep_time)
|
||||
time.sleep(sleep_time)
|
||||
|
||||
# helpers
|
||||
|
||||
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):
|
||||
page.mouse.wheel(0, 1600)
|
||||
time.sleep(pause_seconds)
|
||||
if rounds:
|
||||
page.mouse.wheel(0, -3000)
|
||||
time.sleep(min(0.5, pause_seconds))
|
||||
Reference in New Issue
Block a user