Prepared for production
This commit is contained in:
@@ -1,6 +1,7 @@
|
||||
# Инициализация Celery-приложения и периодических задач.
|
||||
|
||||
import logging
|
||||
import os
|
||||
|
||||
from celery import Celery
|
||||
from celery.signals import worker_process_init, worker_ready, setup_logging as celery_setup_logging
|
||||
@@ -11,6 +12,12 @@ from ..core.logs import setup_logging
|
||||
|
||||
logger = logging.getLogger("iaai_scraper.worker.celery_app")
|
||||
STARTUP_SYNC_DISPATCH_KEY = "iaai:state:startup_sync_dispatched"
|
||||
IAAI_SYNC_QUEUE = "iaai_sync"
|
||||
|
||||
|
||||
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"}
|
||||
|
||||
|
||||
@celery_setup_logging.connect
|
||||
@@ -79,18 +86,19 @@ celery_app.conf.update(
|
||||
worker_hijack_root_logger=False,
|
||||
beat_schedule={
|
||||
"periodic-sync-listing": {
|
||||
"task": "iaai_scraper.worker.tasks.sync_listing_task",
|
||||
"task": "iaai.sync_cars_feed",
|
||||
"schedule": settings.celery.beat_sync_interval_minutes * 60.0,
|
||||
"args": (),
|
||||
"kwargs": {"limit": settings.celery.beat_sync_limit, "only_new": True},
|
||||
"options": {
|
||||
"queue": "scraping",
|
||||
"queue": IAAI_SYNC_QUEUE,
|
||||
"expires": settings.celery.beat_sync_interval_minutes * 60.0,
|
||||
},
|
||||
}
|
||||
},
|
||||
task_routes={
|
||||
"iaai_scraper.worker.tasks.*": {"queue": "scraping"}
|
||||
"iaai.sync_cars_feed": {"queue": IAAI_SYNC_QUEUE},
|
||||
"iaai_scraper.worker.tasks.*": {"queue": IAAI_SYNC_QUEUE},
|
||||
},
|
||||
)
|
||||
|
||||
@@ -100,6 +108,10 @@ celery_app.autodiscover_tasks(["iaai_scraper.worker"])
|
||||
@worker_ready.connect
|
||||
def _on_worker_ready(**kwargs):
|
||||
"""При старте worker отправляем первый sync_listing, если очередь пуста."""
|
||||
if not _env_bool("IAAI_STARTUP_SYNC_ENABLED", True):
|
||||
logger.info("Worker ready: startup sync dispatch disabled by IAAI_STARTUP_SYNC_ENABLED")
|
||||
return
|
||||
|
||||
redis_client = None
|
||||
try:
|
||||
redis_client = Redis.from_url(
|
||||
@@ -121,11 +133,11 @@ def _on_worker_ready(**kwargs):
|
||||
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)
|
||||
queue_len = int(redis_client.llen(IAAI_SYNC_QUEUE) 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)
|
||||
logger.info("Worker ready: iaai_sync queue already has %d task(s); skip startup dispatch", queue_len)
|
||||
return
|
||||
|
||||
should_dispatch = bool(redis_client.set(STARTUP_SYNC_DISPATCH_KEY, "1", nx=True, ex=600))
|
||||
@@ -145,8 +157,8 @@ def _on_worker_ready(**kwargs):
|
||||
|
||||
logger.info("Worker ready — dispatching initial sync_listing task")
|
||||
celery_app.send_task(
|
||||
"iaai_scraper.worker.tasks.sync_listing_task",
|
||||
"iaai.sync_cars_feed",
|
||||
kwargs={"limit": settings.celery.beat_sync_limit, "only_new": False},
|
||||
queue="scraping",
|
||||
queue=IAAI_SYNC_QUEUE,
|
||||
expires=settings.celery.beat_sync_interval_minutes * 60.0,
|
||||
)
|
||||
|
||||
@@ -19,6 +19,8 @@ from ..discovery import SitemapDiscoveryError, discover_vehicle_urls_from_sitema
|
||||
|
||||
logger = logging.getLogger("iaai_scraper.worker.tasks")
|
||||
|
||||
IAAI_SYNC_QUEUE = "iaai_sync"
|
||||
|
||||
# Минимальная пауза между батчами (секунды) — не давит IAAI.
|
||||
_INTER_BATCH_DELAY = max(float(os.getenv("IAAI_INTER_BATCH_DELAY_SECONDS", "0.3")), 0.0)
|
||||
# Если доля failed в батче превышает порог — прерываем (IAAI блокирует).
|
||||
@@ -44,7 +46,7 @@ HOURLY_FAILURE_STREAK_KEY = "iaai:state:hourly_failure_streak"
|
||||
HOURLY_FAILURE_STREAK_LIMIT = 3
|
||||
HOURLY_FAILURE_STREAK_TTL_SECONDS = 6 * 60 * 60 # сброс через 6 часов
|
||||
SYNC_LISTING_FOLLOWUP_PENDING_KEY = "iaai:state:sync_listing_followup_pending"
|
||||
SYNC_LISTING_TASK_NAME = "iaai_scraper.worker.tasks.sync_listing_task"
|
||||
SYNC_LISTING_TASK_NAME = "iaai.sync_cars_feed"
|
||||
SYNC_SEGMENT_LOCK_KEY_FMT = "iaai:locks:sync_segment:{idx}"
|
||||
SYNC_SEGMENTS_PROGRESS_KEY = "iaai:state:sync_segments_progress"
|
||||
SYNC_SEGMENTS_TOTAL_KEY = "iaai:state:sync_segments_total"
|
||||
@@ -482,7 +484,7 @@ def _start_stall_watchdog(
|
||||
celery_app.send_task(
|
||||
SYNC_LISTING_TASK_NAME,
|
||||
kwargs={},
|
||||
queue="scraping",
|
||||
queue=IAAI_SYNC_QUEUE,
|
||||
countdown=15,
|
||||
expires=followup_ttl,
|
||||
)
|
||||
@@ -903,6 +905,7 @@ def _mark_segment_completed(redis_client: Redis, segment_index: int) -> tuple[in
|
||||
|
||||
@shared_task(
|
||||
name="iaai_scraper.worker.tasks.sync_segment_task",
|
||||
queue=IAAI_SYNC_QUEUE,
|
||||
bind=True,
|
||||
max_retries=2,
|
||||
default_retry_delay=60,
|
||||
@@ -1076,6 +1079,7 @@ def sync_segment_task(
|
||||
|
||||
@shared_task(
|
||||
name="iaai_scraper.worker.tasks.sync_vehicle_task",
|
||||
queue=IAAI_SYNC_QUEUE,
|
||||
bind=True,
|
||||
max_retries=2,
|
||||
default_retry_delay=30,
|
||||
@@ -1106,7 +1110,8 @@ def sync_vehicle_task(self, vehicle_url: str, lane: str = "iaai"):
|
||||
|
||||
|
||||
@shared_task(
|
||||
name="iaai_scraper.worker.tasks.sync_listing_task",
|
||||
name=SYNC_LISTING_TASK_NAME,
|
||||
queue=IAAI_SYNC_QUEUE,
|
||||
bind=True,
|
||||
max_retries=3,
|
||||
default_retry_delay=120,
|
||||
@@ -1178,7 +1183,7 @@ def sync_listing_task(
|
||||
try:
|
||||
followup_expires = max(int(lock_ttl), int(delay_seconds) + 300)
|
||||
self.app.send_task(
|
||||
"iaai_scraper.worker.tasks.sync_listing_task",
|
||||
SYNC_LISTING_TASK_NAME,
|
||||
kwargs={
|
||||
"make": make,
|
||||
"model": model,
|
||||
@@ -1186,7 +1191,7 @@ def sync_listing_task(
|
||||
"limit": limit,
|
||||
"only_new": only_new,
|
||||
},
|
||||
queue="scraping",
|
||||
queue=IAAI_SYNC_QUEUE,
|
||||
countdown=max(0, int(delay_seconds)),
|
||||
expires=followup_expires,
|
||||
)
|
||||
@@ -1480,7 +1485,7 @@ def sync_listing_task(
|
||||
"only_new": effective_only_new,
|
||||
"is_bootstrap": force_bootstrap_full_scan,
|
||||
},
|
||||
queue="scraping",
|
||||
queue=IAAI_SYNC_QUEUE,
|
||||
)
|
||||
dispatched += 1
|
||||
except Exception:
|
||||
@@ -1654,7 +1659,7 @@ def sync_listing_task(
|
||||
if _try_set_followup_pending(redis_client, ttl_seconds=followup_ttl):
|
||||
try:
|
||||
self.app.send_task(
|
||||
"iaai_scraper.worker.tasks.sync_listing_task",
|
||||
SYNC_LISTING_TASK_NAME,
|
||||
kwargs={
|
||||
"make": make,
|
||||
"model": model,
|
||||
@@ -1662,7 +1667,7 @@ def sync_listing_task(
|
||||
"limit": limit,
|
||||
"only_new": only_new,
|
||||
},
|
||||
queue="scraping",
|
||||
queue=IAAI_SYNC_QUEUE,
|
||||
countdown=10,
|
||||
expires=followup_ttl,
|
||||
)
|
||||
|
||||
Reference in New Issue
Block a user