fix task duplication on retry
This commit is contained in:
@@ -5,6 +5,7 @@ from threading import Event, Thread
|
|||||||
import uuid
|
import uuid
|
||||||
|
|
||||||
from celery import shared_task
|
from celery import shared_task
|
||||||
|
from celery.exceptions import MaxRetriesExceededError, Retry
|
||||||
from celery.signals import worker_process_init
|
from celery.signals import worker_process_init
|
||||||
from redis import Redis
|
from redis import Redis
|
||||||
|
|
||||||
@@ -24,6 +25,18 @@ _redis: Redis | None = None
|
|||||||
@worker_process_init.connect
|
@worker_process_init.connect
|
||||||
def _reset_globals_after_fork(**kwargs):
|
def _reset_globals_after_fork(**kwargs):
|
||||||
global _persistence, _redis
|
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
|
_persistence = None
|
||||||
_redis = None
|
_redis = None
|
||||||
|
|
||||||
@@ -156,6 +169,10 @@ def encar_sync_listing_task(
|
|||||||
heartbeat_stop: Event | None = None
|
heartbeat_stop: Event | None = None
|
||||||
heartbeat_thread: Thread | None = None
|
heartbeat_thread: Thread | None = None
|
||||||
run_id: int | None = None
|
run_id: int | None = None
|
||||||
|
# Планировать следующий прогон только когда текущий прогон завершён
|
||||||
|
# (success или final failure после исчерпания retries). Retry сам
|
||||||
|
# перепланирует задачу — не дублируем apply_async в этом случае.
|
||||||
|
reschedule_next = True
|
||||||
|
|
||||||
try:
|
try:
|
||||||
heartbeat_stop, heartbeat_thread = _start_lock_heartbeat(
|
heartbeat_stop, heartbeat_thread = _start_lock_heartbeat(
|
||||||
@@ -244,7 +261,14 @@ def encar_sync_listing_task(
|
|||||||
)
|
)
|
||||||
except Exception:
|
except Exception:
|
||||||
logger.warning("Failed to record sync_run failure", exc_info=True)
|
logger.warning("Failed to record sync_run failure", exc_info=True)
|
||||||
|
try:
|
||||||
|
# Retry сам перепланирует задачу — не планируем следующий прогон
|
||||||
|
reschedule_next = False
|
||||||
raise self.retry(exc=exc)
|
raise self.retry(exc=exc)
|
||||||
|
except MaxRetriesExceededError:
|
||||||
|
# Retries исчерпаны — запускаем следующий прогон по расписанию
|
||||||
|
reschedule_next = True
|
||||||
|
raise
|
||||||
finally:
|
finally:
|
||||||
if heartbeat_stop is not None:
|
if heartbeat_stop is not None:
|
||||||
heartbeat_stop.set()
|
heartbeat_stop.set()
|
||||||
@@ -252,6 +276,7 @@ def encar_sync_listing_task(
|
|||||||
heartbeat_thread.join(timeout=max(1.0, min(5.0, lock_ttl / 10)))
|
heartbeat_thread.join(timeout=max(1.0, min(5.0, lock_ttl / 10)))
|
||||||
if lock_acquired:
|
if lock_acquired:
|
||||||
_release_lock_if_owner(redis_client, ENCAR_SYNC_LOCK_KEY, owner_token)
|
_release_lock_if_owner(redis_client, ENCAR_SYNC_LOCK_KEY, owner_token)
|
||||||
|
if reschedule_next:
|
||||||
try:
|
try:
|
||||||
interval_minutes = Settings().celery.encar_beat_interval_minutes
|
interval_minutes = Settings().celery.encar_beat_interval_minutes
|
||||||
countdown = max(60, int(interval_minutes) * 60)
|
countdown = max(60, int(interval_minutes) * 60)
|
||||||
|
|||||||
Reference in New Issue
Block a user