diff --git a/encar_scraper/worker/celery_app.py b/encar_scraper/worker/celery_app.py index 9a373a3..a2d6ce8 100644 --- a/encar_scraper/worker/celery_app.py +++ b/encar_scraper/worker/celery_app.py @@ -40,17 +40,9 @@ celery_app.conf.update( }, result_expires=86400, worker_hijack_root_logger=False, - beat_schedule={ - "periodic-encar-sync": { - "task": "encar_scraper.worker.tasks.encar_sync_listing_task", - "schedule": settings.celery.encar_beat_interval_minutes * 60.0, - "args": (), - "kwargs": { - "car_type": "all", - }, - "options": {"queue": "encar"}, - }, - }, + # Следующий запуск планирует сама задача в finally через apply_async(countdown=...), + # чтобы час отсчитывался от завершения прогона, а не от старта. + beat_schedule={}, task_routes={ "encar_scraper.worker.tasks.*": {"queue": "encar"}, }, diff --git a/encar_scraper/worker/tasks.py b/encar_scraper/worker/tasks.py index fdb1e11..154c7db 100644 --- a/encar_scraper/worker/tasks.py +++ b/encar_scraper/worker/tasks.py @@ -5,6 +5,7 @@ from threading import Event, Thread import uuid from celery import shared_task +from celery.signals import worker_process_init from redis import Redis from ..core.config import Settings @@ -20,6 +21,13 @@ _persistence: PersistenceService | None = None _redis: Redis | None = None +@worker_process_init.connect +def _reset_globals_after_fork(**kwargs): + global _persistence, _redis + _persistence = None + _redis = None + + def _sync_listing_lock_ttl_seconds() -> int: settings = Settings() return max(settings.celery.task_time_limit + 120, 300) @@ -147,6 +155,7 @@ def encar_sync_listing_task( heartbeat_stop: Event | None = None heartbeat_thread: Thread | None = None + run_id: int | None = None try: heartbeat_stop, heartbeat_thread = _start_lock_heartbeat( @@ -162,6 +171,9 @@ def encar_sync_listing_task( effective_limit = limit or rc.sync.limit effective_lane = lane or rc.sync.lane or "encar" + # --- Фиксируем начало sync_run --- + run_id = persistence.start_sync_run(effective_lane) + car_type_map = {"all": "A", "domestic": "Y", "import": "N"} filters = EncarFilters( car_type=car_type_map.get(car_type, "A"), @@ -180,7 +192,7 @@ def encar_sync_listing_task( set(rc.filters.exclude_brands) if rc.filters.exclude_brands else None ) - scraper = EncarScraper() + scraper = EncarScraper(persistence=persistence) result = scraper.sync_listing( limit=effective_limit, filters=filters, @@ -191,6 +203,16 @@ def encar_sync_listing_task( probe_all_photos=rc.sync.probe_all_photos, ) + # --- Запись результата в sync_runs --- + persistence.finish_sync_run( + run_id, + status="success", + ids_fetched=result.get("items_collected", 0), + cars_upserted=result.get("cars_synced", 0), + cars_failed=result.get("cars_failed", 0), + images_upserted=0, + ) + summary = { "task_id": task_id, "source": "encar", @@ -208,6 +230,19 @@ def encar_sync_listing_task( except Exception as exc: logger.error("encar_sync_listing_task failed: %s", exc, exc_info=True) + try: + if run_id is not None: + persistence.finish_sync_run( + run_id, + status="failed", + ids_fetched=0, + cars_upserted=0, + cars_failed=0, + images_upserted=0, + error_summary=str(exc)[:500], + ) + except Exception: + logger.warning("Failed to record sync_run failure", exc_info=True) raise self.retry(exc=exc) finally: if heartbeat_stop is not None: @@ -216,6 +251,20 @@ def encar_sync_listing_task( heartbeat_thread.join(timeout=max(1.0, min(5.0, lock_ttl / 10))) if lock_acquired: _release_lock_if_owner(redis_client, ENCAR_SYNC_LOCK_KEY, owner_token) + try: + interval_minutes = Settings().celery.encar_beat_interval_minutes + countdown = max(60, int(interval_minutes) * 60) + self.apply_async( + kwargs={"car_type": car_type}, + countdown=countdown, + queue="encar", + ) + logger.info( + "Next encar_sync_listing_task scheduled in %d minutes", + countdown // 60, + ) + except Exception: + logger.warning("Failed to schedule next sync run", exc_info=True) @shared_task( @@ -231,7 +280,7 @@ def encar_sync_vehicle_task(self, vehicle_url: str, lane: str = "encar"): persistence.create_tables() try: - scraper = EncarScraper() + scraper = EncarScraper(persistence=persistence) result = scraper.sync_vehicle(vehicle_url, lane=lane) logger.info("encar_sync_vehicle_task completed: %s", vehicle_url) return {