Files
encar/encar_scraper/worker/celery_app.py
2026-04-16 18:10:50 +03:00

71 lines
2.2 KiB
Python

# Инициализация 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)