Files
chunyu_project/utils/safe_task.py
T
2026-08-05 23:59:15 +08:00

88 lines
2.8 KiB
Python

import logging
import threading
import time
logger = logging.getLogger(__name__)
SUBMIT_TIMEOUT = 3
PING_TIMEOUT = 2
WORKER_CHECK_TTL = 30 # worker 探测结果缓存时长(秒),避免每次请求都阻塞 2s
_worker_cache = {'checked_at': 0.0, 'available': None}
def _worker_available(task, timeout=PING_TIMEOUT):
"""探测是否存在存活的 Celery worker(带 TTL 缓存)"""
now = time.time()
if _worker_cache['available'] is not None and now - _worker_cache['checked_at'] < WORKER_CHECK_TTL:
return _worker_cache['available']
try:
inspector = task.app.control.inspect(timeout=timeout)
available = bool(inspector.ping())
except Exception:
available = False
_worker_cache['checked_at'] = now
_worker_cache['available'] = available
return available
def _run_with_timeout(func, args, kwargs, timeout):
result_holder = {}
def target():
try:
result_holder['result'] = func(*args, **kwargs)
except Exception as exc:
result_holder['exception'] = exc
thread = threading.Thread(target=target, daemon=True)
thread.start()
thread.join(timeout)
if thread.is_alive():
raise TimeoutError('task submission exceeded %ds' % timeout)
if 'exception' in result_holder:
raise result_holder['exception']
return result_holder.get('result')
def submit_task(task, *args, **kwargs):
"""发送任务,优先异步提交;无可用 worker 或超时/异常时降级为同步执行"""
timeout = kwargs.pop('_timeout', SUBMIT_TIMEOUT)
if not _worker_available(task):
logger.info(
'submit_task: no celery worker available, run synchronously: task=%s',
getattr(task, 'name', task),
)
try:
return task.apply(args, kwargs)
except Exception as exc:
logger.error(
'submit_task eager fallback failed: task=%s err=%s',
getattr(task, 'name', task), exc,
)
return None
try:
# task.apply_async(args, kwargs) 的正确传参方式
return _run_with_timeout(task.apply_async, (args, kwargs), {}, timeout)
except Exception as exc:
logger.warning(
'submit_task fallback: task=%s args=%r err=%s',
getattr(task, 'name', task), args, exc,
)
try:
return task.apply(args, kwargs)
except Exception as exc2:
logger.error(
'submit_task eager fallback failed: task=%s err=%s',
getattr(task, 'name', task), exc2,
)
return None
def safe_submit(task, *args, fallback=None, **kwargs):
try:
result = submit_task(task, *args, **kwargs)
return result if result is not None else fallback
except Exception:
return fallback