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 7b7efbf525
9 changed files with 245 additions and 53 deletions

View File

@@ -219,6 +219,8 @@ class CeleryConfig:
task_time_limit: int = _env_int("CELERY_TASK_TIME_LIMIT", 3600)
worker_concurrency: int = _env_int("CELERY_WORKER_CONCURRENCY", 4)
worker_max_tasks_per_child: int = _env_int("CELERY_WORKER_MAX_TASKS_PER_CHILD", 5)
# KB; worker перезапустит child при превышении (защита от утечек Chromium/Playwright).
worker_max_memory_per_child_kb: int = _env_int("CELERY_WORKER_MAX_MEMORY_PER_CHILD_KB", 500_000)
broker_visibility_timeout: int = _env_int("CELERY_BROKER_VISIBILITY_TIMEOUT", 7200)
beat_sync_interval_minutes: int = _env_int("CELERY_BEAT_SYNC_INTERVAL_MINUTES", 60)
beat_sync_limit: int | None = _env_int("CELERY_BEAT_SYNC_LIMIT", 0) or None

View File

@@ -12,3 +12,7 @@ class SiteStructureChangedError(ScraperError):
class ListingResumeError(ScraperError):
"""Вызывается, когда resume по checkpoint больше недостижим."""
class ScraperAbortedError(ScraperError):
"""Вызывается при внешнем прерывании (например, потеря Celery lock'а)."""

View File

@@ -1,6 +1,7 @@
import logging
import sys
from contextvars import ContextVar
from logging.handlers import RotatingFileHandler
# Храним trace_id текущего потока/корутины.
TRACE_ID: ContextVar[str] = ContextVar("trace_id", default="-")
@@ -21,7 +22,16 @@ def setup_logging(level: str = "INFO", log_file: str | None = None) -> None:
# stderr — Docker и Celery prefork корректно его подхватывают.
handlers: list[logging.Handler] = [logging.StreamHandler(sys.stderr)]
if log_file:
handlers.append(logging.FileHandler(log_file, encoding="utf-8"))
# Ротация файлового лога: 50MB × 5 файлов, чтобы диск VPS не переполнился
# при многомесячной работе.
handlers.append(
RotatingFileHandler(
log_file,
maxBytes=50 * 1024 * 1024,
backupCount=5,
encoding="utf-8",
)
)
trace_filter = TraceIdFilter()
for handler in handlers:
handler.addFilter(trace_filter)