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