88 lines
2.8 KiB
Python
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
|