add helper scripts

This commit is contained in:
qananasikq
2026-04-24 20:47:18 +03:00
commit cec1288652
70 changed files with 14249 additions and 0 deletions

View File

@@ -0,0 +1,148 @@
# Инициализация 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("dubizzle_scraper.worker.celery_app")
@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(
"dubizzle_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="solo",
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": "dubizzle_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 if settings.discovery.always_full_scan else True,
},
"options": {
"queue": "scraping",
"expires": settings.celery.beat_sync_interval_minutes * 60.0,
},
}
},
task_routes={
"dubizzle_scraper.worker.tasks.*": {"queue": "scraping"}
},
)
celery_app.autodiscover_tasks(["dubizzle_scraper.worker"])
@worker_ready.connect
def _on_worker_ready(**kwargs):
"""При старте worker отправляем первый sync_listing, если очередь пуста."""
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,
)
for stale_key in ("dubizzle:locks:sync_listing",):
try:
ttl = redis_client.ttl(stale_key)
if ttl is not None and ttl != -2:
redis_client.delete(stale_key)
logger.warning("Cleared stale lock on startup: %s (ttl was %s)", stale_key, ttl)
except Exception:
logger.warning("Failed to clear stale lock %s on startup", stale_key, exc_info=True)
try:
queue_len = int(redis_client.llen("scraping") or 0)
except Exception:
queue_len = 0
if queue_len > 0:
logger.info("Worker ready: scraping queue already has %d task(s); skip startup dispatch", queue_len)
return
except Exception:
logger.warning("Worker ready startup sync check failed; skipping immediate dispatch", exc_info=True)
return
finally:
if redis_client is not None:
try:
redis_client.close()
except Exception:
pass
logger.info("Worker ready — dispatching initial sync_listing task")
celery_app.send_task(
"dubizzle_scraper.worker.tasks.sync_listing_task",
kwargs={"limit": settings.celery.beat_sync_limit, "only_new": False},
queue="scraping",
expires=settings.celery.beat_sync_interval_minutes * 60.0,
)