From 0c08fc7f63dcfdaf175890c2fad00c1b723c03d0 Mon Sep 17 00:00:00 2001 From: qananasikq Date: Fri, 17 Apr 2026 23:50:29 +0300 Subject: [PATCH] fix task duplication on retry --- encar_scraper/worker/tasks.py | 55 +++++++++++++++++++++++++---------- 1 file changed, 40 insertions(+), 15 deletions(-) diff --git a/encar_scraper/worker/tasks.py b/encar_scraper/worker/tasks.py index d3c271c..6f3d958 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.exceptions import MaxRetriesExceededError, Retry from celery.signals import worker_process_init from redis import Redis @@ -24,6 +25,18 @@ _redis: Redis | None = None @worker_process_init.connect def _reset_globals_after_fork(**kwargs): global _persistence, _redis + # После fork закрываем унаследованные TCP-сокеты родителя (иначе дети + # делят одни connection pool'ы → повреждение данных). + if _persistence is not None: + try: + _persistence.engine.dispose() + except Exception: + logger.debug("engine.dispose after fork failed", exc_info=True) + if _redis is not None: + try: + _redis.close() + except Exception: + logger.debug("redis.close after fork failed", exc_info=True) _persistence = None _redis = None @@ -156,6 +169,10 @@ def encar_sync_listing_task( heartbeat_stop: Event | None = None heartbeat_thread: Thread | None = None run_id: int | None = None + # Планировать следующий прогон только когда текущий прогон завершён + # (success или final failure после исчерпания retries). Retry сам + # перепланирует задачу — не дублируем apply_async в этом случае. + reschedule_next = True try: heartbeat_stop, heartbeat_thread = _start_lock_heartbeat( @@ -244,7 +261,14 @@ def encar_sync_listing_task( ) except Exception: logger.warning("Failed to record sync_run failure", exc_info=True) - raise self.retry(exc=exc) + try: + # Retry сам перепланирует задачу — не планируем следующий прогон + reschedule_next = False + raise self.retry(exc=exc) + except MaxRetriesExceededError: + # Retries исчерпаны — запускаем следующий прогон по расписанию + reschedule_next = True + raise finally: if heartbeat_stop is not None: heartbeat_stop.set() @@ -252,20 +276,21 @@ 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) + if reschedule_next: + 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(