267 lines
9.1 KiB
Python
267 lines
9.1 KiB
Python
# Celery-приложение и периодические задачи.
|
||
|
||
import json
|
||
import logging
|
||
import os
|
||
import time
|
||
|
||
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
|
||
from .constants import (
|
||
GLOBAL_DB_PROGRESS_TS_KEY,
|
||
GLOBAL_PROGRESS_TS_KEY,
|
||
MOBILEDE_RUNTIME_SEGMENTS_TASK,
|
||
MOBILEDE_SYNC_QUEUE,
|
||
MOBILEDE_SYNC_TASK_NAME,
|
||
)
|
||
|
||
logger = logging.getLogger("mobilede_scraper.worker.celery_app")
|
||
STARTUP_SYNC_DISPATCH_KEY = "mobilede:state:startup_sync_dispatched"
|
||
PROGRESS_KEY_PREFIX = "mobilede:state:task_progress:"
|
||
MOBILEDE_SYNC_DETAIL_TASK = "mobilede.sync_detail"
|
||
MOBILEDE_ENRICH_IMAGES_TASK = "mobilede.enrich_images_batch"
|
||
STARTUP_SYNC_DISPATCH_TTL_SECONDS = 10 * 60
|
||
STARTUP_PROGRESS_MAX_AGE_SECONDS = 180
|
||
STARTUP_GLOBAL_PROGRESS_MAX_AGE_SECONDS = 300
|
||
|
||
|
||
def _env_bool(name: str, default: bool) -> bool:
|
||
raw = os.getenv(name, "true" if default else "false").strip().lower()
|
||
return raw in {"1", "true", "yes", "on"}
|
||
|
||
|
||
MOBILEDE_BEAT_SYNC_ENABLED = _env_bool("MOBILEDE_BEAT_SYNC_ENABLED", True)
|
||
|
||
|
||
def _runtime_sync_kwargs() -> dict[str, bool | float]:
|
||
return {
|
||
"delay_seconds": float(os.getenv("MOBILEDE_REQUEST_DELAY_SECONDS", "0.7")),
|
||
"use_cursor": _env_bool("MOBILEDE_CURSOR_ENABLED", True),
|
||
"continuous": _env_bool("MOBILEDE_CONTINUOUS_SYNC_ENABLED", True),
|
||
}
|
||
|
||
|
||
def _runtime_sync_expires_seconds() -> float:
|
||
return settings.celery.beat_sync_interval_minutes * 60.0
|
||
|
||
|
||
def _has_fresh_active_progress(
|
||
redis_client: Redis,
|
||
*,
|
||
max_age_seconds: int = STARTUP_PROGRESS_MAX_AGE_SECONDS,
|
||
) -> bool:
|
||
now = int(time.time())
|
||
try:
|
||
for raw_key in redis_client.scan_iter(f"{PROGRESS_KEY_PREFIX}*"):
|
||
payload = redis_client.get(raw_key)
|
||
if not payload:
|
||
continue
|
||
try:
|
||
progress = json.loads(payload)
|
||
except (TypeError, ValueError):
|
||
continue
|
||
ts = int(progress.get("ts") or 0)
|
||
stage = str(progress.get("stage") or "")
|
||
if ts > 0 and now - ts <= max_age_seconds and stage not in {"segment_done", "sync_done", "failed"}:
|
||
return True
|
||
except Exception:
|
||
logger.warning("Failed to inspect startup progress keys", exc_info=True)
|
||
return False
|
||
|
||
|
||
def _has_recent_global_progress(
|
||
redis_client: Redis,
|
||
*,
|
||
max_age_seconds: int = STARTUP_GLOBAL_PROGRESS_MAX_AGE_SECONDS,
|
||
) -> bool:
|
||
now = int(time.time())
|
||
try:
|
||
progress_ts = int(redis_client.get(GLOBAL_PROGRESS_TS_KEY) or 0)
|
||
db_progress_ts = int(redis_client.get(GLOBAL_DB_PROGRESS_TS_KEY) or 0)
|
||
except Exception:
|
||
logger.warning("Failed to inspect global startup progress keys", exc_info=True)
|
||
return False
|
||
freshest_ts = max(progress_ts, db_progress_ts)
|
||
return freshest_ts > 0 and now - freshest_ts <= max_age_seconds
|
||
|
||
|
||
def _has_live_startup_progress(redis_client: Redis, *, queue_len: int) -> bool:
|
||
if _has_recent_global_progress(redis_client):
|
||
return True
|
||
if queue_len <= 0:
|
||
return False
|
||
return _has_fresh_active_progress(redis_client)
|
||
|
||
|
||
def _claim_startup_dispatch(redis_client: Redis, *, reset_stale: bool) -> bool:
|
||
if redis_client.set(STARTUP_SYNC_DISPATCH_KEY, "1", nx=True, ex=STARTUP_SYNC_DISPATCH_TTL_SECONDS):
|
||
return True
|
||
if not reset_stale:
|
||
return False
|
||
redis_client.delete(STARTUP_SYNC_DISPATCH_KEY)
|
||
return bool(
|
||
redis_client.set(
|
||
STARTUP_SYNC_DISPATCH_KEY,
|
||
"1",
|
||
nx=True,
|
||
ex=STARTUP_SYNC_DISPATCH_TTL_SECONDS,
|
||
)
|
||
)
|
||
|
||
|
||
@celery_setup_logging.connect
|
||
def _configure_logging(loglevel=None, **kwargs):
|
||
# Пишем логи Celery в stderr.
|
||
level = settings.log_level if settings.log_level else "INFO"
|
||
setup_logging(level, None)
|
||
|
||
|
||
@worker_process_init.connect
|
||
def _on_worker_process_init(**kwargs):
|
||
# Повторно настраиваем логирование после 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(
|
||
"mobilede_scraper",
|
||
broker=_broker_url(),
|
||
backend=_result_backend(),
|
||
)
|
||
|
||
# Если hard limit слишком большой, сжимаем его до soft + 120.
|
||
_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.info(
|
||
"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
|
||
|
||
beat_schedule = {}
|
||
if MOBILEDE_BEAT_SYNC_ENABLED:
|
||
beat_schedule = {
|
||
"periodic-mobilede-sync-search": {
|
||
"task": MOBILEDE_RUNTIME_SEGMENTS_TASK,
|
||
"schedule": _runtime_sync_expires_seconds(),
|
||
"args": (),
|
||
"kwargs": _runtime_sync_kwargs(),
|
||
"options": {
|
||
"queue": MOBILEDE_SYNC_QUEUE,
|
||
"expires": _runtime_sync_expires_seconds(),
|
||
},
|
||
}
|
||
}
|
||
|
||
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=settings.celery.worker_pool,
|
||
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=beat_schedule,
|
||
task_routes={
|
||
MOBILEDE_RUNTIME_SEGMENTS_TASK: {"queue": MOBILEDE_SYNC_QUEUE},
|
||
MOBILEDE_SYNC_TASK_NAME: {"queue": MOBILEDE_SYNC_QUEUE},
|
||
MOBILEDE_SYNC_DETAIL_TASK: {"queue": MOBILEDE_SYNC_QUEUE},
|
||
MOBILEDE_ENRICH_IMAGES_TASK: {"queue": MOBILEDE_SYNC_QUEUE},
|
||
"mobilede_scraper.worker.tasks.*": {"queue": MOBILEDE_SYNC_QUEUE},
|
||
},
|
||
)
|
||
|
||
celery_app.autodiscover_tasks(["mobilede_scraper.worker"])
|
||
|
||
|
||
@worker_ready.connect
|
||
def _on_worker_ready(**kwargs):
|
||
"""При старте worker отправляем первый canonical runtime sync, если очередь пуста."""
|
||
if not _env_bool("MOBILEDE_STARTUP_SYNC_ENABLED", True):
|
||
logger.info("Worker ready: startup sync dispatch disabled by MOBILEDE_STARTUP_SYNC_ENABLED")
|
||
return
|
||
|
||
should_dispatch = False
|
||
redis_client = None
|
||
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,
|
||
)
|
||
|
||
queue_len = int(redis_client.llen(MOBILEDE_SYNC_QUEUE) or 0)
|
||
has_live_progress = _has_live_startup_progress(redis_client, queue_len=queue_len)
|
||
|
||
if queue_len > 0 and has_live_progress:
|
||
logger.info("Worker ready: MOBILEDE_sync queue already has %d task(s); skip startup dispatch", queue_len)
|
||
return
|
||
has_fresh_progress = _has_fresh_active_progress(redis_client)
|
||
if queue_len > 0 and not has_live_progress:
|
||
logger.warning(
|
||
"Worker ready: MOBILEDE_sync queue has %d task(s), but no fresh progress is visible; forcing runtime sync dispatch",
|
||
queue_len,
|
||
)
|
||
elif queue_len <= 0 and has_fresh_progress:
|
||
logger.info("Worker ready: queue is empty but active progress is still visible; relying on startup dedupe")
|
||
|
||
should_dispatch = _claim_startup_dispatch(
|
||
redis_client,
|
||
reset_stale=not has_live_progress,
|
||
)
|
||
if should_dispatch and not has_live_progress:
|
||
logger.info("Worker ready: claimed startup dispatch after stale or missing dedupe state")
|
||
except Exception:
|
||
logger.warning("Worker ready startup sync dedupe check failed; skipping immediate dispatch", exc_info=True)
|
||
return
|
||
finally:
|
||
if redis_client is not None:
|
||
try:
|
||
redis_client.close()
|
||
except Exception:
|
||
pass
|
||
|
||
if not should_dispatch:
|
||
logger.info("Worker ready immediate sync already dispatched recently; skipping duplicate enqueue")
|
||
return
|
||
|
||
logger.info("Worker ready - dispatching initial mobile.de sync_runtime_segments task")
|
||
celery_app.send_task(
|
||
MOBILEDE_RUNTIME_SEGMENTS_TASK,
|
||
kwargs=_runtime_sync_kwargs(),
|
||
queue=MOBILEDE_SYNC_QUEUE,
|
||
expires=_runtime_sync_expires_seconds(),
|
||
)
|