92 lines
2.7 KiB
Python
92 lines
2.7 KiB
Python
# Celery-задачи для синхронизации автомобилей и листинга IAAI.
|
|
|
|
import json
|
|
import logging
|
|
|
|
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 одного автомобиля.
|
|
persistence = _get_persistence()
|
|
persistence.create_tables()
|
|
|
|
try:
|
|
with IAAIScraper() as scraper:
|
|
result = scraper.sync_vehicle(vehicle_url, lane=lane)
|
|
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:
|
|
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 всех найденных машин.
|
|
persistence = _get_persistence()
|
|
persistence.create_tables()
|
|
|
|
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"),
|
|
}
|
|
logger.info(
|
|
"sync_listing_task completed: %d upserted, %d failed",
|
|
summary["cars_upserted"], summary["cars_failed"],
|
|
)
|
|
return {"status": "success", **summary}
|
|
|
|
except Exception as exc:
|
|
logger.error("sync_listing_task failed: %s", exc)
|
|
raise self.retry(exc=exc)
|