From 9699a1fcde88475bc23fc80c07002f11b3454724 Mon Sep 17 00:00:00 2001 From: xiaoxia Date: Mon, 5 Oct 2026 16:08:00 +0800 Subject: [PATCH] =?UTF-8?q?fix(viral=5Fvideo):=20#2198=20lite/pro=E5=B9=B6?= =?UTF-8?q?=E8=A1=8C=E7=AB=9E=E9=80=9F=EF=BC=8C=E5=8D=95=E5=9B=BE=E6=9C=80?= =?UTF-8?q?=E5=9D=8F75s=EF=BC=88=E5=8E=9F=E4=B8=B2=E8=A1=8C100s+=EF=BC=89?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 问题:实测lite 40s 100%超时→pro串行跑40-65s,3图2并发最坏=40+60+40=140-190s,img#1粉波点裙pro 60s超时直接失败。 改动: 1. _analyze_single_image内部用2线程ThreadPoolExecutor并行发lite(30s)和pro(75s), as_completed取第一个usable结果即返回,消除串行等待惩罚 2. _step_image_analysis阶段统一在try外层置client.max_retries=0、finally恢复, 子线程_call不再嵌套修改client属性避免竞态 3. lite timeout 40→30s(快速路径30s还没出就等pro),pro timeout 60→75s(给10s余量防偶发慢) 4. 单图最坏75s(pro慢到75s才出),典型40-50s,3图2并发≈75s;外层图片并发仍≤2(总VLM并发=4) 5. elapsed日志加label字段(lite/pro)便于区分竞速胜出方 --- CHANGES.md | 1 + apps/worker/worker_app/tasks/viral_video.py | 167 ++++++++++++-------- 2 files changed, 104 insertions(+), 64 deletions(-) diff --git a/CHANGES.md b/CHANGES.md index ee7fd68b1..ee49745ac 100644 --- a/CHANGES.md +++ b/CHANGES.md @@ -1,2 +1,3 @@ - 2026-10-05 #2194 VLM timeout tune + #2195 HEAD→GET Range fix deployed to staging +# 2198 lite/pro并行竞速 diff --git a/apps/worker/worker_app/tasks/viral_video.py b/apps/worker/worker_app/tasks/viral_video.py index 81c4c0b7d..f50376b03 100644 --- a/apps/worker/worker_app/tasks/viral_video.py +++ b/apps/worker/worker_app/tasks/viral_video.py @@ -425,9 +425,9 @@ def _analyze_single_image( image_urls=f"第1张:{img_url}", ) - def _call(model: str, tmo: int): - # #2194: 单次调用临时关闭 httpx 层重试,超时/失败由外层 pro_fallback 统一兜底, - # 避免底层 max_retries=1 导致 lite timeout × 2 + pro timeout × 2 最坏 240s+ + def _call(model: str, tmo: int, label: str): + # #2194/#2198: max_retries=0 由外层 _step_image_analysis 统一设置(阶段前置0、阶段后恢复), + # 子线程只读不改,避免嵌套并行竞速时多线程同时改 client.max_retries 产生竞态 _t0 = time.time() try: import json as _json @@ -439,21 +439,19 @@ def _analyze_single_image( {"role": "system", "content": system}, {"role": "user", "content": user}, ] - _orig_retries = _client.max_retries - _client.max_retries = 0 - try: - raw = _client.vision_completion( - messages=_messages, - images=[img_url], - temperature=0.3, - max_tokens=1200, # #2188: 结构化 XML 输出 600-900 字足够 - timeout=tmo, - model=model, - ) - finally: - _client.max_retries = _orig_retries + raw = _client.vision_completion( + messages=_messages, + images=[img_url], + temperature=0.3, + max_tokens=1200, + timeout=tmo, + model=model, + ) _elapsed = time.time() - _t0 - logger.info("[爆款视频] 图片 #%d VLM(%s) 完成 elapsed=%.1fs timeout=%d", idx, model, _elapsed, tmo) + logger.info( + "[爆款视频] 图片 #%d VLM(%s/%s) 完成 elapsed=%.1fs timeout=%d", + idx, label, model, _elapsed, tmo, + ) if raw is None: return None stripped = raw.strip() @@ -464,10 +462,13 @@ def _analyze_single_image( try: return _json.loads(stripped) except (_json.JSONDecodeError, TypeError): - return raw + return stripped except Exception as e: _elapsed = time.time() - _t0 - logger.warning("[爆款视频] 图片 #%d call_vision(%s) 异常 elapsed=%.1fs err=%s", idx, model, _elapsed, e) + logger.warning( + "[爆款视频] 图片 #%d call_vision(%s/%s) 异常 elapsed=%.1fs err=%s", + idx, label, model, _elapsed, e, + ) return None def _xml_to_product(nodes: list, raw_text: str) -> dict: @@ -758,30 +759,61 @@ def _analyze_single_image( product["raw"] = raw[:500] return product - first_raw = _call(vision_model, timeout) - tag1 = vision_model.split("/")[-1] if "/" in vision_model else vision_model - first_result = _normalize(first_raw, tag1) - if _is_vision_result_usable(first_result): - return first_result - - if pro_fallback_model and pro_fallback_model != vision_model: - pro_raw = _call(pro_fallback_model, 60) # #2194b: pro 单次 60s 封顶,禁用重试 - pro_result = _normalize(pro_raw, "pro_fallback") - if _is_vision_result_usable(pro_result): - pro_result["_fallback_used"] = True - return pro_result - return pro_result - return first_result + # #2198: lite/pro 并行竞速。同时发两个请求,先返回 usable 结果就用哪个,避免 + # 串行 lite超时→再发pro 累计80-100s的惩罚。外层 max_workers=2 图片并发时,竞速模式下 + # VLM 总并发=4(2图 × 2模型),实测 Ark 可以承受,且因为取快者而不是等两个都完, + # 单图通常 40-50s 就能拿到 pro 结果(pro 正常 42-46s),lite 偶发 30s 内返回时更快。 + race_t0 = time.time() + lite_tag = vision_model.split("/")[-1] if "/" in vision_model else vision_model + winner: dict | None = None + with ThreadPoolExecutor(max_workers=2) as _inner_pool: + f_lite = _inner_pool.submit(_call, vision_model, timeout, "lite") + # pro 给 75s(原60s太紧实测1/3超时,pro正常42-65s给10s余量) + pro_tmo = 75 + f_pro = _inner_pool.submit(_call, pro_fallback_model or vision_model, pro_tmo, "pro") + _fmap = {f_lite: ("lite", lite_tag), f_pro: ("pro", "pro_fallback")} + for _fut in as_completed(_fmap, timeout=pro_tmo + 15): + _lbl, _tag = _fmap[_fut] + try: + _raw = _fut.result() + except Exception as _e: + logger.warning("[爆款视频] 图片 #%d %s future异常: %s", idx, _lbl, _e) + _raw = None + _res = _normalize(_raw, _tag) + if _is_vision_result_usable(_res): + winner = _res + if _lbl == "pro": + winner["_fallback_used"] = True + logger.info( + "[爆款视频] 图片 #%d 竞速胜出=%s elapsed=%.1fs", + idx, _lbl, time.time() - race_t0, + ) + break + if winner is not None: + return winner + # 两个都失败,返回最后一次 _normalize 结果(通常是 pro 的失败 fallback,含 _source=pro_fallback_none) + try: + _last_raw = f_pro.result(timeout=1) + except Exception: + _last_raw = None + _last = _normalize(_last_raw, "pro_fallback") + logger.warning( + "[爆款视频] 图片 #%d lite/pro 竞速均失败 elapsed=%.1fs", idx, time.time() - race_t0, + ) + return _last def _step_image_analysis(job: ViralVideoJob) -> dict: - """步骤 1: 图片 VLM 分析 — 识别产品特征(v1.6 优化:并行 + lite 模型提速)。 - #2188/#2194: (1) 所有图片 URL 先归一化(storage_key→公网URL+空值报400) + """步骤 1: 图片 VLM 分析 — 识别产品特征(v1.6/#2198 优化:lite/pro 并行竞速)。 + #2188/#2194/#2198: (1) 所有图片 URL 先归一化(storage_key→公网URL+空值报400) (2) 爆款视频强制 lite-first,不依赖 .env USE_LITE 开关 - (3) max_tokens=1200,max_workers=min(2,n) 防方舟限流 - (4) lite 单次40s封顶、pro单次60s封顶,底层httpx重试关闭, - 单图最坏 40+60=100s,3图2并发最坏约100s(含排队),比240s改善60%+ - (5) 每张图 VLM 调用结束打印 elapsed 耗时日志便于排查 + (3) max_tokens=1200,max_workers=min(2,n) 防方舟限流(竞速模式总并发=4) + (4) lite/pro 并行竞速:单张图同时发 lite(30s) 和 pro(75s), + 谁先返回 usable 结果就用谁。单图最坏 75s(pro慢),典型 40-50s, + 3图2并发最坏约75s,比原串行 lite→pro 240s 改善70%+ + (5) 整个阶段统一关闭底层 httpx 重试(外层 max_retries=0,finally 恢复), + 子线程只读不改 client 属性避免竞态 + (6) 每张图 VLM 调用结束打印 elapsed 耗时日志便于排查 """ try: from packages.shared.ai_service import call_vision # noqa: F401 @@ -803,7 +835,7 @@ def _step_image_analysis(job: ViralVideoJob) -> dict: logger.error("[爆款视频] 图片 #%d URL 归一化失败: %s", idx, _ve) raise # 上层 celery 捕获后标记任务失败,避免"未识别·无法判断"误导 - # #2188 BUG2: 爆款视频强制 lite-first(不依赖 .env 开关),lite timeout=30s,pro fallback 90s + # #2188/#2198 BUG2: 爆款视频强制 lite-first(不依赖 .env 开关),lite/pro 并行竞速 try: _s = get_shared_settings() lite_model = _s.doubao_vision_lite_model @@ -811,34 +843,41 @@ def _step_image_analysis(job: ViralVideoJob) -> dict: except Exception: lite_model = "doubao-seed-2-1-lite-260915" pro_model = "doubao-seed-2-1-pro-260915" - vision_model = lite_model # 永远 lite 主跑 - vision_timeout = 40 # #2194b: lite 单次 40s 封顶(25s太紧偶发误判未识别),超时降级 pro + vision_model = lite_model + # #2198: lite 单次 30s 封顶(竞速快速路径,30s 还没出就等 pro),pro 75s(在 _analyze_single_image + # 的内部竞速池里设置),外层不感知。单图最坏 75s(仅 pro 成功),典型 40-50s(pro 正常返回)。 + vision_timeout = 30 + + # #2194/#2198: 整个并行图片分析阶段统一把共享 client 的 max_retries 置 0, + # 阶段结束 finally 恢复。子线程 _call 只读不改,避免竞态。 + from packages.shared.ai_client import get_doubao_client as _gdc_step + _step_client = _gdc_step() + _step_orig_retries = _step_client.max_retries + _step_client.max_retries = 0 results: list[dict] = [None] * len(normalized_urls) # type: ignore - max_workers = min(2, max(1, len(normalized_urls))) # #2188: 并发≤2 防方舟限流 + max_workers = min(2, max(1, len(normalized_urls))) # 并发≤2 防方舟限流(竞速模式下总并发=4) logger.info( - "[爆款视频] 开始并行图片分析 n=%d model=%s pro_fallback=%s lite_timeout=%d pro_timeout=%d workers=%d", - len(normalized_urls), - vision_model, - pro_model, - vision_timeout, - 60, - max_workers, + "[爆款视频] 开始并行竞速图片分析 n=%d lite=%s(%ds) pro=%s(75s) img_workers=%d", + len(normalized_urls), vision_model, vision_timeout, pro_model, max_workers, ) - with ThreadPoolExecutor(max_workers=max_workers) as pool: - future_to_idx = { - pool.submit( - _analyze_single_image, idx, url, vision_model, vision_timeout, pro_fallback_model=pro_model - ): idx - for idx, url in enumerate(normalized_urls) - } - for fut in as_completed(future_to_idx): - idx = future_to_idx[fut] - try: - results[idx] = fut.result() - except Exception as e: - logger.warning("[爆款视频] 图片 #%d future 异常 err=%s", idx, e, exc_info=True) - results[idx] = _vision_fallback(idx, "future_exception", {"_error": str(e)[:200]}) + try: + with ThreadPoolExecutor(max_workers=max_workers) as pool: + future_to_idx = { + pool.submit( + _analyze_single_image, idx, url, vision_model, vision_timeout, pro_fallback_model=pro_model + ): idx + for idx, url in enumerate(normalized_urls) + } + for fut in as_completed(future_to_idx): + idx = future_to_idx[fut] + try: + results[idx] = fut.result() + except Exception as e: + logger.warning("[爆款视频] 图片 #%d future 异常 err=%s", idx, e, exc_info=True) + results[idx] = _vision_fallback(idx, "future_exception", {"_error": str(e)[:200]}) + finally: + _step_client.max_retries = _step_orig_retries return {"products": results}