initial commit
This commit is contained in:
3
encar_scraper/worker/__init__.py
Normal file
3
encar_scraper/worker/__init__.py
Normal file
@@ -0,0 +1,3 @@
|
||||
from .celery_app import celery_app
|
||||
|
||||
__all__ = ["celery_app"]
|
||||
71
encar_scraper/worker/celery_app.py
Normal file
71
encar_scraper/worker/celery_app.py
Normal file
@@ -0,0 +1,71 @@
|
||||
# Инициализация Celery-приложения и периодических задач для Encar.
|
||||
|
||||
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(
|
||||
"encar_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,
|
||||
task_acks_late=True,
|
||||
task_reject_on_worker_lost=True,
|
||||
task_track_started=True,
|
||||
worker_concurrency=settings.celery.worker_concurrency,
|
||||
worker_max_tasks_per_child=settings.celery.worker_max_tasks_per_child,
|
||||
worker_pool="prefork",
|
||||
worker_prefetch_multiplier=1,
|
||||
broker_connection_retry_on_startup=True,
|
||||
broker_transport_options={
|
||||
"visibility_timeout": settings.celery.broker_visibility_timeout,
|
||||
},
|
||||
result_expires=86400,
|
||||
worker_hijack_root_logger=False,
|
||||
beat_schedule={
|
||||
"periodic-encar-sync": {
|
||||
"task": "encar_scraper.worker.tasks.encar_sync_listing_task",
|
||||
"schedule": settings.celery.encar_beat_interval_minutes * 60.0,
|
||||
"args": (),
|
||||
"kwargs": {
|
||||
"car_type": "all",
|
||||
},
|
||||
"options": {"queue": "encar"},
|
||||
},
|
||||
},
|
||||
task_routes={
|
||||
"encar_scraper.worker.tasks.*": {"queue": "encar"},
|
||||
},
|
||||
)
|
||||
|
||||
celery_app.autodiscover_tasks(["encar_scraper.worker"])
|
||||
|
||||
|
||||
from celery.signals import setup_logging as celery_setup_logging
|
||||
|
||||
|
||||
@celery_setup_logging.connect
|
||||
def _configure_logging(loglevel, logfile, format, colorize, **kwargs):
|
||||
"""Перехватываем настройку логирования Celery, чтобы использовать свой формат."""
|
||||
from ..core.logs import setup_logging
|
||||
import logging
|
||||
level_name = logging.getLevelName(loglevel) if isinstance(loglevel, int) else str(loglevel)
|
||||
setup_logging(level=level_name, log_file=logfile)
|
||||
244
encar_scraper/worker/tasks.py
Normal file
244
encar_scraper/worker/tasks.py
Normal file
@@ -0,0 +1,244 @@
|
||||
# Celery tasks для синхронизации автомобилей Encar.
|
||||
|
||||
import logging
|
||||
from threading import Event, Thread
|
||||
import uuid
|
||||
|
||||
from celery import shared_task
|
||||
from redis import Redis
|
||||
|
||||
from ..core.config import Settings
|
||||
from ..core.runtime_config import load_runtime_config
|
||||
from ..encar import EncarScraper, EncarFilters, expand_allowed_brands
|
||||
from ..storage.db import PersistenceService
|
||||
|
||||
logger = logging.getLogger("encar_scraper.worker.tasks")
|
||||
|
||||
ENCAR_SYNC_LOCK_KEY = "encar:locks:sync_listing"
|
||||
|
||||
_persistence: PersistenceService | None = None
|
||||
_redis: Redis | None = None
|
||||
|
||||
|
||||
def _sync_listing_lock_ttl_seconds() -> int:
|
||||
settings = Settings()
|
||||
return max(settings.celery.task_time_limit + 120, 300)
|
||||
|
||||
|
||||
def _get_persistence() -> PersistenceService:
|
||||
global _persistence
|
||||
if _persistence is None:
|
||||
_persistence = PersistenceService(Settings())
|
||||
return _persistence
|
||||
|
||||
|
||||
def _get_redis() -> Redis:
|
||||
global _redis
|
||||
if _redis is None:
|
||||
settings = Settings()
|
||||
_redis = Redis.from_url(
|
||||
settings.redis.url,
|
||||
decode_responses=True,
|
||||
socket_timeout=10,
|
||||
socket_connect_timeout=5,
|
||||
)
|
||||
return _redis
|
||||
|
||||
|
||||
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
|
||||
|
||||
|
||||
_REFRESH_LOCK_SCRIPT = """
|
||||
if redis.call('GET', KEYS[1]) == ARGV[1] then
|
||||
return redis.call('EXPIRE', KEYS[1], tonumber(ARGV[2]))
|
||||
end
|
||||
return 0
|
||||
"""
|
||||
|
||||
_RELEASE_LOCK_SCRIPT = """
|
||||
if redis.call('GET', KEYS[1]) == ARGV[1] then
|
||||
return redis.call('DEL', KEYS[1])
|
||||
end
|
||||
return 0
|
||||
"""
|
||||
|
||||
|
||||
def _refresh_lock_if_owner(redis_client: Redis, key: str, owner_token: str, ttl_seconds: int) -> bool | None:
|
||||
try:
|
||||
script = redis_client.register_script(_REFRESH_LOCK_SCRIPT)
|
||||
refreshed = script(keys=[key], args=[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:
|
||||
script = redis_client.register_script(_RELEASE_LOCK_SCRIPT)
|
||||
script(keys=[key], args=[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="encar-lock-heartbeat", daemon=True)
|
||||
thread.start()
|
||||
return stop_event, thread
|
||||
|
||||
|
||||
@shared_task(
|
||||
name="encar_scraper.worker.tasks.encar_sync_listing_task",
|
||||
bind=True,
|
||||
max_retries=2,
|
||||
default_retry_delay=60,
|
||||
acks_late=True,
|
||||
)
|
||||
def encar_sync_listing_task(
|
||||
self,
|
||||
car_type: str = "all",
|
||||
manufacturer: str | None = None,
|
||||
year_from: int | None = None,
|
||||
year_to: int | None = None,
|
||||
price_max: int | None = None,
|
||||
lane: str = "encar",
|
||||
limit: int | None = None,
|
||||
):
|
||||
"""
|
||||
Полная синхронизация листинга Encar.
|
||||
Собирает ВСЕ авто, обновляет цены/данные, помечает проданные.
|
||||
"""
|
||||
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_ttl = _sync_listing_lock_ttl_seconds()
|
||||
lock_acquired = _acquire_lock(redis_client, ENCAR_SYNC_LOCK_KEY, owner_token, lock_ttl)
|
||||
|
||||
if not lock_acquired:
|
||||
logger.info("encar_sync_listing_task skipped: another sync is already running")
|
||||
return {
|
||||
"status": "skipped",
|
||||
"reason": "sync_already_running",
|
||||
"task_id": task_id,
|
||||
}
|
||||
|
||||
heartbeat_stop: Event | None = None
|
||||
heartbeat_thread: Thread | None = None
|
||||
|
||||
try:
|
||||
heartbeat_stop, heartbeat_thread = _start_lock_heartbeat(
|
||||
redis_client,
|
||||
ENCAR_SYNC_LOCK_KEY,
|
||||
owner_token,
|
||||
lock_ttl,
|
||||
)
|
||||
self.update_state(state="STARTED", meta={"stage": "encar_sync_started", "task_id": task_id})
|
||||
|
||||
# --- Загрузка runtime_config ---
|
||||
rc = load_runtime_config()
|
||||
effective_limit = limit or rc.sync.limit
|
||||
effective_lane = lane or rc.sync.lane or "encar"
|
||||
|
||||
car_type_map = {"all": "A", "domestic": "Y", "import": "N"}
|
||||
filters = EncarFilters(
|
||||
car_type=car_type_map.get(car_type, "A"),
|
||||
manufacturer=manufacturer,
|
||||
year_from=year_from,
|
||||
year_to=year_to,
|
||||
price_to=price_max,
|
||||
)
|
||||
|
||||
# Набор разрешённых брендов (lowercased) из runtime_config
|
||||
# Расширяем английские имена alias'ами (корейский + русский перевод)
|
||||
allowed_brands = expand_allowed_brands(
|
||||
set(rc.filters.brands) if rc.filters.brands else None
|
||||
)
|
||||
excluded_brands = expand_allowed_brands(
|
||||
set(rc.filters.exclude_brands) if rc.filters.exclude_brands else None
|
||||
)
|
||||
|
||||
scraper = EncarScraper()
|
||||
result = scraper.sync_listing(
|
||||
limit=effective_limit,
|
||||
filters=filters,
|
||||
lane=effective_lane,
|
||||
redis_client=redis_client,
|
||||
allowed_brands=allowed_brands,
|
||||
excluded_brands=excluded_brands,
|
||||
probe_all_photos=rc.sync.probe_all_photos,
|
||||
)
|
||||
|
||||
summary = {
|
||||
"task_id": task_id,
|
||||
"source": "encar",
|
||||
"total_available": result.get("total_available", 0),
|
||||
"items_collected": result.get("items_collected", 0),
|
||||
"cars_synced": result.get("cars_synced", 0),
|
||||
"cars_failed": result.get("cars_failed", 0),
|
||||
"cars_marked_sold": result.get("cars_marked_sold", 0),
|
||||
}
|
||||
logger.info(
|
||||
"encar_sync_listing_task completed: %d synced, %d failed, %d marked sold",
|
||||
summary["cars_synced"], summary["cars_failed"], summary["cars_marked_sold"],
|
||||
)
|
||||
return {"status": "success", **summary}
|
||||
|
||||
except Exception as exc:
|
||||
logger.error("encar_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, ENCAR_SYNC_LOCK_KEY, owner_token)
|
||||
|
||||
|
||||
@shared_task(
|
||||
name="encar_scraper.worker.tasks.encar_sync_vehicle_task",
|
||||
bind=True,
|
||||
max_retries=2,
|
||||
default_retry_delay=30,
|
||||
acks_late=True,
|
||||
)
|
||||
def encar_sync_vehicle_task(self, vehicle_url: str, lane: str = "encar"):
|
||||
"""Синхронизация одного авто Encar."""
|
||||
persistence = _get_persistence()
|
||||
persistence.create_tables()
|
||||
|
||||
try:
|
||||
scraper = EncarScraper()
|
||||
result = scraper.sync_vehicle(vehicle_url, lane=lane)
|
||||
logger.info("encar_sync_vehicle_task completed: %s", vehicle_url)
|
||||
return {
|
||||
"status": "success",
|
||||
"vehicle_url": vehicle_url,
|
||||
"origin_id": result.get("origin_id"),
|
||||
}
|
||||
except Exception as exc:
|
||||
logger.error("encar_sync_vehicle_task failed: %s — %s", vehicle_url, exc, exc_info=True)
|
||||
raise self.retry(exc=exc)
|
||||
Reference in New Issue
Block a user