diff --git a/.env.example b/.env.example index a8563ea..3b193ec 100644 --- a/.env.example +++ b/.env.example @@ -61,3 +61,4 @@ CELERY_TASK_SOFT_TIME_LIMIT=600 CELERY_TASK_TIME_LIMIT=900 CELERY_WORKER_CONCURRENCY=1 CELERY_BEAT_SYNC_INTERVAL_MINUTES=60 +CELERY_BEAT_SYNC_LIMIT=26 diff --git a/README.md b/README.md index 9d15930..fd467b6 100644 --- a/README.md +++ b/README.md @@ -23,6 +23,8 @@ docker compose up -d Поднимутся 5 контейнеров: postgres, redis, api, worker, beat. API на `http://localhost:8000`. +По умолчанию `beat` запускает сбор листинга **раз в 1 час** и обрабатывает **до 26 машин за запуск** (`CELERY_BEAT_SYNC_INTERVAL_MINUTES=60`, `CELERY_BEAT_SYNC_LIMIT=26`). + Swagger-документация: `http://localhost:8000/docs` ## API diff --git a/iaai_scraper/api/routes/tasks.py b/iaai_scraper/api/routes/tasks.py index d0b2695..bb350b1 100644 --- a/iaai_scraper/api/routes/tasks.py +++ b/iaai_scraper/api/routes/tasks.py @@ -37,7 +37,10 @@ def start_sync_vehicle( persistence: PersistenceService = Depends(get_persistence), ): """Запустить задачу скрапинга одного автомобиля через Celery.""" - result = sync_vehicle_task.delay(body.vehicle_url, lane=body.lane) + result = sync_vehicle_task.apply_async( + kwargs={"vehicle_url": body.vehicle_url, "lane": body.lane}, + queue="scraping", + ) # Регистрируем задачу в БД. persistence.create_scrape_task( @@ -59,12 +62,15 @@ def start_sync_listing( persistence: PersistenceService = Depends(get_persistence), ): """Запустить задачу полного цикла листинга через Celery.""" - result = sync_listing_task.delay( - make=body.make, - model=body.model, - lane=body.lane, - limit=body.limit, - only_new=body.only_new, + result = sync_listing_task.apply_async( + kwargs={ + "make": body.make, + "model": body.model, + "lane": body.lane, + "limit": body.limit, + "only_new": body.only_new, + }, + queue="scraping", ) persistence.create_scrape_task( diff --git a/iaai_scraper/core/config.py b/iaai_scraper/core/config.py index 4ec3c78..7c19e55 100644 --- a/iaai_scraper/core/config.py +++ b/iaai_scraper/core/config.py @@ -91,6 +91,7 @@ class CeleryConfig: task_time_limit: int = int(os.getenv("CELERY_TASK_TIME_LIMIT", "900")) worker_concurrency: int = int(os.getenv("CELERY_WORKER_CONCURRENCY", "1")) beat_sync_interval_minutes: int = int(os.getenv("CELERY_BEAT_SYNC_INTERVAL_MINUTES", "60")) + beat_sync_limit: int = int(os.getenv("CELERY_BEAT_SYNC_LIMIT", "26")) @dataclass(slots=True) diff --git a/iaai_scraper/worker/celery_app.py b/iaai_scraper/worker/celery_app.py index 3384ff0..52c91a3 100644 --- a/iaai_scraper/worker/celery_app.py +++ b/iaai_scraper/worker/celery_app.py @@ -45,6 +45,7 @@ celery_app.conf.update( "task": "iaai_scraper.worker.tasks.sync_listing_task", "schedule": settings.celery.beat_sync_interval_minutes * 60.0, "args": (), + "kwargs": {"limit": settings.celery.beat_sync_limit}, "options": {"queue": "scraping"}, }, },