diff --git a/apps/worker/worker_app/tasks/viral_video.py b/apps/worker/worker_app/tasks/viral_video.py
index 340f3cf76..6c652ca78 100644
--- a/apps/worker/worker_app/tasks/viral_video.py
+++ b/apps/worker/worker_app/tasks/viral_video.py
@@ -1,7 +1,9 @@
-"""爆款视频 Celery 编排器 — ViralVideoOrchestrator (v1.6 单次 Seedance 出片版).
+"""爆款视频 Celery 编排器 — ViralVideoOrchestrator.
-v1.6 重大简化(Seedance 2.5 单次最长 30 秒,直接出片):
- 1. _step_image_analysis 图片 VLM 分析(保留)
+V2 图片分析(10-05):火山OCR专用API + doubao-lite强约束JSON并行,单图<3s,8图<15s;pro VLM单次兜底。输出字段兼容旧格式,下游信任链/t2i零改动。
+
+流水线步骤:
+ 1. _step_image_analysis 图片分析(V2: OCR+lite VLM并行 + pro兜底)
1.5 _step_video_analysis 参考视频风格分析(可选)
2. _step_intent_parsing 用户文案意图解析
3. _step_script_generation 编导分镜脚本生成(融合原 copy_fusion+storyboard+review,输出 copy_result 结构 + voiceover_script)
@@ -26,7 +28,6 @@ import re
import tempfile
import threading
import time
-from concurrent.futures import ThreadPoolExecutor, as_completed
from pathlib import Path
from typing import Any
@@ -330,633 +331,78 @@ def _vision_fallback(idx: int, reason: str, extra: dict | None = None) -> dict:
return d
-def _is_vision_result_usable(result: dict) -> bool:
- """判断 VLM 返回是否有效。
- #2198b: 判定条件放宽——有人像(portrait_prompt非空/非'无人像')即视为usable(爆款视频核心
- 是要人物描述给信任链t2i用,product name/brand/features识别不准是次要的)。
- 非人像场景才要求name+summary+features有效。
- """
- if not isinstance(result, dict):
- return False
- # 有人像描述(爆款视频最核心需求,portrait_prompt给信任链t2i做参考)就视为usable
- pp = (result.get("portrait_prompt") or "").strip()
- if pp and pp not in ("无人像", "无法判断", "未识别"):
- return True
- # 非人像场景:要求name+summary有效
- name = (result.get("name") or "").strip()
- if not name or name in ("未识别", "无法判断", "未知"):
- return False
- summary = (result.get("summary") or "").strip()
- if len(summary) < 5 or summary in ("无法判断", "未识别"):
- return False
- category = (result.get("category") or "").strip()
- if category == "非产品图":
- return True
- feats = result.get("key_features") or []
- if not isinstance(feats, list) or len(feats) == 0 or feats == ["无法判断"]:
- return False
- return True
-
-
def _normalize_image_url(raw: str, idx: int) -> str:
- """#2188: 将 job.images 中的 storage_key/相对路径/空值统一归一化为可公网访问 URL。
- - 以 http:// 或 https:// 开头 → 视为公网 URL
- - 其他 → 视为 storage_key,用 SharedStorageService.get_url() 转公网 URL
- - 空值/None/非字符串 → 抛 ValueError(上层 catch 后走 400 错误)
- 返回前做 HTTP 可达性检查(GET+Range:0-1024 避免 OSS 签名 URL 对 HEAD 返回 403 的假阴性)。
+ """将 job.images 中的 storage_key/相对路径/空值统一归一化为可公网访问 URL。
+ - http(s):// → 直接用
+ - 其他 → storage_key,通过 SharedStorageService.get_url() 转公网 URL
+ - 空值/非字符串 → 抛 ValueError
"""
- import requests as _req
-
if not raw or not isinstance(raw, str):
raise ValueError(f"图片 #{idx} URL 为空或类型错误: {type(raw).__name__}={raw!r}")
url = raw.strip()
if not url:
raise ValueError(f"图片 #{idx} URL 为空白字符串")
- # storage_key 判定:不以 http 开头
- if not url.startswith("http://") and not url.startswith("https://"):
- # 去掉可能的前导斜杠
- storage_key = url.lstrip("/")
- try:
- from packages.shared.storage import get_storage_service
-
- _svc = get_storage_service()
- url = _svc.get_url(storage_key)
- except Exception as _e:
- raise ValueError(f"图片 #{idx} storage_key={storage_key!r} 转公网URL失败: {_e}") from _e
- logger.info("[爆款视频] 图片 #%d storage_key 已转公网 URL: %s", idx, url[:120])
- # #2194: 用 GET+Range 代替 HEAD。
- # Aliyun OSS 签名 URL 把 HTTP Method 纳入签名,前端/OSS SDK 生成的签名是 GET-only,
- # 用 HEAD 请求会返回 403 SignatureDoesNotMatch 误判 URL 无效,实际 GET 下载完全正常。
- # Range: bytes=0-1024 只取前1KB,开销极小。
+ if url.startswith("http://") or url.startswith("https://"):
+ return url
+ storage_key = url.lstrip("/")
try:
- _r = _req.get(url, timeout=5, allow_redirects=True, stream=True, headers={"Range": "bytes=0-1024"})
- if _r.status_code >= 400:
- logger.warning("[爆款视频] 图片 #%d URL 可达性检查返回 %d: %s", idx, _r.status_code, url[:120])
- _r.close()
+ from packages.shared.storage import get_storage_service
+ url = get_storage_service().get_url(storage_key)
except Exception as _e:
- logger.warning("[爆款视频] 图片 #%d URL 可达性检查异常: %s url=%s", idx, _e, url[:120])
+ raise ValueError(f"图片 #{idx} storage_key={storage_key!r} 转公网URL失败: {_e}") from _e
+ logger.info("[爆款视频] 图片 #%d storage_key → 公网URL: %s", idx, url[:120])
return url
-def _analyze_single_image(
- idx: int,
- img_url: str,
- vision_model: str,
- timeout: int,
- *,
- pro_fallback_model: str | None = None,
-) -> dict:
- """单张图片 VLM 分析(#2040:改为从 prompt_loader 读模板 + XML 解析)。
-
- lite 失败/不可用时用 pro 降级重试 1 次。失败/None 最终返回含默认字段的 dict。
- """
- try:
- from packages.application.viral_video import xml_parser as xp
- from packages.application.viral_video.prompt_loader import (
- get_template,
- render_system_prompt,
- render_user_prompt,
- )
- from packages.shared.ai_service import call_vision
- except ImportError as e:
- logger.warning("[爆款视频] prompt 模板/解析模块不可用: %s", e)
- return _vision_fallback(idx, f"fallback_import_error:{e}")
-
- if not img_url or not isinstance(img_url, str):
- return _vision_fallback(idx, "invalid_url")
-
- template = get_template("image_analysis")
- system = render_system_prompt(template)
- user = render_user_prompt(
- template,
- image_count=1,
- industry="通用",
- image_urls=f"第1张:{img_url}",
- )
-
- 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
-
- from packages.shared.ai_client import get_doubao_client as _gdc
-
- _client = _gdc()
- _messages = [
- {"role": "system", "content": system},
- {"role": "user", "content": user},
- ]
- 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/%s) 完成 elapsed=%.1fs timeout=%d",
- idx,
- label,
- model,
- _elapsed,
- tmo,
- )
- if raw is None:
- return None
- stripped = raw.strip()
- if stripped.startswith("```"):
- stripped = stripped.strip("`")
- if stripped.startswith("json"):
- stripped = stripped[4:].lstrip()
- try:
- return _json.loads(stripped)
- except (_json.JSONDecodeError, TypeError):
- return stripped
- except Exception as e:
- _elapsed = time.time() - _t0
- 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:
- product_nodes = [n for n in nodes if n["tag"] == "product"]
- scene = xp.text_of(raw_text, "scene") or "通用"
- mood = xp.text_of(raw_text, "mood") or ""
- # #2184: #2177 XML 重构后人物信息放在顶层 ,
- # 不再是 的 portrait_prompt 属性。需从顶层 people 标签提取并拼装 portrait_prompt。
- portrait_prompt = "无人像"
- try:
- people_node = xp.find_first(raw_text, "people")
- if people_node:
- _pa = people_node.get("attrs") or {}
- _has_person = xp.attr_bool(_pa.get("has_person"), False)
- if _has_person:
- _gender = _pa.get("gender", "无法判断") or "无法判断"
- _age = _pa.get("age_range", "无法判断") or "无法判断"
- _hair = _pa.get("hair", "无法判断") or "无法判断"
- _skin = _pa.get("skin_tone", "无法判断") or "无法判断"
- _face = _pa.get("face_shape", "无法判断") or "无法判断"
- _outfit = _pa.get("outfit", "无法判断") or "无法判断"
- _pose = _pa.get("pose", "无法判断") or "无法判断"
- _expr = _pa.get("expression", "无法判断") or "无法判断"
- _count = xp.attr_int(_pa.get("count"), 1)
- # #2185: VLM有时对外貌属性输出"无法判断",用通用兜底值确保portrait_prompt始终有完整外貌描述
- if _hair == "无法判断":
- _hair = "自然发型"
- if _skin == "无法判断":
- _skin = "自然"
- if _face == "无法判断":
- _face = "标准"
- if _outfit == "无法判断":
- _outfit = "日常服装"
- _parts = []
- if _gender != "无法判断":
- _g = _gender + ("性" if not _gender.endswith("性") else "")
- _parts.append(_g)
- else:
- _parts.append("成年人")
- if _age != "无法判断":
- _parts.append(_age)
- _parts.append("人物")
- _parts.append(_hair)
- _parts.append(f"{_skin}肤色")
- _parts.append(f"{_face}脸型")
- _parts.append(f"身着{_outfit}")
- if _pose != "无法判断":
- _parts.append(f"姿态{_pose}")
- if _expr != "无法判断":
- _parts.append(f"表情{_expr}")
- else:
- _parts.append("表情自然")
- # #2186: 智能回填——VLM有时省略hair/outfit等外貌属性,但product.name/features/colors里已有相关信息
- # 从product名字和features中提取服装关键词回填outfit
- if _outfit in ("日常服装", "无法判断"):
- for _ppn in product_nodes:
- _pn = (_ppn.get("attrs") or {}).get("name", "") or ""
- _pf = (_ppn.get("attrs") or {}).get("features", "") or ""
- _ptxt = _pn + " " + _pf
- # 服装关键词识别(常见上装/下装/裙装/套装)
- _cloth_kws = [
- # 衬衫/T恤类
- "衬衫",
- "T恤",
- "POLO衫",
- "polo衫",
- "Polo衫",
- "打底衫",
- "雪纺衫",
- "罩衫",
- "针织衫",
- # 毛衣/卫衣/针织类
- "毛衣",
- "卫衣",
- "帽衫",
- "针织",
- "毛衫",
- "开衫",
- # 外套/西装/夹克/风衣类
- "外套",
- "西装",
- "西服",
- "夹克",
- "皮衣",
- "皮夹克",
- "风衣",
- "大衣",
- "羽绒服",
- "棉服",
- "棉服",
- "马甲",
- "背心",
- "开衫外套",
- # 裙装
- "连衣裙",
- "半身裙",
- "短裙",
- "长裙",
- "百褶裙",
- "A字裙",
- "旗袍",
- "汉服",
- "JK裙",
- # 裤装
- "牛仔裤",
- "休闲裤",
- "西裤",
- "运动裤",
- "短裤",
- "阔腿裤",
- "打底裤",
- # 制服/套装
- "制服",
- "套装",
- "职业装",
- "工装",
- # 通用上装/下装词(兜底)
- "上衣",
- "短袖",
- "长袖",
- "无袖",
- "半袖",
- "吊带",
- "背心",
- "网纱",
- "雪纺",
- "真丝",
- "纯棉",
- "亚麻",
- ]
- for _ckw in _cloth_kws:
- if _ckw in _ptxt:
- _ci = _ptxt.find(_ckw)
- # 向前找颜色/材质/款式形容词(白/黑/米/红/蓝/灰/棉/麻/长/短/厚/薄/长袖/短袖/翻领/圆领/V领/印花/条纹等)
- _start = max(0, _ci - 12)
- # 向后包含款式词(长袖/短袖/外套/套装/上衣等后续修饰)
- _end = min(len(_ptxt), _ci + len(_ckw) + 8)
- _outfit_extract = _ptxt[_start:_end].strip(" ,,。.、")
- # 仅清理明确的品牌/产品类前缀(不清理颜色/款式/尺寸形容词)
- _outfit_extract = re.sub(
- r"^(\S{0,4}牌|\S{0,3}品牌|\S{0,3}款|产品|商品|的)", "", _outfit_extract
- ).strip()
- # 尾部清理:去掉残留的品牌字/型号字(如"标""ml""g""装"等单字杂字)
- _outfit_extract = re.sub(
- r"(标[0-9a-zA-Z]*|\d+\s*(?:ml|g|L|斤|件|个|瓶|盒|包|袋|装)|\s+\d+\s*)$",
- "",
- _outfit_extract,
- flags=re.IGNORECASE,
- ).strip()
- if len(_outfit_extract) >= 2:
- _outfit = _outfit_extract
- break
- if _outfit not in ("日常服装", "无法判断"):
- break
- # 从color标签中提取头发颜色回填hair
- if _hair in ("自然发型", "无法判断"):
- _hair_color = ""
- _color_nodes = [n for n in nodes if n["tag"] == "color"]
- _hair_kws_map = {
- "黑": "黑色",
- "棕": "棕色",
- "金": "金色",
- "栗": "栗色",
- "红": "红色",
- "白": "白色",
- "灰": "灰色",
- "蓝": "蓝色",
- "黄": "黄色",
- "紫": "紫色",
- }
- for _cn in _color_nodes:
- _cname = (_cn.get("attrs") or {}).get("name", "") or ""
- # 小占比颜色更可能是发色(非主色的小面积色),且名称含头发/黑/棕/金等
- _ccov = 0.0
- try:
- _ccov = float((_cn.get("attrs") or {}).get("coverage", "0") or 0)
- except Exception:
- pass
- for _hk, _hv in _hair_kws_map.items():
- if _hk in _cname and _ccov < 0.3:
- _hair_color = _hv
- break
- if _hair_color:
- break
- if _hair_color:
- _hair = f"{_hair_color}头发"
- else:
- _hair = "自然发型"
- # 重新拼装_parts(回填后)
- _parts = []
- if _gender != "无法判断":
- _g = _gender + ("性" if not _gender.endswith("性") else "")
- _parts.append(_g)
- else:
- _parts.append("成年人")
- if _age != "无法判断":
- _parts.append(_age)
- _parts.append("人物")
- _parts.append(_hair)
- _parts.append(f"{_skin}肤色")
- _parts.append(f"{_face}脸型")
- _parts.append(f"身着{_outfit}")
- if _pose != "无法判断":
- _parts.append(f"姿态{_pose}")
- if _expr != "无法判断":
- _parts.append(f"表情{_expr}")
- else:
- _parts.append("表情自然")
- portrait_prompt = ",".join(_parts)
- logger.info(
- "[爆款视频] 图片 #%d 解析(回填后): count=%d gender=%s age=%s hair=%s skin=%s face=%s outfit=%s pose=%s expr=%s → %s",
- idx,
- _count,
- _gender,
- _age,
- _hair,
- _skin,
- _face,
- _outfit,
- _pose,
- _expr,
- portrait_prompt,
- )
- except Exception as _pe:
- logger.warning("[爆款视频] 图片 #%d 解析标签异常: %s,回退无人像", idx, _pe)
- for p in product_nodes:
- a = p["attrs"]
- text_on_pkg = a.get("text_on_package", "")
- p_body = p.get("text", "") or ""
- if not text_on_pkg and p_body:
- text_on_pkg = xp.text_of(p_body, "text_on_package") or ""
- text_list = [x.strip() for x in re.split(r"[,,;;]", text_on_pkg) if x.strip()] if text_on_pkg else []
- features = a.get("features", "")
- feat_list = [x.strip() for x in re.split(r"[,,;;]", features) if x.strip()] if features else []
- name = a.get("name", "") or "未识别"
- brand = a.get("brand", "") or "无法判断"
- category = a.get("category", "") or "无法判断"
- appearance = a.get("appearance", "") or "无法判断"
- packaging = a.get("packaging", "") or "无法判断"
- summary = a.get("summary", "") or f"{brand} {name}"
- # 优先取 product 属性上的 portrait_prompt(兼容旧schema),否则用顶层 解析结果
- _pp_from_attr = a.get("portrait_prompt", "")
- if _pp_from_attr and _pp_from_attr != "无人像":
- portrait_prompt = _pp_from_attr
- return {
- "name": name,
- "brand": brand,
- "category": category,
- "appearance": appearance,
- "packaging": packaging,
- "text_on_package": text_list,
- "key_features": feat_list or [features] if features else ["无法判断"],
- "scene": scene,
- "mood": mood,
- "portrait_prompt": portrait_prompt,
- "summary": summary,
- "_source": "xml",
- }
- # 没有 product 标签但有 也要能取到人物描述(兜底)
- if portrait_prompt != "无人像":
- return {
- "name": "未识别",
- "brand": "无法判断",
- "category": "无法判断",
- "appearance": "无法判断",
- "packaging": "无法判断",
- "text_on_package": [],
- "key_features": ["无法判断"],
- "scene": scene,
- "mood": mood,
- "portrait_prompt": portrait_prompt,
- "summary": "未识别",
- "_source": "xml_no_product",
- }
- return _vision_fallback(idx, "no_product_tag")
-
- def _normalize(raw, source: str) -> dict:
- if raw is None:
- return _vision_fallback(idx, f"{source}_none")
- # VLM 偶尔直接返回 JSON 对象(不包裹```json),_call 里 json.loads 后已是 dict
- if isinstance(raw, dict):
- _prod = {
- "name": raw.get("name") or "未识别",
- "brand": raw.get("brand") or "无法判断",
- "category": raw.get("category") or "无法判断",
- "appearance": raw.get("appearance") or "无法判断",
- "packaging": raw.get("packaging") or "无法判断",
- "text_on_package": raw.get("text_on_package") or [],
- "key_features": raw.get("key_features") or raw.get("features") or ["无法判断"],
- "scene": raw.get("scene") or "通用",
- "mood": raw.get("mood") or "",
- "portrait_prompt": raw.get("portrait_prompt") or "无人像",
- "summary": raw.get("summary") or f"{raw.get('brand','')} {raw.get('name','')}",
- "_source": source,
- }
- return _prod
- if not isinstance(raw, str):
- return _vision_fallback(idx, f"{source}_badtype")
- nodes = xp.parse_tags(raw)
- if not nodes:
- # 不是 XML 也不是 dict:尝试当作纯 JSON 字符串再解析一次
- try:
- import json as _j2
-
- _jd = _j2.loads(raw)
- if isinstance(_jd, dict):
- return _normalize(_jd, source)
- except Exception:
- pass
- logger.warning("[爆款视频] 图片 #%d XML/JSON 解析都失败 source=%s raw_head=%s", idx, source, raw[:200])
- return _vision_fallback(idx, f"{source}_parse_fail", {"_raw": raw[:500]})
- product = _xml_to_product(nodes, raw)
- product.setdefault("_source", source)
- product["raw"] = raw[:500]
- return product
-
- # #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: 图片分析。
+ """步骤 1: 图片分析(V2 主路径)。
- V2(VISION_V2_ENABLED=true,10-05 新方案):
- - 每图并行 2 路:火山 MediaKit OCR(专用API)+ doubao-seed-2.1-lite 强约束 JSON
- (火山云端无人体属性/商品检测/图像标签公开 HTTP API,用 lite JSON-only VLM 弥补),
- 目标单图 <3s;
- - 外层 8 图全并发,目标 8 图 <15s;
- - 置信度低/全失败时降级 doubao-seed-2.1-pro 完整 VLM 兜底(复用旧竞速逻辑);
- - 输出 dict 格式与旧 _normalize() 完全一致,下游信任链/t2i 零改动。
-
- V1(默认,#2198/#2199 lite/pro 并行竞速):
- - 过渡版兜底,单图 lite(30s)/pro(75s) 竞速,外层 max_workers=2,典型 40-75s/图。
+ 架构:
+ - 主力:火山 MediaKit OCR(专用API)+ doubao-seed-2.1-lite 强约束 JSON(弥补火山云端缺失的
+ 人体属性/商品检测/图像标签专用HTTP API),每图2路并行,目标<3s;
+ - 外层全并发(workers=8),目标8图<15s;
+ - 兜底:fast 结果不可用时单次调用 doubao-seed-2.1-pro VLM(简单、无竞速)。
+ 输出 dict 字段(name/brand/category/appearance/key_features/scene/mood/portrait_prompt/summary/_source)
+ 与旧版格式完全一致,下游信任链/t2i/intent_parsing/script_generation 零改动。
"""
- try:
- from packages.shared.ai_service import call_vision # noqa: F401
- except ImportError:
- logger.warning("[爆款视频] ai_service.call_vision 不可用,使用占位结果")
- return {"products": [_vision_fallback(0, "fallback_import_error")]}
-
if not job.images:
logger.warning("[爆款视频] 任务无 images,跳过图片分析")
return {"products": []}
- # URL 归一化 — storage_key→公网URL + 空值报400
+ # URL 归一化(storage_key→公网URL;空值直接400)
normalized_urls: list[str] = []
for idx, raw in enumerate(job.images):
- try:
- normalized_urls.append(_normalize_image_url(raw, idx))
- except ValueError as _ve:
- logger.error("[爆款视频] 图片 #%d URL 归一化失败: %s", idx, _ve)
- raise
+ normalized_urls.append(_normalize_image_url(raw, idx))
- # 模型配置(V1/V2 共用)
try:
- _s = get_shared_settings()
- lite_model = _s.doubao_vision_lite_model
- pro_model = _s.doubao_vision_model
+ from worker_app.tasks.vision import analyze_images_v2 as _aiv2
+ except ImportError:
+ try:
+ from tasks.vision import analyze_images_v2 as _aiv2 # type: ignore
+ except ImportError as e:
+ logger.error("[爆款视频] vision 模块导入失败: %s", e)
+ return {"products": [_vision_fallback(0, f"vision_import_error:{e}")]}
+
+ # 整个阶段关闭底层 httpx 重试,避免线程里出现不可控等待
+ try:
+ from packages.shared.ai_client import get_doubao_client as _gdc
+ _cli = _gdc()
+ _orig_retries = _cli.max_retries
+ _cli.max_retries = 0
except Exception:
- lite_model = "doubao-seed-2-1-lite-260915"
- pro_model = "doubao-seed-2-1-pro-260915"
+ _cli = None
+ _orig_retries = 0
- # ========== V2 路径(VISION_V2_ENABLED=true)==========
- import os as _os_v2
-
- _v2_enabled = _os_v2.environ.get("VISION_V2_ENABLED", "false").lower() in ("1", "true", "yes", "on")
- if _v2_enabled:
- try:
- from worker_app.tasks.vision import analyze_images_v2 as _aiv2
- except ImportError:
- try:
- from tasks.vision import analyze_images_v2 as _aiv2 # type: ignore
- except ImportError:
- logger.warning("[vision.v2] 模块导入失败,回退 V1 路径")
- _aiv2 = None # type: ignore
- if _aiv2 is not None:
- from packages.shared.ai_client import get_doubao_client as _gdc_v2
-
- _cli = _gdc_v2()
- _orig_retries_v2 = _cli.max_retries
- _cli.max_retries = 0
- try:
- results_v2 = _aiv2(normalized_urls, lite_model=lite_model, pro_model=pro_model)
- finally:
- _cli.max_retries = _orig_retries_v2
- return {"products": list(results_v2)}
- # 导入失败 → fallthrough 走 V1
-
- # ========== V1 路径(默认,lite/pro 并行竞速)==========
- vision_model = lite_model
- vision_timeout = 30
-
- 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))) # 并发≤2 防方舟限流(竞速模式下总并发=4)
- logger.info(
- "[爆款视频] 开始并行竞速图片分析(V1) n=%d lite=%s(%ds) pro=%s(75s) img_workers=%d",
- len(normalized_urls),
- vision_model,
- vision_timeout,
- pro_model,
- max_workers,
- )
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]})
+ results = _aiv2(normalized_urls)
finally:
- _step_client.max_retries = _step_orig_retries
+ if _cli is not None:
+ try:
+ _cli.max_retries = _orig_retries
+ except Exception:
+ pass
- return {"products": results}
+ return {"products": list(results)}
def _step_video_analysis(job: ViralVideoJob) -> dict | None:
diff --git a/apps/worker/worker_app/tasks/vision/__init__.py b/apps/worker/worker_app/tasks/vision/__init__.py
index c4672627f..acb6960ff 100644
--- a/apps/worker/worker_app/tasks/vision/__init__.py
+++ b/apps/worker/worker_app/tasks/vision/__init__.py
@@ -1,16 +1,3 @@
# -*- coding: utf-8 -*-
-"""V2 图片分析:专用 API 组合路径(OCR + lite JSON VLM 并行 + pro VLM 兜底)。
-
-灵应指令(10-05):
-- 优先火山引擎视觉智能 API:OCR 是真实专用云端 API;
-- 人体属性/商品检测/图像标签:火山云端无公开 HTTP API(仅有移动端 SDK),
- 采用 doubao-seed-2.1-lite + 强约束 JSON-only prompt 作为"伪专用 API",
- 目标 1-3s 返回结构化字段;
-- VLM(doubao-seed-2.1-pro)保留为终极兜底(置信度低/全失败时降级);
-- 单图并行 2 路(OCR + lite JSON VLM),外层 8 图全并发,目标 8 图 <15s。
-
-输出 dict 格式与 viral_video._normalize() 完全一致,下游信任链/t2i 零改动。
-灰度开关:VISION_V2_ENABLED=true(默认 false,走旧 #2198/#2199 竞速逻辑)。
-"""
-
+"""V2 图片分析:火山OCR专用API + doubao-lite强约束JSON并行,单次pro VLM兜底。"""
from .fast_path import analyze_image_v2, analyze_images_v2 # noqa: F401
diff --git a/apps/worker/worker_app/tasks/vision/assembler.py b/apps/worker/worker_app/tasks/vision/assembler.py
index 79be272e9..b45a4d60e 100644
--- a/apps/worker/worker_app/tasks/vision/assembler.py
+++ b/apps/worker/worker_app/tasks/vision/assembler.py
@@ -21,8 +21,6 @@ def _join_parts(*parts: str | None) -> str:
_AGE_PREFIX = {
- "儿童": "小女孩" if None else "儿童",
- "青少年": "少女" if None else "少年",
"青年": "年轻",
"中年": "中年",
"老年": "老年",
@@ -155,7 +153,8 @@ def _build_portrait_prompt(fj: dict[str, Any]) -> str:
if detail_parts:
pieces.append(",".join(detail_parts))
if style_parts:
- pieces.append(",".join(style_parts) + "风格")
+ # 风格词之间不用逗号,用空格紧凑
+ pieces.append("".join(style_parts) + "风格")
else:
pieces.append("人像写真")
diff --git a/apps/worker/worker_app/tasks/vision/fast_path.py b/apps/worker/worker_app/tasks/vision/fast_path.py
index b8226a4c5..6b6585ebf 100644
--- a/apps/worker/worker_app/tasks/vision/fast_path.py
+++ b/apps/worker/worker_app/tasks/vision/fast_path.py
@@ -1,9 +1,10 @@
# -*- coding: utf-8 -*-
-"""V2 快速路径:每图并行 OCR + lite JSON VLM,失败降级 pro VLM。
+"""V2 图片分析主路径:每图并行 OCR(火山专用API)+ lite JSON VLM,失败时单次 pro VLM 兜底。
-单图并行 2 路(OCR + lite JSON VLM),目标 <3s。
-外层 8 图全并发,目标 8 图 <15s。
-终极兜底:复用旧 _analyze_single_image 完整 pro VLM 逻辑。
+设计原则(灵应10-05要求):
+- 主力路径简洁:单图2路并行,外层N图全并发
+- 兜底简单:单次 pro VLM 调用,无竞速/重试/复杂超时
+- 输出 dict 格式与旧版完全一致,下游零改动
"""
from __future__ import annotations
@@ -14,230 +15,114 @@ import time
from concurrent.futures import ThreadPoolExecutor, as_completed
from typing import Any
-from . import assembler, ocr_volc, vlm_fast_json
+from . import assembler, ocr_volc, vlm_fallback, vlm_fast_json
logger = logging.getLogger(__name__)
-# ---------- 配置项(可通过环境变量覆盖) ----------
-VISION_V2_ENABLED = os.environ.get("VISION_V2_ENABLED", "false").lower() in ("1", "true", "yes", "on")
-# 单图 fast 路径总超时(包含 OCR + fast_json 并行)
-V2_FAST_TIMEOUT = float(os.environ.get("VISION_V2_FAST_TIMEOUT", "10"))
-# 外层图片并发(默认 8,即全并行)
-V2_IMG_WORKERS = int(os.environ.get("VISION_V2_IMG_WORKERS", "8"))
-# fast_json 单次超时
-V2_FAST_JSON_TIMEOUT = float(os.environ.get("VISION_V2_FAST_JSON_TIMEOUT", "8"))
-# OCR 单次超时
-V2_OCR_TIMEOUT = float(os.environ.get("VISION_V2_OCR_TIMEOUT", "8"))
-# pro VLM 兜底超时(仅在 fast 路径完全失败时触发)
-V2_PRO_TIMEOUT = float(os.environ.get("VISION_V2_PRO_TIMEOUT", "45"))
-V2_LITE_TIMEOUT = float(os.environ.get("VISION_V2_LITE_TIMEOUT", "20"))
+# 可通过环境变量调参(有默认值,无需配置即可跑)
+_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"))
+_PRO_TIMEOUT = float(os.environ.get("VISION_V2_PRO_TIMEOUT", "45"))
+
+_FALLBACK_RESULT = {
+ "name": "未识别", "brand": "无法判断", "category": "非产品图",
+ "appearance": "无法判断", "packaging": "无法判断", "text_on_package": [],
+ "key_features": ["无法判断"], "scene": "通用", "mood": "",
+ "portrait_prompt": "无法判断", "summary": "未识别",
+}
-def _is_result_usable(result: dict[str, Any]) -> bool:
- """与 viral_video._is_vision_result_usable 对齐的可用判定。"""
- pp = (result.get("portrait_prompt") or "").strip()
+def _is_usable(r: dict[str, Any]) -> bool:
+ """结果可用判定:portrait_prompt 是核心,有效就算 usable。"""
+ pp = (r.get("portrait_prompt") or "").strip()
if pp and pp not in ("无人像", "无法判断", "未识别"):
return True
- name = result.get("name") or ""
+ name = (r.get("name") or "").strip()
if name and name not in ("未识别", "无法判断", "未知"):
return True
- summary = result.get("summary") or ""
- if len(summary) >= 5 and summary not in ("无法判断", "未识别"):
- return True
- cat = result.get("category") or ""
- if cat == "非产品图" and pp != "无人像":
- return True
- kf = result.get("key_features") or []
- if kf and kf != ["无法判断"]:
- # 只要有非默认特征且非空
- return True
return False
-def _call_pro_fallback(img_url: str, idx: int, lite_model: str, pro_model: str) -> dict[str, Any] | None:
- """fast 路径失败时,调用旧的 lite/pro 竞速 VLM。
-
- 复用 viral_video._analyze_single_image 的实现,避免重复代码。
- """
- try:
- from worker_app.tasks.viral_video import _analyze_single_image
- except ImportError:
- try:
- from tasks.viral_video import _analyze_single_image # type: ignore
- except ImportError:
- logger.warning("[vision.v2] 无法 import _analyze_single_image,跳过 pro 兜底")
- return None
- try:
- return _analyze_single_image(
- idx,
- img_url,
- vision_model=lite_model,
- timeout=int(V2_LITE_TIMEOUT),
- pro_fallback_model=pro_model,
- )
- except Exception as e:
- logger.warning("[vision.v2] 图片 #%d pro 兜底异常 err=%s", idx, e, exc_info=True)
- return None
-
-
-def analyze_image_v2(
- idx: int,
- img_url: str,
- *,
- lite_model: str | None = None,
- pro_model: str | None = None,
-) -> dict[str, Any]:
- """单张图片 V2 分析:OCR + lite JSON VLM 并行,必要时降级 pro VLM。
-
- 返回的 dict 与 viral_video._normalize() 输出格式完全一致。
- """
+def analyze_image_v2(idx: int, img_url: str) -> dict[str, Any]:
+ """单张图片 V2 分析。"""
t0 = time.time()
- # ---- 第 1 层:fast 路径并行 ----
- fast_json_result: dict[str, Any] | None = None
+
+ # 第1层:OCR + lite JSON VLM 并行
+ fj_result: dict[str, Any] | None = None
ocr_result: list[str] = []
-
with ThreadPoolExecutor(max_workers=2) as pool:
- f_fj = pool.submit(
- vlm_fast_json.call_fast_json,
- img_url,
- model=lite_model,
- timeout=V2_FAST_JSON_TIMEOUT,
- )
- f_ocr = pool.submit(ocr_volc.call_ocr, img_url, timeout=V2_OCR_TIMEOUT)
-
- # 等全部完成或超时
- for fut in as_completed([f_fj, f_ocr], timeout=V2_FAST_TIMEOUT):
+ 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 fast 子任务异常: %s", idx, e)
+ logger.warning("[vision.v2] 图片 #%d 子任务异常: %s", idx, e)
continue
- if fut is f_fj:
- fast_json_result = res if isinstance(res, dict) else None
- elif fut is f_ocr:
- ocr_result = res if isinstance(res, list) else []
+ if fut is f_fj and isinstance(res, dict):
+ fj_result = res
+ elif fut is f_ocr and isinstance(res, list):
+ ocr_result = res
fast_elapsed = time.time() - t0
- # ---- 组装 fast 结果 ----
- assembled: dict[str, Any] | None = None
- if fast_json_result:
- assembled = assembler.assemble_result(idx, fast_json_result, ocr_result)
- if _is_result_usable(assembled):
+ # 组装 fast 结果
+ if fj_result:
+ assembled = assembler.assemble_result(idx, fj_result, ocr_result)
+ if _is_usable(assembled):
assembled["_fast_elapsed"] = round(fast_elapsed, 2)
logger.info(
- "[vision.v2] 图片 #%d fast 路径命中 elapsed=%.2fs portrait_prompt=%s",
- idx,
- fast_elapsed,
- (assembled.get("portrait_prompt") or "")[:40],
+ "[vision.v2] 图片 #%d fast命中 elapsed=%.2fs pp=%s",
+ idx, fast_elapsed, (assembled.get("portrait_prompt") or "")[:40],
)
return assembled
- logger.info(
- "[vision.v2] 图片 #%d fast 结果不可用 portrait_prompt=%s,走 pro 兜底",
- idx,
- (assembled.get("portrait_prompt") or "")[:40],
- )
- else:
- logger.info("[vision.v2] 图片 #%d fast_json 返回空 elapsed=%.2fs,走 pro 兜底", idx, fast_elapsed)
- # ---- 第 2 层:pro VLM 兜底(复用旧竞速逻辑)----
+ # 第2层:pro VLM 单次兜底
pro_t0 = time.time()
- use_lite = lite_model or vlm_fast_json.DEFAULT_LITE_MODEL
- use_pro = pro_model or "doubao-seed-2-1-pro-260915"
- pro_result = _call_pro_fallback(img_url, idx, use_lite, use_pro)
- if pro_result and _is_result_usable(pro_result):
+ pro_result = vlm_fallback.call_pro_vlm(img_url, idx, timeout=_PRO_TIMEOUT)
+ if pro_result and _is_usable(pro_result):
pro_result["_fallback_used"] = True
pro_result["_fast_elapsed"] = round(fast_elapsed, 2)
pro_result["_pro_elapsed"] = round(time.time() - pro_t0, 2)
- logger.info(
- "[vision.v2] 图片 #%d pro 兜底命中 total_elapsed=%.2fs",
- idx,
- time.time() - t0,
- )
+ if ocr_result and not pro_result.get("text_on_package"):
+ pro_result["text_on_package"] = ocr_result[:8]
+ logger.info("[vision.v2] 图片 #%d pro兜底命中 total=%.2fs", idx, time.time()-t0)
return pro_result
- # ---- 第 3 层:兜底失败,返回 assembled 或标准 fallback ----
- if assembled:
- assembled["_source"] = "v2_fast_json_degraded"
- logger.warning(
- "[vision.v2] 图片 #%d pro 兜底也失败,返回降级 fast 结果 elapsed=%.2fs",
- idx,
- time.time() - t0,
- )
- return assembled
-
- # 最后的最后:返回最小可用结构
- logger.warning("[vision.v2] 图片 #%d 所有路径均失败 elapsed=%.2fs", idx, time.time() - t0)
- return {
- "name": "未识别",
- "brand": "无法判断",
- "category": "非产品图",
- "appearance": "无法判断",
- "packaging": "无法判断",
- "text_on_package": ocr_result[:8],
- "key_features": ["无法判断"],
- "scene": "通用",
- "mood": "",
- "portrait_prompt": "无法判断",
- "summary": "未识别",
- "_source": "v2_all_failed",
- }
+ # 最终:返回最小可用结果
+ logger.warning("[vision.v2] 图片 #%d 全路径失败 elapsed=%.2fs", idx, time.time()-t0)
+ out = dict(_FALLBACK_RESULT)
+ out["_source"] = "v2_all_failed"
+ out["text_on_package"] = ocr_result[:8]
+ out["_fast_elapsed"] = round(fast_elapsed, 2)
+ return out
-def analyze_images_v2(
- img_urls: list[str],
- *,
- lite_model: str | None = None,
- pro_model: str | None = None,
- max_workers: int | None = None,
-) -> list[dict[str, Any]]:
- """批量图片 V2 分析(外层全并行)。"""
+def analyze_images_v2(img_urls: list[str]) -> list[dict[str, Any]]:
+ """批量图片 V2 分析,外层全并发。"""
if not img_urls:
return []
- workers = max_workers if max_workers and max_workers > 0 else V2_IMG_WORKERS
- workers = min(workers, len(img_urls), 16) # 安全上限 16
+ workers = min(_IMG_WORKERS, len(img_urls), 16)
results: list[dict[str, Any] | None] = [None] * len(img_urls)
- logger.info(
- "[vision.v2] 开始 V2 并行图片分析 n=%d workers=%d fast_timeout=%.0fs",
- len(img_urls),
- workers,
- V2_FAST_TIMEOUT,
- )
+ logger.info("[vision.v2] 开始图片分析 n=%d workers=%d fast_timeout=%.0fs", len(img_urls), workers, _FAST_TIMEOUT)
t0 = time.time()
with ThreadPoolExecutor(max_workers=workers) as pool:
- future_to_idx = {
- pool.submit(analyze_image_v2, idx, url, lite_model=lite_model, pro_model=pro_model): idx
- for idx, url in enumerate(img_urls)
- }
+ future_to_idx = {pool.submit(analyze_image_v2, idx, url): idx for idx, url in enumerate(img_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("[vision.v2] 图片 #%d future 异常 err=%s", idx, e, exc_info=True)
- results[idx] = {
- "name": "未识别",
- "brand": "无法判断",
- "category": "非产品图",
- "appearance": "无法判断",
- "packaging": "无法判断",
- "text_on_package": [],
- "key_features": ["无法判断"],
- "scene": "通用",
- "mood": "",
- "portrait_prompt": "无法判断",
- "summary": "未识别",
- "_source": "v2_future_exception",
- }
+ logger.warning("[vision.v2] 图片 #%d future异常: %s", idx, e, exc_info=True)
+ r = dict(_FALLBACK_RESULT)
+ r["_source"] = "v2_future_exception"
+ results[idx] = r
+
elapsed = time.time() - t0
- succ = sum(1 for r in results if r and _is_result_usable(r))
+ succ = sum(1 for r in results if r and _is_usable(r))
fb = sum(1 for r in results if r and r.get("_fallback_used"))
- logger.info(
- "[vision.v2] V2 图片分析完成 n=%d success=%d pro_fallback=%d elapsed=%.2fs",
- len(img_urls),
- succ,
- fb,
- elapsed,
- )
+ logger.info("[vision.v2] 完成 n=%d usable=%d pro_fallback=%d elapsed=%.2fs", len(img_urls), succ, fb, elapsed)
return [r for r in results if r is not None]
diff --git a/apps/worker/worker_app/tasks/vision/vlm_fallback.py b/apps/worker/worker_app/tasks/vision/vlm_fallback.py
new file mode 100644
index 000000000..90701630b
--- /dev/null
+++ b/apps/worker/worker_app/tasks/vision/vlm_fallback.py
@@ -0,0 +1,197 @@
+# -*- coding: utf-8 -*-
+"""VLM 兜底:专用API路径失败时的最后一道防线,单次调用 doubao-seed-2.1-pro。
+
+设计原则:简单、直接、无竞速、无复杂超时逻辑。只在 fast_json 结果不可用时调用。
+"""
+from __future__ import annotations
+
+import json
+import logging
+import re
+import time
+from typing import Any
+
+logger = logging.getLogger(__name__)
+
+DEFAULT_PRO_MODEL = "doubao-seed-2-1-pro-260915"
+DEFAULT_TIMEOUT = 45
+DEFAULT_MAX_TOKENS = 800
+
+
+def _strip_code_fence(s: str) -> str:
+ s = s.strip()
+ if s.startswith("```"):
+ lines = s.split("\n")
+ if lines and lines[0].startswith("```"):
+ lines = lines[1:]
+ if lines and lines[-1].strip().startswith("```"):
+ lines = lines[:-1]
+ s = "\n".join(lines).strip()
+ return s
+
+
+def _xml_text(tag: str, xml: str) -> str:
+ m = re.search(rf"<{tag}[^>]*>(.*?){tag}>", xml, re.S)
+ return (m.group(1) if m else "").strip()
+
+
+def _xml_attr(tag: str, attr: str, xml: str) -> str:
+ m = re.search(rf"<{tag}[^>]*\b{attr}\s*=\s*[\"']([^\"']*)[\"']", xml)
+ return (m.group(1) if m else "").strip()
+
+
+def _xml_to_product(raw: str, idx: int) -> dict[str, Any]:
+ """解析 VLM 输出的 XML 格式(简化版)。"""
+ scene = _xml_text("scene", raw) or "通用"
+ mood = _xml_text("mood", raw) or ""
+
+ portrait_prompt = "无人像"
+ p_has = _xml_attr("people", "has_person", raw)
+ if p_has and p_has.lower() != "false":
+ gender = _xml_attr("people", "gender", raw) or ""
+ age = _xml_attr("people", "age_range", raw) or ""
+ outfit = _xml_attr("people", "outfit", raw) or ""
+ hair = _xml_attr("people", "hair", raw) or "自然发型"
+ pose = _xml_attr("people", "pose", raw) or ""
+ expr = _xml_attr("people", "expression", raw) or "自然"
+ parts: list[str] = []
+ if gender:
+ parts.append(gender + ("性" if not gender.endswith("性") else ""))
+ if age:
+ parts.append(age)
+ parts.append("人物")
+ parts.append(hair)
+ if outfit:
+ parts.append(f"身着{outfit}")
+ if pose:
+ parts.append(f"姿态{pose}")
+ parts.append(f"表情{expr}")
+ portrait_prompt = ",".join(parts)
+
+ m = re.search(r"]*>(.*?)", raw, re.S)
+ if m:
+ pbody = m.group(1)
+ name = _xml_attr("product", "name", raw) or _xml_text("name", pbody) or "未识别"
+ brand = _xml_attr("product", "brand", raw) or _xml_text("brand", pbody) or "无法判断"
+ category = _xml_attr("product", "category", raw) or _xml_text("category", pbody) or "无法判断"
+ appearance = _xml_attr("product", "appearance", raw) or _xml_text("appearance", pbody) or "无法判断"
+ packaging = _xml_attr("product", "packaging", raw) or _xml_text("packaging", pbody) or "无法判断"
+ feat = _xml_attr("product", "features", raw) or _xml_text("features", pbody) or ""
+ feat_list = [x.strip() for x in re.split(r"[,,;;]", feat) if x.strip()] if feat else ["无法判断"]
+ top_text = _xml_attr("product", "text_on_package", raw) or _xml_text("text_on_package", pbody) or ""
+ text_list = [x.strip() for x in re.split(r"[,,;;]", top_text) if x.strip()] if top_text else []
+ summary = _xml_attr("product", "summary", raw) or _xml_text("summary", pbody) or f"{brand} {name}"
+ pp_attr = _xml_attr("product", "portrait_prompt", raw)
+ if pp_attr and pp_attr != "无人像":
+ portrait_prompt = pp_attr
+ return {
+ "name": name, "brand": brand, "category": category,
+ "appearance": appearance, "packaging": packaging, "text_on_package": text_list,
+ "key_features": feat_list, "scene": scene, "mood": mood,
+ "portrait_prompt": portrait_prompt, "summary": summary,
+ "_source": "vlm_pro_xml",
+ }
+
+ if portrait_prompt != "无人像":
+ return {
+ "name": "未识别", "brand": "无法判断", "category": "无法判断",
+ "appearance": "无法判断", "packaging": "无法判断", "text_on_package": [],
+ "key_features": ["无法判断"], "scene": scene, "mood": mood,
+ "portrait_prompt": portrait_prompt, "summary": "未识别",
+ "_source": "vlm_pro_no_product",
+ }
+ return {
+ "name": "未识别", "brand": "无法判断", "category": "无法判断",
+ "appearance": "无法判断", "packaging": "无法判断", "text_on_package": [],
+ "key_features": ["无法判断"], "scene": scene, "mood": mood,
+ "portrait_prompt": "无人像", "summary": "未识别", "_source": "vlm_pro_no_tag",
+ }
+
+
+def call_pro_vlm(
+ img_url: str,
+ idx: int,
+ *,
+ model: str | None = None,
+ timeout: int = DEFAULT_TIMEOUT,
+) -> dict[str, Any] | None:
+ """单次调用 pro VLM,解析后返回 product dict;失败返回 None。"""
+ t0 = time.time()
+ try:
+ from packages.application.viral_video.prompt_loader import (
+ get_template,
+ render_system_prompt,
+ render_user_prompt,
+ )
+ from packages.shared.ai_client import get_doubao_client
+ except ImportError as e:
+ logger.warning("[vision.vlm] 导入失败: %s", e)
+ return None
+
+ try:
+ template = get_template("image_analysis")
+ system = render_system_prompt(template)
+ user = render_user_prompt(template, image_count=1, industry="通用", image_urls=f"第1张:{img_url}")
+ except Exception as e:
+ logger.warning("[vision.vlm] 模板加载失败: %s", e)
+ return None
+
+ client = get_doubao_client()
+ if not client.is_available:
+ return None
+
+ use_model = model or DEFAULT_PRO_MODEL
+ try:
+ raw = client.vision_completion(
+ messages=[{"role": "system", "content": system}, {"role": "user", "content": user}],
+ images=[img_url],
+ temperature=0.3,
+ max_tokens=DEFAULT_MAX_TOKENS,
+ timeout=timeout,
+ model=use_model,
+ )
+ except Exception as e:
+ logger.warning("[vision.vlm] 图片 #%d pro VLM 调用失败 elapsed=%.1fs err=%s", idx, time.time()-t0, e)
+ return None
+
+ elapsed = time.time() - t0
+ if not raw:
+ logger.warning("[vision.vlm] 图片 #%d pro VLM 返回空 elapsed=%.1fs", idx, elapsed)
+ return None
+
+ text = _strip_code_fence(raw)
+ l, r = text.find("{"), text.rfind("}")
+ if l >= 0 and r > l:
+ try:
+ obj = json.loads(text[l:r+1])
+ if isinstance(obj, dict):
+ logger.info("[vision.vlm] 图片 #%d pro VLM JSON 完成 elapsed=%.1fs", idx, elapsed)
+ return {
+ "name": obj.get("name") or "未识别",
+ "brand": obj.get("brand") or "无法判断",
+ "category": obj.get("category") or "无法判断",
+ "appearance": obj.get("appearance") or "无法判断",
+ "packaging": obj.get("packaging") or "无法判断",
+ "text_on_package": obj.get("text_on_package") or [],
+ "key_features": obj.get("key_features") or obj.get("features") or ["无法判断"],
+ "scene": obj.get("scene") or "通用",
+ "mood": obj.get("mood") or "",
+ "portrait_prompt": obj.get("portrait_prompt") or "无人像",
+ "summary": obj.get("summary") or f"{obj.get('brand','')} {obj.get('name','')}",
+ "_source": "vlm_pro_json",
+ }
+ except json.JSONDecodeError:
+ pass
+
+ try:
+ result = _xml_to_product(text, idx)
+ result["_fallback_used"] = True
+ result["_pro_elapsed"] = round(elapsed, 2)
+ logger.info(
+ "[vision.vlm] 图片 #%d pro VLM XML 完成 elapsed=%.2fs pp=%s",
+ idx, elapsed, (result.get("portrait_prompt") or "")[:40],
+ )
+ return result
+ except Exception as e:
+ logger.warning("[vision.vlm] 图片 #%d 解析失败 elapsed=%.1fs err=%s head=%s", idx, elapsed, e, raw[:200])
+ return None