improve batch sync add postgres upsert fix sync locking improve listing sync speed up scraper clean up project prepare for github update docker setup
222 lines
7.0 KiB
Python
222 lines
7.0 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 _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:
|
|
return PersistenceService(Settings())
|
|
|
|
|
|
def _get_redis() -> Redis:
|
|
settings = Settings()
|
|
return Redis.from_url(settings.redis.url, decode_responses=True)
|
|
|
|
|
|
def _acquire_lock(redis_client: Redis, key: str, owner_token: str, ttl_seconds: int) -> bool:
|
|
try:
|
|
return bool(redis_client.set(key, owner_token, nx=True, ex=ttl_seconds))
|
|
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)
|