# Инициализация Celery-приложения и периодических задач для Encar. import logging from celery import Celery from celery.signals import beat_init from redis import Redis from sqlalchemy import func, select from ..core.config import settings logger = logging.getLogger("encar_scraper.worker.celery_app") _INITIAL_SYNC_BOOTSTRAP_KEY = "encar:bootstrap:initial_sync_enqueued" 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, # Следующий запуск планирует сама задача в finally через apply_async(countdown=...), # чтобы час отсчитывался от завершения прогона, а не от старта. beat_schedule={}, task_routes={ "encar_scraper.worker.tasks.*": {"queue": "encar"}, }, ) celery_app.autodiscover_tasks(["encar_scraper.worker"]) @beat_init.connect def _schedule_initial_sync_on_first_start(**kwargs): """Ставит первый sync сразу после первого запуска проекта на пустой БД.""" from ..storage.db import PersistenceService from ..storage.models import SyncRun from .tasks import encar_sync_listing_task persistence = PersistenceService(settings) redis_client: Redis | None = None try: with persistence.session_scope() as session: total_runs = session.execute(select(func.count(SyncRun.id))).scalar() or 0 if total_runs: return redis_client = Redis.from_url( _broker_url(), decode_responses=True, socket_timeout=10, socket_connect_timeout=5, ) marked = redis_client.set(_INITIAL_SYNC_BOOTSTRAP_KEY, "1", nx=True, ex=600) if not marked: logger.info("Initial Encar sync already bootstrapped on this startup") return encar_sync_listing_task.apply_async(kwargs={"car_type": "all"}, queue="encar") logger.info("Scheduled initial Encar sync because sync_runs table is empty") except Exception: logger.warning("Failed to schedule initial Encar sync", exc_info=True) finally: persistence.engine.dispose() if redis_client is not None: try: redis_client.close() except Exception: pass 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)