diff --git a/apps/worker/worker_app/tasks/vision/fast_path.py b/apps/worker/worker_app/tasks/vision/fast_path.py index 6bedcd237..715eb2194 100644 --- a/apps/worker/worker_app/tasks/vision/fast_path.py +++ b/apps/worker/worker_app/tasks/vision/fast_path.py @@ -21,9 +21,9 @@ logger = logging.getLogger(__name__) # 可通过环境变量调参(有默认值,无需配置即可跑) _IMG_WORKERS = int(os.environ.get("VISION_V2_IMG_WORKERS", "8")) -_FAST_TIMEOUT = float(os.environ.get("VISION_V2_FAST_TIMEOUT", "8")) -_FAST_JSON_TIMEOUT = float(os.environ.get("VISION_V2_FAST_JSON_TIMEOUT", "6")) -_OCR_TIMEOUT = float(os.environ.get("VISION_V2_OCR_TIMEOUT", "6")) +_FAST_TIMEOUT = float(os.environ.get("VISION_V2_FAST_TIMEOUT", "6")) +_FAST_JSON_TIMEOUT = float(os.environ.get("VISION_V2_FAST_JSON_TIMEOUT", "5")) +_OCR_TIMEOUT = float(os.environ.get("VISION_V2_OCR_TIMEOUT", "5")) _PRO_TIMEOUT = float(os.environ.get("VISION_V2_PRO_TIMEOUT", "45")) _FALLBACK_RESULT = { @@ -62,16 +62,23 @@ def analyze_image_v2(idx: int, img_url: str) -> dict[str, Any]: with ThreadPoolExecutor(max_workers=2) as pool: f_fj = pool.submit(vlm_fast_json.call_fast_json, img_url, timeout=_FAST_JSON_TIMEOUT) f_ocr = pool.submit(ocr_volc.call_ocr, img_url, timeout=_OCR_TIMEOUT) - for fut in as_completed([f_fj, f_ocr], timeout=_FAST_TIMEOUT + 2): - try: - res = fut.result(timeout=1) - except Exception as e: - logger.warning("[vision.v2] 图片 #%d 子任务异常: %s", idx, e) - continue - if fut is f_fj and isinstance(res, dict): - fj_result = res - elif fut is f_ocr and isinstance(res, list): - ocr_result = res + try: + for fut in as_completed([f_fj, f_ocr], timeout=_FAST_TIMEOUT): + try: + res = fut.result(timeout=1) + except Exception as e: + logger.warning("[vision.v2] 图片 #%d 子任务异常: %s", idx, e) + continue + if fut is f_fj and isinstance(res, dict): + fj_result = res + elif fut is f_ocr and isinstance(res, list): + ocr_result = res + except TimeoutError: + # fast 整体超时,取消还没跑完的子任务,继续走 pro 兜底 + for f in (f_fj, f_ocr): + if not f.done(): + f.cancel() + logger.warning("[vision.v2] 图片 #%d fast路径超时(%.0fs),走pro兜底", idx, _FAST_TIMEOUT) fast_elapsed = time.time() - t0 diff --git a/apps/worker/worker_app/tasks/vision/vlm_fallback.py b/apps/worker/worker_app/tasks/vision/vlm_fallback.py index bb6e0481a..558500c17 100644 --- a/apps/worker/worker_app/tasks/vision/vlm_fallback.py +++ b/apps/worker/worker_app/tasks/vision/vlm_fallback.py @@ -164,6 +164,8 @@ def call_pro_vlm( return None use_model = model or DEFAULT_PRO_MODEL + _orig_retries = client.max_retries + client.max_retries = 0 try: raw = client.vision_completion( messages=[{"role": "system", "content": system}, {"role": "user", "content": user}], @@ -175,7 +177,9 @@ def call_pro_vlm( ) except Exception as e: logger.warning("[vision.vlm] 图片 #%d pro VLM 调用失败 elapsed=%.1fs err=%s", idx, time.time() - t0, e) + client.max_retries = _orig_retries return None + client.max_retries = _orig_retries elapsed = time.time() - t0 if not raw: diff --git a/apps/worker/worker_app/tasks/vision/vlm_fast_json.py b/apps/worker/worker_app/tasks/vision/vlm_fast_json.py index 6534061ff..3b93b6580 100644 --- a/apps/worker/worker_app/tasks/vision/vlm_fast_json.py +++ b/apps/worker/worker_app/tasks/vision/vlm_fast_json.py @@ -93,17 +93,23 @@ def call_fast_json( return None use_model = model or DEFAULT_LITE_MODEL - raw = client.vision_completion( - messages=[ - {"role": "system", "content": _FAST_SYSTEM}, - {"role": "user", "content": _FAST_USER}, - ], - images=[img_url], - temperature=0.1, - max_tokens=max_tokens, - timeout=timeout, - model=use_model, - ) + # 强制不重试:lite 是快速路径,失败直接走外层 pro 兜底 + _orig_retries = client.max_retries + client.max_retries = 0 + try: + raw = client.vision_completion( + messages=[ + {"role": "system", "content": _FAST_SYSTEM}, + {"role": "user", "content": _FAST_USER}, + ], + images=[img_url], + temperature=0.1, + max_tokens=max_tokens, + timeout=timeout, + model=use_model, + ) + finally: + client.max_retries = _orig_retries elapsed = time.time() - t0 if raw is None: logger.warning("[vision.v2] fast_json 返回 None elapsed=%.1fs model=%s", elapsed, use_model)