add fastapi rest api and celery workers

This commit is contained in:
qananasikq
2026-04-08 22:54:33 +03:00
parent 6e97991c68
commit f4b9d20191
15 changed files with 982 additions and 0 deletions

View File

@@ -0,0 +1,3 @@
from .app import create_app
__all__ = ["create_app"]

41
iaai_scraper/api/app.py Normal file
View File

@@ -0,0 +1,41 @@
"""FastAPI app."""
from contextlib import asynccontextmanager
from fastapi import FastAPI
from ..core.config import Settings
from ..storage.db import PersistenceService
from .routes import cars, health, tasks
@asynccontextmanager
async def lifespan(app: FastAPI):
"""Жизненный цикл."""
persistence: PersistenceService = app.state.persistence
persistence.create_tables()
yield
def create_app(settings: Settings | None = None) -> FastAPI:
_settings = settings or Settings()
app = FastAPI(
title="IAAI Scraper API",
description="REST API для управления задачами скрапинга IAAI и просмотра данных",
version="1.0.0",
lifespan=lifespan,
)
app.state.settings = _settings
app.state.persistence = PersistenceService(_settings)
app.include_router(health.router, tags=["health"])
app.include_router(cars.router, prefix="/api/v1", tags=["cars"])
app.include_router(tasks.router, prefix="/api/v1", tags=["tasks"])
return app
# Для запуска через uvicorn.
app = create_app()

9
iaai_scraper/api/deps.py Normal file
View File

@@ -0,0 +1,9 @@
"""FastAPI deps."""
from fastapi import Request
from ..storage.db import PersistenceService
def get_persistence(request: Request) -> PersistenceService:
return request.app.state.persistence

View File

View File

@@ -0,0 +1,144 @@
"""Cars endpoints."""
from fastapi import APIRouter, Depends, HTTPException, Query
from sqlalchemy import func, select
from ..deps import get_persistence
from ...storage.db import PersistenceService
from ...storage.models import Car, Image
router = APIRouter()
@router.get("/cars")
def list_cars(
page: int = Query(1, ge=1),
per_page: int = Query(20, ge=1, le=100),
brand: str | None = None,
model: str | None = None,
year_min: int | None = None,
year_max: int | None = None,
is_sold: bool | None = None,
persistence: PersistenceService = Depends(get_persistence),
):
"""Список автомобилей с пагинацией и фильтрами."""
with persistence.session_scope() as session:
query = select(Car)
if brand:
query = query.where(Car.brand.ilike(f"%{brand}%"))
if model:
query = query.where(Car.model.ilike(f"%{model}%"))
if year_min is not None:
query = query.where(Car.year >= year_min)
if year_max is not None:
query = query.where(Car.year <= year_max)
if is_sold is not None:
query = query.where(Car.is_sold == is_sold)
# Общее количество.
count_query = select(func.count()).select_from(query.subquery())
total = session.execute(count_query).scalar() or 0
# Пагинация.
offset = (page - 1) * per_page
cars = session.execute(
query.order_by(Car.last_seen_at.desc()).offset(offset).limit(per_page)
).scalars().all()
return {
"total": total,
"page": page,
"per_page": per_page,
"pages": (total + per_page - 1) // per_page if per_page else 0,
"items": [_car_to_dict(car) for car in cars],
}
@router.get("/cars/{car_id}")
def get_car(
car_id: int,
persistence: PersistenceService = Depends(get_persistence),
):
"""Детальная информация об автомобиле с изображениями."""
with persistence.session_scope() as session:
car = session.get(Car, car_id)
if car is None:
raise HTTPException(status_code=404, detail="Car not found")
return _car_to_dict(car, include_images=True)
@router.get("/cars/by-origin/{origin_id}")
def get_car_by_origin(
origin_id: str,
persistence: PersistenceService = Depends(get_persistence),
):
"""Поиск автомобиля по origin_id."""
with persistence.session_scope() as session:
car = session.execute(
select(Car).where(Car.origin_id == origin_id)
).scalars().first()
if car is None:
raise HTTPException(status_code=404, detail="Car not found")
return _car_to_dict(car, include_images=True)
@router.get("/stats")
def get_stats(persistence: PersistenceService = Depends(get_persistence)):
"""Общая статистика по БД."""
with persistence.session_scope() as session:
total_cars = session.execute(select(func.count(Car.id))).scalar() or 0
total_images = session.execute(select(func.count(Image.id))).scalar() or 0
brands = session.execute(
select(Car.brand, func.count(Car.id))
.group_by(Car.brand)
.order_by(func.count(Car.id).desc())
.limit(20)
).all()
return {
"total_cars": total_cars,
"total_images": total_images,
"top_brands": [{"brand": b, "count": c} for b, c in brands],
}
def _car_to_dict(car: Car, include_images: bool = False) -> dict:
"""Сериализация Car в dict."""
result = {
"id": car.id,
"parser_id": car.parser_id,
"brand": car.brand,
"model": car.model,
"year": car.year,
"price": car.price,
"currency": car.currency,
"mileage": car.mileage,
"country": car.country,
"is_sold": car.is_sold,
"color": car.color,
"drive": car.drive,
"gearbox": car.gearbox,
"steering_wheel": car.steering_wheel,
"body_type": car.body_type,
"engine_volume": car.engine_volume,
"selling_type": car.selling_type,
"origin": car.origin,
"origin_url": car.origin_url,
"origin_id": car.origin_id,
"is_damaged": car.is_damaged,
"slug": car.slug,
"last_seen_at": car.last_seen_at.isoformat() if car.last_seen_at else None,
"content_hash": car.content_hash,
}
if include_images:
result["images"] = [
{
"id": img.id,
"fullres_image": img.fullres_image,
"preview_image": img.preview_image,
"order_index": img.order_index,
}
for img in sorted(car.images, key=lambda i: i.order_index)
]
return result

View File

@@ -0,0 +1,27 @@
"""Health endpoint."""
from fastapi import APIRouter, Depends
from sqlalchemy import text
from ..deps import get_persistence
from ...storage.db import PersistenceService
router = APIRouter()
@router.get("/health")
def health_check(persistence: PersistenceService = Depends(get_persistence)):
"""Проверка API и БД."""
db_ok = False
try:
with persistence.session_scope() as session:
session.execute(text("SELECT 1"))
db_ok = True
except Exception:
pass
return {
"status": "ok" if db_ok else "degraded",
"service": "iaai-scraper-api",
"database": "connected" if db_ok else "unavailable",
}

View File

@@ -0,0 +1,174 @@
"""Task endpoints."""
from datetime import datetime, timezone
from fastapi import APIRouter, Depends, HTTPException, Query
from pydantic import BaseModel
from sqlalchemy import select, func
from ..deps import get_persistence
from ...storage.db import PersistenceService
from ...storage.models import ScrapeTask, SyncRun
from ...worker.tasks import sync_vehicle_task, sync_listing_task
router = APIRouter()
# --- Схемы.
class SyncVehicleRequest(BaseModel):
vehicle_url: str
lane: str = "iaai"
class SyncListingRequest(BaseModel):
make: str | None = None
model: str | None = None
lane: str = "iaai_cars"
limit: int | None = None
only_new: bool | None = None
# --- Эндпоинты.
@router.post("/tasks/sync-vehicle")
def start_sync_vehicle(
body: SyncVehicleRequest,
persistence: PersistenceService = Depends(get_persistence),
):
"""Запустить задачу скрапинга одного автомобиля через Celery."""
result = sync_vehicle_task.delay(body.vehicle_url, lane=body.lane)
# Регистрируем задачу в БД.
persistence.create_scrape_task(
celery_task_id=result.id,
task_type="sync_vehicle",
vehicle_url=body.vehicle_url,
)
return {
"task_id": result.id,
"status": "queued",
"vehicle_url": body.vehicle_url,
}
@router.post("/tasks/sync-listing")
def start_sync_listing(
body: SyncListingRequest,
persistence: PersistenceService = Depends(get_persistence),
):
"""Запустить задачу полного цикла листинга через Celery."""
result = sync_listing_task.delay(
make=body.make,
model=body.model,
lane=body.lane,
limit=body.limit,
only_new=body.only_new,
)
persistence.create_scrape_task(
celery_task_id=result.id,
task_type="sync_listing",
)
return {
"task_id": result.id,
"status": "queued",
}
@router.get("/tasks/{task_id}")
def get_task_status(
task_id: str,
persistence: PersistenceService = Depends(get_persistence),
):
"""Статус Celery-задачи."""
task_info = persistence.get_scrape_task(task_id)
if task_info is None:
raise HTTPException(status_code=404, detail="Task not found")
return task_info
@router.get("/tasks")
def list_tasks(
page: int = Query(1, ge=1),
per_page: int = Query(20, ge=1, le=100),
status: str | None = None,
task_type: str | None = None,
persistence: PersistenceService = Depends(get_persistence),
):
"""Список задач с пагинацией."""
with persistence.session_scope() as session:
query = select(ScrapeTask)
if status:
query = query.where(ScrapeTask.status == status)
if task_type:
query = query.where(ScrapeTask.task_type == task_type)
total = session.execute(
select(func.count()).select_from(query.subquery())
).scalar() or 0
offset = (page - 1) * per_page
tasks = session.execute(
query.order_by(ScrapeTask.created_at.desc()).offset(offset).limit(per_page)
).scalars().all()
return {
"total": total,
"page": page,
"per_page": per_page,
"items": [
{
"id": t.id,
"celery_task_id": t.celery_task_id,
"task_type": t.task_type,
"status": t.status,
"vehicle_url": t.vehicle_url,
"created_at": t.created_at.isoformat() if t.created_at else None,
"started_at": t.started_at.isoformat() if t.started_at else None,
"finished_at": t.finished_at.isoformat() if t.finished_at else None,
"result_summary": t.result_summary,
"error_message": t.error_message,
}
for t in tasks
],
}
@router.get("/sync-runs")
def list_sync_runs(
page: int = Query(1, ge=1),
per_page: int = Query(20, ge=1, le=100),
persistence: PersistenceService = Depends(get_persistence),
):
"""История запусков синхронизации."""
with persistence.session_scope() as session:
total = session.execute(select(func.count(SyncRun.id))).scalar() or 0
offset = (page - 1) * per_page
runs = session.execute(
select(SyncRun).order_by(SyncRun.started_at.desc()).offset(offset).limit(per_page)
).scalars().all()
return {
"total": total,
"page": page,
"per_page": per_page,
"items": [
{
"id": r.id,
"started_at": r.started_at.isoformat() if r.started_at else None,
"finished_at": r.finished_at.isoformat() if r.finished_at else None,
"status": r.status,
"lane": r.lane,
"ids_fetched": r.ids_fetched,
"cars_upserted": r.cars_upserted,
"cars_failed": r.cars_failed,
"images_upserted": r.images_upserted,
"error_summary": r.error_summary,
}
for r in runs
],
}

View File

@@ -0,0 +1,157 @@
import logging
import re
from dataclasses import asdict, dataclass, field
from typing import Any
from urllib.parse import urljoin
from playwright.sync_api import Page
from .pace import HumanPacer
from ..core.config import Settings
from ..core.utils import first_non_empty
logger = logging.getLogger("iaai_scraper.listing")
VEHICLE_HREF_RE = re.compile(r"/VehicleDetail/\d+(?:~[A-Z]{2})?", re.IGNORECASE)
@dataclass(slots=True)
class ListingVehicleLink:
# Ссылка на карточку.
href: str
title: str = ""
lot_number: str | None = None
@dataclass(slots=True)
class ListingPageResult:
# Результат страницы листинга.
source_url: str
page_number: int
vehicle_links: list[ListingVehicleLink] = field(default_factory=list)
pagination_available: bool = False
next_page_detected: bool = False
class ListingCollector:
def __init__(self, settings: Settings, pacer: HumanPacer) -> None:
# Сборщик ссылок листинга.
self.settings = settings
self.pacer = pacer
def open_cars_listing(self, page: Page) -> None:
# Открываем листинг и ждем ссылки.
logger.info("Opening cars listing page: %s", self.settings.listing.cars_url)
page.goto(self.settings.listing.cars_url, wait_until="domcontentloaded")
try:
page.wait_for_load_state("networkidle", timeout=15000)
except Exception:
logger.debug("networkidle timeout on listing page, continuing with current state")
try:
page.wait_for_selector("a[href*='/VehicleDetail/']", timeout=20000)
except Exception:
logger.warning("Vehicle links did not appear within timeout; page may not have rendered fully")
self.pacer.after_listing_open()
def apply_filters(self, page: Page, make: str | None = None, model: str | None = None) -> dict[str, str | None]:
# Пробуем make/model фильтры.
applied = {"make": None, "model": None}
if make and self._try_fill_filter_input(page, ["input[placeholder*='Make']", "input[aria-label*='Make']"], make):
applied["make"] = make
self.pacer.after_filter_action()
if model and self._try_fill_filter_input(page, ["input[placeholder*='Model']", "input[aria-label*='Model']"], model):
applied["model"] = model
self.pacer.after_filter_action()
return applied
def collect_current_page(self, page: Page, page_number: int = 1) -> ListingPageResult:
# Собираем ссылки текущей страницы.
anchors = page.locator("a[href*='/VehicleDetail/']")
total = min(anchors.count(), self.settings.listing.page_link_limit)
links: list[ListingVehicleLink] = []
seen: set[str] = set()
for idx in range(total):
anchor = anchors.nth(idx)
href = anchor.get_attribute("href") or ""
match = VEHICLE_HREF_RE.search(href)
if not match:
continue
absolute = urljoin(self.settings.home_url, match.group(0))
if absolute in seen:
continue
seen.add(absolute)
title = first_non_empty([anchor.get_attribute("title"), anchor.text_content(), ""]) or ""
links.append(ListingVehicleLink(href=absolute, title=str(title).strip()))
if len(links) >= self.settings.listing.max_vehicles_per_run:
break
next_page_detected = self._has_next_page(page)
return ListingPageResult(source_url=page.url, page_number=page_number, vehicle_links=links, pagination_available=next_page_detected, next_page_detected=next_page_detected)
def go_to_next_page(self, page: Page) -> bool:
# Переход на следующую страницу.
selectors = ["a[aria-label*='Next']", "button[aria-label*='Next']", "a.pagination-next", "button.pagination-next", "a:has-text('Next')", "button:has-text('Next')"]
for selector in selectors:
locator = page.locator(selector).first
if locator.count() == 0:
continue
disabled = (locator.get_attribute("disabled") or "").lower()
aria_disabled = (locator.get_attribute("aria-disabled") or "").lower()
classes = (locator.get_attribute("class") or "").lower()
if disabled or aria_disabled == "true" or "disabled" in classes:
continue
self.pacer.move_mouse_to(page, locator)
locator.click()
try:
page.wait_for_load_state("networkidle", timeout=15000)
except Exception:
pass
self.pacer.after_page_change()
return True
return False
def collect_listing_links(self, page: Page, *, make: str | None = None, model: str | None = None) -> dict[str, Any]:
# Последовательный сбор с лимитами.
self.open_cars_listing(page)
applied_filters = self.apply_filters(page, make=make, model=model)
pages: list[dict[str, object]] = []
all_links: list[str] = []
for page_number in range(1, max(1, self.settings.listing.max_pages_per_run) + 1):
page_result = self.collect_current_page(page, page_number=page_number)
pages.append({"page_number": page_result.page_number, "source_url": page_result.source_url, "links_found": len(page_result.vehicle_links), "vehicle_links": [asdict(item) for item in page_result.vehicle_links], "next_page_detected": page_result.next_page_detected})
for item in page_result.vehicle_links:
if item.href not in all_links:
all_links.append(item.href)
if len(all_links) >= self.settings.listing.max_vehicles_per_run:
break
if len(all_links) >= self.settings.listing.max_vehicles_per_run or self.settings.listing.collect_current_page_only or not self.settings.listing.include_pagination or not page_result.next_page_detected:
break
if not self.go_to_next_page(page):
break
return {"listing_url": self.settings.listing.cars_url, "applied_filters": applied_filters, "pages_collected": len(pages), "vehicles_collected": len(all_links), "vehicle_urls": all_links, "pages": pages, "strategy": {"sequential": True, "collect_current_page_only": self.settings.listing.collect_current_page_only, "include_pagination": self.settings.listing.include_pagination, "max_pages_per_run": self.settings.listing.max_pages_per_run, "max_vehicles_per_run": self.settings.listing.max_vehicles_per_run}}
@staticmethod
def _try_fill_filter_input(page: Page, selectors: list[str], value: str) -> bool:
# Несколько селекторов для фильтра.
for selector in selectors:
locator = page.locator(selector).first
if locator.count() == 0:
continue
try:
locator.click()
locator.fill(value)
page.keyboard.press("Enter")
try:
page.wait_for_load_state("networkidle", timeout=15000)
except Exception:
pass
return True
except Exception:
continue
return False
@staticmethod
def _has_next_page(page: Page) -> bool:
# Проверка кнопки next.
for selector in ["a[aria-label*='Next']", "button[aria-label*='Next']", "a.pagination-next", "button.pagination-next", "a:has-text('Next')", "button:has-text('Next')"]:
if page.locator(selector).count() > 0:
return True
return False

View File

@@ -0,0 +1,3 @@
from .celery_app import celery_app
__all__ = ["celery_app"]

View File

@@ -0,0 +1,59 @@
"""Celery app factory."""
from celery import Celery
from ..core.config import settings
def _broker_url() -> str:
return settings.celery.broker_url or settings.redis.url
def _result_backend() -> str:
return settings.celery.result_backend or settings.redis.url
celery_app = Celery(
"iaai_scraper",
broker=_broker_url(),
backend=_result_backend(),
)
celery_app.conf.update(
# Сериализация.
task_serializer="json",
accept_content=["json"],
result_serializer="json",
timezone="UTC",
enable_utc=True,
# Лимиты.
task_soft_time_limit=settings.celery.task_soft_time_limit,
task_time_limit=settings.celery.task_time_limit,
worker_concurrency=settings.celery.worker_concurrency,
# Playwright sync API: prefork + 1 процесс.
worker_pool="prefork",
worker_prefetch_multiplier=1,
# Retry policy брокера.
broker_connection_retry_on_startup=True,
# Beat schedule.
beat_schedule={
"periodic-sync-listing": {
"task": "iaai_scraper.worker.tasks.sync_listing_task",
"schedule": settings.celery.beat_sync_interval_minutes * 60.0,
"args": (),
"options": {"queue": "scraping"},
},
},
# Роутинг задач.
task_routes={
"iaai_scraper.worker.tasks.*": {"queue": "scraping"},
},
)
# Автообнаружение задач.
celery_app.autodiscover_tasks(["iaai_scraper.worker"])

View File

@@ -0,0 +1,137 @@
"""Celery tasks для IAAI."""
import json
import logging
from datetime import datetime, timezone
from celery import shared_task
from ..core.config import Settings
from ..scraper import IAAIScraper
from ..storage.db import PersistenceService
logger = logging.getLogger("iaai_scraper.worker.tasks")
def _get_persistence() -> PersistenceService:
return PersistenceService(Settings())
@shared_task(
name="iaai_scraper.worker.tasks.sync_vehicle_task",
bind=True,
max_retries=2,
default_retry_delay=30,
acks_late=True,
)
def sync_vehicle_task(self, vehicle_url: str, lane: str = "iaai"):
"""Скрапинг и upsert одного автомобиля."""
task_id = self.request.id
persistence = _get_persistence()
persistence.create_tables()
# Регистрируем запуск.
persistence.update_scrape_task(
task_id,
status="running",
started_at=datetime.now(timezone.utc),
)
try:
with IAAIScraper() as scraper:
result = scraper.sync_vehicle(vehicle_url, lane=lane)
persistence.update_scrape_task(
task_id,
status="success",
finished_at=datetime.now(timezone.utc),
result_summary=json.dumps({
"car_upserted": True,
"db_action": result.get("db_action"),
"images_upserted": result.get("images_upserted", 0),
"elapsed_seconds": result.get("elapsed_seconds"),
}, default=str),
)
logger.info("sync_vehicle_task completed: %s", vehicle_url)
return {
"status": "success",
"vehicle_url": vehicle_url,
"db_action": result.get("db_action"),
"images_upserted": result.get("images_upserted", 0),
}
except Exception as exc:
persistence.update_scrape_task(
task_id,
status="failed",
finished_at=datetime.now(timezone.utc),
error_message=str(exc)[:2000],
)
logger.error("sync_vehicle_task failed: %s%s", vehicle_url, exc)
raise self.retry(exc=exc)
@shared_task(
name="iaai_scraper.worker.tasks.sync_listing_task",
bind=True,
max_retries=1,
default_retry_delay=120,
acks_late=True,
)
def sync_listing_task(
self,
make: str | None = None,
model: str | None = None,
lane: str = "iaai_cars",
limit: int | None = None,
only_new: bool | None = None,
):
"""Полный цикл: листинг + sync всех найденных машин."""
task_id = self.request.id
persistence = _get_persistence()
persistence.create_tables()
persistence.update_scrape_task(
task_id,
status="running",
started_at=datetime.now(timezone.utc),
)
try:
with IAAIScraper() as scraper:
result = scraper.sync_listing(
make=make,
model=model,
lane=lane,
limit=limit,
only_new=only_new,
)
summary = {
"cars_upserted": result.get("cars_upserted", 0),
"cars_failed": result.get("cars_failed", 0),
"images_upserted": result.get("images_upserted", 0),
"skipped_existing": result.get("skipped_existing", 0),
"elapsed_seconds": result.get("elapsed_seconds"),
}
persistence.update_scrape_task(
task_id,
status="success",
finished_at=datetime.now(timezone.utc),
result_summary=json.dumps(summary, default=str),
)
logger.info(
"sync_listing_task completed: %d upserted, %d failed",
summary["cars_upserted"], summary["cars_failed"],
)
return {"status": "success", **summary}
except Exception as exc:
persistence.update_scrape_task(
task_id,
status="failed",
finished_at=datetime.now(timezone.utc),
error_message=str(exc)[:2000],
)
logger.error("sync_listing_task failed: %s", exc)
raise self.retry(exc=exc)