improve runtime sync flow
This commit is contained in:
@@ -1,4 +1,4 @@
|
||||
# Инициализация Celery-приложения и периодических задач.
|
||||
# Celery-приложение и периодические задачи.
|
||||
|
||||
import json
|
||||
import logging
|
||||
@@ -11,12 +11,22 @@ 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
|
||||
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"
|
||||
MOBILEDE_SYNC_QUEUE = "mobilede_sync"
|
||||
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:
|
||||
@@ -27,7 +37,23 @@ def _env_bool(name: str, default: bool) -> bool:
|
||||
MOBILEDE_BEAT_SYNC_ENABLED = _env_bool("MOBILEDE_BEAT_SYNC_ENABLED", True)
|
||||
|
||||
|
||||
def _has_fresh_active_progress(redis_client: Redis, *, max_age_seconds: int = 180) -> bool:
|
||||
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}*"):
|
||||
@@ -47,7 +73,11 @@ def _has_fresh_active_progress(redis_client: Redis, *, max_age_seconds: int = 18
|
||||
return False
|
||||
|
||||
|
||||
def _has_recent_global_progress(redis_client: Redis, *, max_age_seconds: int = 300) -> bool:
|
||||
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)
|
||||
@@ -59,17 +89,40 @@ def _has_recent_global_progress(redis_client: Redis, *, max_age_seconds: int = 3
|
||||
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 (Docker logs).
|
||||
# Пишем логи 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):
|
||||
# Повторно настраиваем логирование в каждом дочернем prefork-процессе,
|
||||
# чтобы StreamHandler(stderr) корректно работал после fork.
|
||||
# Повторно настраиваем логирование после fork.
|
||||
level = settings.log_level if settings.log_level else "INFO"
|
||||
setup_logging(level, None)
|
||||
|
||||
@@ -88,8 +141,7 @@ celery_app = Celery(
|
||||
backend=_result_backend(),
|
||||
)
|
||||
|
||||
# Auto-clamp: если hard limit слишком далёк от soft (> soft + 120),
|
||||
# ограничиваем, чтобы зависший worker не жил вечно.
|
||||
# Если 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
|
||||
@@ -105,17 +157,13 @@ beat_schedule = {}
|
||||
if MOBILEDE_BEAT_SYNC_ENABLED:
|
||||
beat_schedule = {
|
||||
"periodic-mobilede-sync-search": {
|
||||
"task": "mobilede.sync_runtime_segments",
|
||||
"schedule": settings.celery.beat_sync_interval_minutes * 60.0,
|
||||
"task": MOBILEDE_RUNTIME_SEGMENTS_TASK,
|
||||
"schedule": _runtime_sync_expires_seconds(),
|
||||
"args": (),
|
||||
"kwargs": {
|
||||
"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),
|
||||
},
|
||||
"kwargs": _runtime_sync_kwargs(),
|
||||
"options": {
|
||||
"queue": MOBILEDE_SYNC_QUEUE,
|
||||
"expires": settings.celery.beat_sync_interval_minutes * 60.0,
|
||||
"expires": _runtime_sync_expires_seconds(),
|
||||
},
|
||||
}
|
||||
}
|
||||
@@ -144,9 +192,10 @@ celery_app.conf.update(
|
||||
worker_hijack_root_logger=False,
|
||||
beat_schedule=beat_schedule,
|
||||
task_routes={
|
||||
"mobilede.sync_runtime_segments": {"queue": MOBILEDE_SYNC_QUEUE},
|
||||
"mobilede.sync_search": {"queue": MOBILEDE_SYNC_QUEUE},
|
||||
"mobilede.sync_detail": {"queue": MOBILEDE_SYNC_QUEUE},
|
||||
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},
|
||||
},
|
||||
)
|
||||
@@ -161,6 +210,7 @@ def _on_worker_ready(**kwargs):
|
||||
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(
|
||||
@@ -172,29 +222,27 @@ def _on_worker_ready(**kwargs):
|
||||
retry_on_timeout=True,
|
||||
)
|
||||
|
||||
has_fresh_progress = _has_fresh_active_progress(redis_client)
|
||||
has_recent_global_progress = _has_recent_global_progress(redis_client)
|
||||
has_live_progress = bool(has_fresh_progress or has_recent_global_progress)
|
||||
queue_len = int(redis_client.llen(MOBILEDE_SYNC_QUEUE) or 0)
|
||||
has_live_progress = _has_live_startup_progress(redis_client, queue_len=queue_len)
|
||||
|
||||
try:
|
||||
queue_len = int(redis_client.llen(MOBILEDE_SYNC_QUEUE) or 0)
|
||||
except Exception:
|
||||
queue_len = 0
|
||||
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 = bool(redis_client.set(STARTUP_SYNC_DISPATCH_KEY, "1", nx=True, ex=600))
|
||||
if not should_dispatch and not has_live_progress:
|
||||
redis_client.delete(STARTUP_SYNC_DISPATCH_KEY)
|
||||
should_dispatch = bool(redis_client.set(STARTUP_SYNC_DISPATCH_KEY, "1", nx=True, ex=600))
|
||||
if should_dispatch:
|
||||
logger.info("Worker ready: stale startup dedupe key ignored because queue is empty and no fresh active progress exists")
|
||||
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
|
||||
@@ -211,12 +259,8 @@ def _on_worker_ready(**kwargs):
|
||||
|
||||
logger.info("Worker ready - dispatching initial mobile.de sync_runtime_segments task")
|
||||
celery_app.send_task(
|
||||
"mobilede.sync_runtime_segments",
|
||||
kwargs={
|
||||
"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),
|
||||
},
|
||||
MOBILEDE_RUNTIME_SEGMENTS_TASK,
|
||||
kwargs=_runtime_sync_kwargs(),
|
||||
queue=MOBILEDE_SYNC_QUEUE,
|
||||
expires=settings.celery.beat_sync_interval_minutes * 60.0,
|
||||
expires=_runtime_sync_expires_seconds(),
|
||||
)
|
||||
|
||||
Reference in New Issue
Block a user