harden: cache PG/Redis singletons, safe orphan-lock cleanup, ScraperAbortedError on lock loss

This commit is contained in:
qananasikq
2026-04-18 00:33:10 +03:00
parent a8c5f8aa41
commit 6a259bbfc3
9 changed files with 245 additions and 53 deletions

View File

@@ -27,6 +27,12 @@ def _on_worker_process_init(**kwargs):
level = settings.log_level if settings.log_level else "INFO"
setup_logging(level, None)
# Сбрасываем cached engine/redis после fork — иначе child наследует
# коннекты родителя, что небезопасно (особенно с psycopg2).
from . import tasks as _tasks
_tasks._CACHED_PERSISTENCE = None
_tasks._CACHED_REDIS = None
def _broker_url() -> str:
return settings.celery.broker_url or settings.redis.url
@@ -66,15 +72,17 @@ celery_app.conf.update(
task_acks_late=True,
task_reject_on_worker_lost=True,
task_track_started=True,
task_ignore_result=True,
worker_concurrency=settings.celery.worker_concurrency,
worker_max_tasks_per_child=settings.celery.worker_max_tasks_per_child,
worker_max_memory_per_child=settings.celery.worker_max_memory_per_child_kb,
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,
result_expires=3600,
worker_redirect_stdouts=False,
worker_hijack_root_logger=False,
beat_schedule={
@@ -107,7 +115,10 @@ def _on_worker_ready(**kwargs):
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))
# Дедуп TTL = интервал beat (по умолчанию 1ч), чтобы не плодить дубликаты
# при flapping-рестартах контейнера.
dedupe_ttl = max(600, settings.celery.beat_sync_interval_minutes * 60)
should_dispatch = bool(redis_client.set(STARTUP_SYNC_DISPATCH_KEY, "1", nx=True, ex=dedupe_ttl))
except Exception:
logger.warning("Worker ready startup sync dedupe check failed; skipping immediate dispatch", exc_info=True)
return