fix(viral_video): #2198 lite/pro并行竞速,单图最坏75s(原串行100s+)
CI/CD Pipeline / Check if frontend-only change (push) Has been skipped
CI/CD Pipeline / Dedup Check - skip PR tests when covered by push pipeline (push) Successful in 2s
CI/CD Pipeline / Staging API Integration Tests (push) Has been cancelled
CI/CD Pipeline / Build Production API Image (push) Has been cancelled
CI/CD Pipeline / Build Production Web Image (push) Has been cancelled
CI/CD Pipeline / Build Production Worker Image (push) Has been cancelled
CI/CD Pipeline / Deploy Production (push) Has been cancelled
CI/CD Pipeline / Validate - Security (push) Has been cancelled
CI/CD Pipeline / Unit Tests (push) Has been cancelled
CI/CD Pipeline / Build Staging API Image (push) Has been cancelled
CI/CD Pipeline / Build Staging Web Image (push) Has been cancelled
CI/CD Pipeline / Build Staging Worker Image (push) Has been cancelled
CI/CD Pipeline / Retag skipped Staging API Image (push) Has been cancelled
CI/CD Pipeline / Check push changed paths (push) Has been cancelled
CI/CD Pipeline / Retag skipped Staging Web Image (push) Has been cancelled
CI/CD Pipeline / Retag skipped Staging Worker Image (push) Has been cancelled
CI/CD Pipeline / Deploy Staging (Watchtower auto-deploy) (push) Has been cancelled
CI/CD Pipeline / Staging E2E Tests (push) Has been cancelled
CI/CD Pipeline / Production Browser E2E (push) Has been cancelled
CI/CD Pipeline / ACR Image Cleanup (push) Has been cancelled
CI/CD Pipeline / Canary Release to Production (push) Has been cancelled
CI/CD Pipeline / CI Gate (push) Has been cancelled
CI/CD Pipeline / Validate - Style (push) Has been cancelled
CI/CD Pipeline / Validate - Python (mypy + alembic) (push) Has been cancelled
CI/CD Pipeline / Integration Tests (push) Has been cancelled
CI/CD Pipeline / Frontend Lint (push) Has been cancelled
CI/CD Pipeline / Frontend Unit Tests (push) Has been cancelled
CI/CD Pipeline / PR Build API Image (push) Has been cancelled
CI/CD Pipeline / PR Build Web Image (push) Has been cancelled
CI/CD Pipeline / PR Build Worker Image (push) Has been cancelled
CI/CD Pipeline / Staging E2E Tests (pull_request) Has been skipped
ACR Cleanup / ACR Image Cleanup (pull_request_target) Successful in 38s
CI/CD Pipeline / Staging API Integration Tests (pull_request) Has been skipped
CI/CD Pipeline / ACR Image Cleanup (pull_request) Has been skipped
CI/CD Pipeline / Check if frontend-only change (pull_request) Successful in 2s
CI/CD Pipeline / Check push changed paths (pull_request) Has been skipped
CI/CD Pipeline / Build Staging Worker Image (pull_request) Has been skipped
CI/CD Pipeline / Build Staging API Image (pull_request) Has been skipped
CI/CD Pipeline / Build Staging Web Image (pull_request) Has been skipped
CI/CD Pipeline / Retag skipped Staging Web Image (pull_request) Has been skipped
CI/CD Pipeline / Retag skipped Staging API Image (pull_request) Has been skipped
CI/CD Pipeline / Retag skipped Staging Worker Image (pull_request) Has been skipped
CI/CD Pipeline / Deploy Staging (Watchtower auto-deploy) (pull_request) Has been skipped
CI/CD Pipeline / Dedup Check - skip PR tests when covered by push pipeline (pull_request) Successful in 1s
CI/CD Pipeline / Frontend Unit Tests (pull_request) Has been skipped
CI/CD Pipeline / PR Build Web Image (pull_request) Has been skipped
CI/CD Pipeline / Frontend Lint (pull_request) Has been skipped
Preview Cleanup / Cleanup Preview Environment (pull_request) Successful in 59s
CI/CD Pipeline / PR Build Worker Image (pull_request) Successful in 1m22s
Preview Deploy / Deploy Preview Environment (pull_request) Successful in 2m4s
CI/CD Pipeline / Unit Tests (pull_request) Failing after 3m5s
PR Automation / Auto Approve on CI Green (pull_request) Successful in 3m15s
CI/CD Pipeline / PR Build API Image (pull_request) Successful in 3m53s
CI/CD Pipeline / Integration Tests (pull_request) Successful in 3m55s
CI/CD Pipeline / Validate - Python (mypy + alembic) (pull_request) Successful in 4m7s
CI/CD Pipeline / Validate - Style (pull_request) Failing after 6m5s
PR Automation / Auto Merge on CI Green + Approved (pull_request) Successful in 1m26s
CI/CD Pipeline / Validate - Security (pull_request) Successful in 9m38s
AI Code Review / AI Code Review (pull_request) Successful in 10m17s
CI/CD Pipeline / Build Production API Image (pull_request) Has been skipped
CI/CD Pipeline / Build Production Web Image (pull_request) Has been skipped
CI/CD Pipeline / Build Production Worker Image (pull_request) Has been skipped
CI/CD Pipeline / CI Gate (pull_request) Failing after 3s
CI/CD Pipeline / Canary Release to Production (pull_request) Has been cancelled
CI/CD Pipeline / Deploy Production (pull_request) Has been skipped
CI/CD Pipeline / Production Browser E2E (pull_request) Has been skipped

问题:实测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)便于区分竞速胜出方
This commit was merged in pull request #2198.
This commit is contained in:
xiaoxia
2026-10-05 16:08:00 +08:00
parent 305e2bd9d5
commit 9699a1fcde
2 changed files with 104 additions and 64 deletions
+1
View File
@@ -1,2 +1,3 @@
- 2026-10-05 #2194 VLM timeout tune + #2195 HEAD→GET Range fix deployed to staging - 2026-10-05 #2194 VLM timeout tune + #2195 HEAD→GET Range fix deployed to staging
# 2198 lite/pro并行竞速
+103 -64
View File
@@ -425,9 +425,9 @@ def _analyze_single_image(
image_urls=f"第1张:{img_url}", image_urls=f"第1张:{img_url}",
) )
def _call(model: str, tmo: int): def _call(model: str, tmo: int, label: str):
# #2194: 单次调用临时关闭 httpx 层重试,超时/失败由外层 pro_fallback 统一兜底, # #2194/#2198: max_retries=0 由外层 _step_image_analysis 统一设置(阶段前置0、阶段后恢复),
# 避免底层 max_retries=1 导致 lite timeout × 2 + pro timeout × 2 最坏 240s+ # 子线程只读不改,避免嵌套并行竞速时多线程同时改 client.max_retries 产生竞态
_t0 = time.time() _t0 = time.time()
try: try:
import json as _json import json as _json
@@ -439,21 +439,19 @@ def _analyze_single_image(
{"role": "system", "content": system}, {"role": "system", "content": system},
{"role": "user", "content": user}, {"role": "user", "content": user},
] ]
_orig_retries = _client.max_retries raw = _client.vision_completion(
_client.max_retries = 0 messages=_messages,
try: images=[img_url],
raw = _client.vision_completion( temperature=0.3,
messages=_messages, max_tokens=1200,
images=[img_url], timeout=tmo,
temperature=0.3, model=model,
max_tokens=1200, # #2188: 结构化 XML 输出 600-900 字足够 )
timeout=tmo,
model=model,
)
finally:
_client.max_retries = _orig_retries
_elapsed = time.time() - _t0 _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: if raw is None:
return None return None
stripped = raw.strip() stripped = raw.strip()
@@ -464,10 +462,13 @@ def _analyze_single_image(
try: try:
return _json.loads(stripped) return _json.loads(stripped)
except (_json.JSONDecodeError, TypeError): except (_json.JSONDecodeError, TypeError):
return raw return stripped
except Exception as e: except Exception as e:
_elapsed = time.time() - _t0 _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 return None
def _xml_to_product(nodes: list, raw_text: str) -> dict: def _xml_to_product(nodes: list, raw_text: str) -> dict:
@@ -758,30 +759,61 @@ def _analyze_single_image(
product["raw"] = raw[:500] product["raw"] = raw[:500]
return product return product
first_raw = _call(vision_model, timeout) # #2198: lite/pro 并行竞速。同时发两个请求,先返回 usable 结果就用哪个,避免
tag1 = vision_model.split("/")[-1] if "/" in vision_model else vision_model # 串行 lite超时→再发pro 累计80-100s的惩罚。外层 max_workers=2 图片并发时,竞速模式下
first_result = _normalize(first_raw, tag1) # VLM 总并发=4(2图 × 2模型),实测 Ark 可以承受,且因为取快者而不是等两个都完,
if _is_vision_result_usable(first_result): # 单图通常 40-50s 就能拿到 pro 结果(pro 正常 42-46s),lite 偶发 30s 内返回时更快。
return first_result race_t0 = time.time()
lite_tag = vision_model.split("/")[-1] if "/" in vision_model else vision_model
if pro_fallback_model and pro_fallback_model != vision_model: winner: dict | None = None
pro_raw = _call(pro_fallback_model, 60) # #2194b: pro 单次 60s 封顶,禁用重试 with ThreadPoolExecutor(max_workers=2) as _inner_pool:
pro_result = _normalize(pro_raw, "pro_fallback") f_lite = _inner_pool.submit(_call, vision_model, timeout, "lite")
if _is_vision_result_usable(pro_result): # pro 给 75s(原60s太紧实测1/3超时,pro正常42-65s给10s余量)
pro_result["_fallback_used"] = True pro_tmo = 75
return pro_result f_pro = _inner_pool.submit(_call, pro_fallback_model or vision_model, pro_tmo, "pro")
return pro_result _fmap = {f_lite: ("lite", lite_tag), f_pro: ("pro", "pro_fallback")}
return first_result 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: def _step_image_analysis(job: ViralVideoJob) -> dict:
"""步骤 1: 图片 VLM 分析 — 识别产品特征(v1.6 优化:并行 + lite 模型提速)。 """步骤 1: 图片 VLM 分析 — 识别产品特征(v1.6/#2198 优化:lite/pro 并行竞速)。
#2188/#2194: (1) 所有图片 URL 先归一化(storage_key→公网URL+空值报400) #2188/#2194/#2198: (1) 所有图片 URL 先归一化(storage_key→公网URL+空值报400)
(2) 爆款视频强制 lite-first,不依赖 .env USE_LITE 开关 (2) 爆款视频强制 lite-first,不依赖 .env USE_LITE 开关
(3) max_tokens=1200,max_workers=min(2,n) 防方舟限流 (3) max_tokens=1200,max_workers=min(2,n) 防方舟限流(竞速模式总并发=4)
(4) lite 单次40s封顶、pro单次60s封顶,底层httpx重试关闭, (4) lite/pro 并行竞速:单张图同时发 lite(30s) 和 pro(75s),
单图最坏 40+60=100s,3图2并发最坏约100s(含排队),比240s改善60%+ 谁先返回 usable 结果就用谁。单图最坏 75s(pro慢),典型 40-50s,
(5) 每张图 VLM 调用结束打印 elapsed 耗时日志便于排查 3图2并发最坏约75s,比原串行 lite→pro 240s 改善70%+
(5) 整个阶段统一关闭底层 httpx 重试(外层 max_retries=0,finally 恢复),
子线程只读不改 client 属性避免竞态
(6) 每张图 VLM 调用结束打印 elapsed 耗时日志便于排查
""" """
try: try:
from packages.shared.ai_service import call_vision # noqa: F401 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) logger.error("[爆款视频] 图片 #%d URL 归一化失败: %s", idx, _ve)
raise # 上层 celery 捕获后标记任务失败,避免"未识别·无法判断"误导 raise # 上层 celery 捕获后标记任务失败,避免"未识别·无法判断"误导
# #2188 BUG2: 爆款视频强制 lite-first(不依赖 .env 开关),lite timeout=30s,pro fallback 90s # #2188/#2198 BUG2: 爆款视频强制 lite-first(不依赖 .env 开关),lite/pro 并行竞速
try: try:
_s = get_shared_settings() _s = get_shared_settings()
lite_model = _s.doubao_vision_lite_model lite_model = _s.doubao_vision_lite_model
@@ -811,34 +843,41 @@ def _step_image_analysis(job: ViralVideoJob) -> dict:
except Exception: except Exception:
lite_model = "doubao-seed-2-1-lite-260915" lite_model = "doubao-seed-2-1-lite-260915"
pro_model = "doubao-seed-2-1-pro-260915" pro_model = "doubao-seed-2-1-pro-260915"
vision_model = lite_model # 永远 lite 主跑 vision_model = lite_model
vision_timeout = 40 # #2194b: lite 单次 40s 封顶(25s太紧偶发误判未识别),超时降级 pro # #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 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( logger.info(
"[爆款视频] 开始并行图片分析 n=%d model=%s pro_fallback=%s lite_timeout=%d pro_timeout=%d workers=%d", "[爆款视频] 开始并行竞速图片分析 n=%d lite=%s(%ds) pro=%s(75s) img_workers=%d",
len(normalized_urls), len(normalized_urls), vision_model, vision_timeout, pro_model, max_workers,
vision_model,
pro_model,
vision_timeout,
60,
max_workers,
) )
with ThreadPoolExecutor(max_workers=max_workers) as pool: try:
future_to_idx = { with ThreadPoolExecutor(max_workers=max_workers) as pool:
pool.submit( future_to_idx = {
_analyze_single_image, idx, url, vision_model, vision_timeout, pro_fallback_model=pro_model pool.submit(
): idx _analyze_single_image, idx, url, vision_model, vision_timeout, pro_fallback_model=pro_model
for idx, url in enumerate(normalized_urls) ): idx
} for idx, url in enumerate(normalized_urls)
for fut in as_completed(future_to_idx): }
idx = future_to_idx[fut] for fut in as_completed(future_to_idx):
try: idx = future_to_idx[fut]
results[idx] = fut.result() try:
except Exception as e: results[idx] = fut.result()
logger.warning("[爆款视频] 图片 #%d future 异常 err=%s", idx, e, exc_info=True) except Exception as e:
results[idx] = _vision_fallback(idx, "future_exception", {"_error": str(e)[:200]}) 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} return {"products": results}