269 lines
8.6 KiB
Python
269 lines
8.6 KiB
Python
# Задачи Celery для синхронизации автомобилей и листинга IAAI.
|
|
|
|
from concurrent.futures import ThreadPoolExecutor
|
|
import logging
|
|
from threading import Event, Thread
|
|
import time
|
|
import uuid
|
|
|
|
from celery import shared_task
|
|
from redis import Redis
|
|
|
|
from ..core.config import Settings
|
|
from ..scraper import IAAIScraper
|
|
from ..storage.db import PersistenceService
|
|
|
|
logger = logging.getLogger("iaai_scraper.worker.tasks")
|
|
|
|
SYNC_LISTING_LOCK_KEY = "iaai:locks:sync_listing"
|
|
|
|
|
|
def _retry_with_backoff(func, *, attempts: int = 5, base_delay_s: float = 1.0):
|
|
last_exc: Exception | None = None
|
|
for attempt in range(1, attempts + 1):
|
|
try:
|
|
return func()
|
|
except Exception as exc:
|
|
last_exc = exc
|
|
if attempt >= attempts:
|
|
break
|
|
delay = base_delay_s * (2 ** (attempt - 1))
|
|
logger.warning(
|
|
"Operation failed (attempt %d/%d): %s. Retrying in %.1fs",
|
|
attempt,
|
|
attempts,
|
|
exc,
|
|
delay,
|
|
)
|
|
time.sleep(delay)
|
|
if last_exc is not None:
|
|
raise last_exc
|
|
|
|
|
|
def _run_browser_job(func, *args, **kwargs):
|
|
# Браузерный код запускаем в отдельном потоке без активного loop.
|
|
with ThreadPoolExecutor(max_workers=1, thread_name_prefix="iaai-browser") as executor:
|
|
future = executor.submit(func, *args, **kwargs)
|
|
return future.result()
|
|
|
|
|
|
def _sync_listing_lock_ttl_seconds() -> int:
|
|
settings = Settings()
|
|
# Небольшой запас к лимиту времени, чтобы lock снимался после сбоев.
|
|
return max(settings.celery.task_time_limit + 120, 300)
|
|
|
|
|
|
def _get_persistence() -> PersistenceService:
|
|
settings = Settings()
|
|
persistence = PersistenceService(settings)
|
|
|
|
def _ping_db() -> None:
|
|
with persistence.engine.connect() as conn:
|
|
conn.exec_driver_sql("SELECT 1")
|
|
|
|
_retry_with_backoff(_ping_db, attempts=5, base_delay_s=1.0)
|
|
return persistence
|
|
|
|
|
|
def _get_redis() -> Redis:
|
|
settings = Settings()
|
|
redis_client = Redis.from_url(settings.redis.url, decode_responses=True)
|
|
|
|
def _ping_redis() -> None:
|
|
redis_client.ping()
|
|
|
|
_retry_with_backoff(_ping_redis, attempts=5, base_delay_s=1.0)
|
|
return redis_client
|
|
|
|
|
|
def _acquire_lock(redis_client: Redis, key: str, owner_token: str, ttl_seconds: int) -> bool:
|
|
try:
|
|
acquired = bool(redis_client.set(key, owner_token, nx=True, ex=ttl_seconds))
|
|
if acquired:
|
|
return True
|
|
|
|
# Автовосстановление: если lock завис без TTL, считаем stale и пересоздаём.
|
|
ttl = redis_client.ttl(key)
|
|
if ttl is not None and ttl < 0:
|
|
logger.warning("Detected stale lock without TTL, removing: %s", key)
|
|
redis_client.delete(key)
|
|
return bool(redis_client.set(key, owner_token, nx=True, ex=ttl_seconds))
|
|
|
|
return False
|
|
except Exception as exc:
|
|
logger.warning("Failed to acquire lock %s", key, exc_info=True)
|
|
return False
|
|
|
|
|
|
def _refresh_lock_if_owner(redis_client: Redis, key: str, owner_token: str, ttl_seconds: int) -> bool | None:
|
|
try:
|
|
refreshed = redis_client.eval(
|
|
"""
|
|
if redis.call('GET', KEYS[1]) == ARGV[1] then
|
|
return redis.call('EXPIRE', KEYS[1], tonumber(ARGV[2]))
|
|
end
|
|
return 0
|
|
""",
|
|
1,
|
|
key,
|
|
owner_token,
|
|
int(ttl_seconds),
|
|
)
|
|
return bool(refreshed)
|
|
except Exception as exc:
|
|
logger.warning("Failed to refresh lock %s", key, exc_info=True)
|
|
return None
|
|
|
|
|
|
def _release_lock_if_owner(redis_client: Redis, key: str, owner_token: str) -> None:
|
|
try:
|
|
redis_client.eval(
|
|
"""
|
|
if redis.call('GET', KEYS[1]) == ARGV[1] then
|
|
return redis.call('DEL', KEYS[1])
|
|
end
|
|
return 0
|
|
""",
|
|
1,
|
|
key,
|
|
owner_token,
|
|
)
|
|
except Exception as exc:
|
|
logger.warning("Failed to release lock %s", key, exc_info=True)
|
|
|
|
|
|
def _start_lock_heartbeat(
|
|
redis_client: Redis,
|
|
key: str,
|
|
owner_token: str,
|
|
ttl_seconds: int,
|
|
) -> tuple[Event, Thread]:
|
|
stop_event = Event()
|
|
interval_seconds = max(5.0, min(30.0, ttl_seconds / 3))
|
|
|
|
def _heartbeat() -> None:
|
|
while not stop_event.wait(interval_seconds):
|
|
refreshed = _refresh_lock_if_owner(redis_client, key, owner_token, ttl_seconds)
|
|
if refreshed is False:
|
|
logger.warning("Lost sync_listing lock ownership for %s", owner_token)
|
|
return
|
|
|
|
thread = Thread(target=_heartbeat, name="sync-listing-lock-heartbeat", daemon=True)
|
|
thread.start()
|
|
return stop_event, thread
|
|
|
|
|
|
@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:
|
|
def _job():
|
|
with IAAIScraper() as scraper:
|
|
return scraper.sync_vehicle(vehicle_url, lane=lane)
|
|
|
|
result = _run_browser_job(_job)
|
|
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, exc_info=True)
|
|
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()
|
|
task_id = self.request.id or "unknown"
|
|
owner_token = f"{task_id}:{uuid.uuid4().hex}"
|
|
redis_client = _get_redis()
|
|
|
|
lock_acquired = False
|
|
lock_ttl = _sync_listing_lock_ttl_seconds()
|
|
heartbeat_stop: Event | None = None
|
|
heartbeat_thread: Thread | None = None
|
|
|
|
lock_acquired = _acquire_lock(redis_client, SYNC_LISTING_LOCK_KEY, owner_token, lock_ttl)
|
|
|
|
if not lock_acquired:
|
|
logger.info("sync_listing_task skipped: another sync is already running")
|
|
return {
|
|
"status": "skipped",
|
|
"reason": "sync_already_running",
|
|
"task_id": task_id,
|
|
}
|
|
|
|
try:
|
|
heartbeat_stop, heartbeat_thread = _start_lock_heartbeat(
|
|
redis_client,
|
|
SYNC_LISTING_LOCK_KEY,
|
|
owner_token,
|
|
lock_ttl,
|
|
)
|
|
self.update_state(state="STARTED", meta={"stage": "sync_listing_started", "task_id": task_id})
|
|
def _job():
|
|
with IAAIScraper() as scraper:
|
|
return scraper.sync_listing(
|
|
make=make,
|
|
model=model,
|
|
lane=lane,
|
|
limit=limit,
|
|
only_new=only_new,
|
|
)
|
|
|
|
result = _run_browser_job(_job)
|
|
|
|
summary = {
|
|
"task_id": task_id,
|
|
"run_id": result.get("run_id"),
|
|
"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, exc_info=True)
|
|
raise self.retry(exc=exc)
|
|
finally:
|
|
if heartbeat_stop is not None:
|
|
heartbeat_stop.set()
|
|
if heartbeat_thread is not None:
|
|
heartbeat_thread.join(timeout=max(1.0, min(5.0, lock_ttl / 10)))
|
|
if lock_acquired:
|
|
_release_lock_if_owner(redis_client, SYNC_LISTING_LOCK_KEY, owner_token)
|