124 lines
4.4 KiB
Python
124 lines
4.4 KiB
Python
# Инициализация Celery-приложения и периодических задач.
|
|
|
|
import logging
|
|
|
|
from celery import Celery
|
|
from celery.signals import worker_process_init, worker_ready, setup_logging as celery_setup_logging
|
|
from redis import Redis
|
|
|
|
from ..core.config import settings
|
|
from ..core.logs import setup_logging
|
|
|
|
logger = logging.getLogger("openlane_scraper.worker.celery_app")
|
|
STARTUP_SYNC_DISPATCH_KEY = "openlane:state:startup_sync_dispatched"
|
|
|
|
|
|
@celery_setup_logging.connect
|
|
def _configure_logging(loglevel=None, **kwargs):
|
|
# Перехватываем логирование Celery и пишем только в stderr (Docker logs).
|
|
level = settings.log_level if settings.log_level else "INFO"
|
|
setup_logging(level, None)
|
|
|
|
|
|
@worker_process_init.connect
|
|
def _on_worker_process_init(**kwargs):
|
|
# Повторно настраиваем логирование в каждом дочернем prefork-процессе,
|
|
# чтобы StreamHandler(stderr) корректно работал после fork.
|
|
level = settings.log_level if settings.log_level else "INFO"
|
|
setup_logging(level, None)
|
|
|
|
|
|
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(
|
|
"openlane_scraper",
|
|
broker=_broker_url(),
|
|
backend=_result_backend(),
|
|
)
|
|
|
|
# Auto-clamp: если hard limit слишком далёк от soft (> soft + 120),
|
|
# ограничиваем, чтобы зависший worker не жил вечно.
|
|
_soft = settings.celery.task_soft_time_limit
|
|
_hard = settings.celery.task_time_limit
|
|
_max_hard = _soft + 120 if _soft else _hard
|
|
if _hard > _max_hard:
|
|
logger.warning(
|
|
"CELERY_TASK_TIME_LIMIT=%d too far from CELERY_TASK_SOFT_TIME_LIMIT=%d; "
|
|
"clamping hard limit to %d",
|
|
_hard, _soft, _max_hard,
|
|
)
|
|
_hard = _max_hard
|
|
|
|
celery_app.conf.update(
|
|
task_serializer="json",
|
|
accept_content=["json"],
|
|
result_serializer="json",
|
|
timezone="UTC",
|
|
enable_utc=True,
|
|
task_soft_time_limit=_soft,
|
|
task_time_limit=_hard,
|
|
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_redirect_stdouts=False,
|
|
worker_hijack_root_logger=False,
|
|
beat_schedule={
|
|
"periodic-sync-listing": {
|
|
"task": "openlane_scraper.worker.tasks.sync_listing_task",
|
|
"schedule": settings.celery.beat_sync_interval_minutes * 60.0,
|
|
"args": (),
|
|
"kwargs": {"limit": settings.celery.beat_sync_limit, "only_new": False},
|
|
"options": {"queue": "scraping"},
|
|
}
|
|
},
|
|
task_routes={
|
|
"openlane_scraper.worker.tasks.*": {"queue": "scraping"}
|
|
},
|
|
)
|
|
|
|
celery_app.autodiscover_tasks(["openlane_scraper.worker"])
|
|
|
|
|
|
@worker_ready.connect
|
|
def _on_worker_ready(**kwargs):
|
|
"""Сразу при старте worker отправляем первую задачу sync_listing,
|
|
чтобы не ждать час до первого beat-цикла."""
|
|
try:
|
|
redis_client = Redis.from_url(
|
|
settings.redis.url,
|
|
decode_responses=True,
|
|
socket_connect_timeout=settings.redis.socket_connect_timeout_seconds,
|
|
socket_timeout=settings.redis.socket_timeout_seconds,
|
|
health_check_interval=settings.redis.health_check_interval_seconds,
|
|
retry_on_timeout=True,
|
|
)
|
|
should_dispatch = bool(redis_client.set(STARTUP_SYNC_DISPATCH_KEY, "1", nx=True, ex=600))
|
|
except Exception:
|
|
logger.warning("Worker ready startup sync dedupe check failed; skipping immediate dispatch", exc_info=True)
|
|
return
|
|
|
|
if not should_dispatch:
|
|
logger.info("Worker ready immediate sync already dispatched recently; skipping duplicate enqueue")
|
|
return
|
|
|
|
logger.info("Worker ready — dispatching initial sync_listing task")
|
|
celery_app.send_task(
|
|
"openlane_scraper.worker.tasks.sync_listing_task",
|
|
kwargs={"limit": settings.celery.beat_sync_limit, "only_new": False},
|
|
queue="scraping",
|
|
) |