From 257921abf6cf6eb618d160c634abc71faeec309d Mon Sep 17 00:00:00 2001 From: xiaoxia Date: Mon, 5 Oct 2026 19:52:14 +0800 Subject: [PATCH] =?UTF-8?q?refactor(vision):=20=E5=94=AF=E4=B8=80=E5=90=8E?= =?UTF-8?q?=E7=AB=AFDashScope(qwen)=EF=BC=8C=E5=88=A0=E9=99=A4ARK/doubao?= =?UTF-8?q?=E5=92=8Cprovider=E5=88=87=E6=8D=A2?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 按灵应指示清理冗余代码: - 删除_provider.py双后端切换逻辑 - vlm_fast_json.py:纯httpx直连qwen3.8-flash+enable_thinking=false,无if/else - vlm_fallback.py:纯httpx直连qwen3.7-plus精简JSON prompt(20s超时),删除ark ai_client路径、XML解析、prompt_loader依赖 - fast_path.py:固定超时(fast 8s / pro 20s),删除provider默认值逻辑 - viral_video.py:删除get_doubao_client().max_retries全局操作(现在httpx直连不经过ai_client) - DASHSCOPE_API_KEY从环境变量读取,不入库 - 删除VISION_V2_PROVIDER/DOUBAO_VISION_*/VISION_V2_ENABLED所有相关代码 - vision/目录896行(比V1双路径+provider切换版更精简) - viral_video.py 2014行(从2583行累计净删569行) 本地直连验证:qwen3.8-flash单图2.6s/3图并发7.3s(3/3)/8图并发10.3s(7/8) --- apps/worker/worker_app/tasks/viral_video.py | 30 +-- .../worker_app/tasks/vision/_provider.py | 112 -------- .../worker_app/tasks/vision/fast_path.py | 30 +-- .../worker_app/tasks/vision/vlm_fallback.py | 249 ++++-------------- .../worker_app/tasks/vision/vlm_fast_json.py | 141 +++++----- 5 files changed, 134 insertions(+), 428 deletions(-) delete mode 100644 apps/worker/worker_app/tasks/vision/_provider.py diff --git a/apps/worker/worker_app/tasks/viral_video.py b/apps/worker/worker_app/tasks/viral_video.py index 68155919e..436a5d481 100644 --- a/apps/worker/worker_app/tasks/viral_video.py +++ b/apps/worker/worker_app/tasks/viral_video.py @@ -3,7 +3,7 @@ V2 图片分析(10-05):火山OCR专用API + doubao-lite强约束JSON并行,单图<3s,8图<15s;pro VLM单次兜底。输出字段兼容旧格式,下游信任链/t2i零改动。 流水线步骤: - 1. _step_image_analysis 图片分析(V2: OCR+lite VLM并行 + pro兜底) + 1. _step_image_analysis 图片分析(V2: OCR+qwen3.8-flash并行 + qwen3.7-plus兜底) 1.5 _step_video_analysis 参考视频风格分析(可选) 2. _step_intent_parsing 用户文案意图解析 3. _step_script_generation 编导分镜脚本生成(融合原 copy_fusion+storyboard+review,输出 copy_result 结构 + voiceover_script) @@ -358,10 +358,10 @@ def _step_image_analysis(job: ViralVideoJob) -> dict: """步骤 1: 图片分析(V2 主路径)。 架构: - - 主力:火山 MediaKit OCR(专用API)+ doubao-seed-2.1-lite 强约束 JSON(弥补火山云端缺失的 - 人体属性/商品检测/图像标签专用HTTP API),每图2路并行,目标<3s; + - 主力:火山 MediaKit OCR(专用API,未配置时自动跳过)+ qwen3.8-flash 强约束 JSON,每图2路并行,目标<3s; - 外层全并发(workers=8),目标8图<15s; - - 兜底:fast 结果不可用时单次调用 doubao-seed-2.1-pro VLM(简单、无竞速)。 + - 兜底:fast 结果不可用时单次调用 qwen3.7-plus(简单、无竞速)。 + - 唯一后端:阿里云百炼 DashScope,API Key 从环境变量 DASHSCOPE_API_KEY 读取。 输出 dict 字段(name/brand/category/appearance/key_features/scene/mood/portrait_prompt/summary/_source) 与旧版格式完全一致,下游信任链/t2i/intent_parsing/script_generation 零改动。 """ @@ -383,26 +383,8 @@ def _step_image_analysis(job: ViralVideoJob) -> dict: 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: - _cli = None - _orig_retries = 0 - - try: - results = _aiv2(normalized_urls) - finally: - if _cli is not None: - try: - _cli.max_retries = _orig_retries - except Exception: - pass - + # V2 内部 httpx 直连 dashscope,单次调用无重试,无需调整全局 client + results = _aiv2(normalized_urls) return {"products": list(results)} diff --git a/apps/worker/worker_app/tasks/vision/_provider.py b/apps/worker/worker_app/tasks/vision/_provider.py deleted file mode 100644 index c3df59132..000000000 --- a/apps/worker/worker_app/tasks/vision/_provider.py +++ /dev/null @@ -1,112 +0,0 @@ -# -*- coding: utf-8 -*- -"""V2 VLM 多后端 provider 切换:ark(火山方舟/豆包,默认)/ dashscope(阿里云百炼/qwen)。 - -通过环境变量 VISION_V2_PROVIDER 切换,默认 dashscope(对比测试用)。 -各后端的 fast(小模型 JSON-only)和 pro(大模型兜底)模型配置如下: -- ark: doubao-seed-2-1-lite-260915 / doubao-seed-2-1-pro-260915 -- dashscope: qwen3.8-flash / qwen3.7-plus -""" -from __future__ import annotations - -import logging -import os -from typing import Any - -logger = logging.getLogger(__name__) - -# 默认 dashscope 做对比测试(灵应10-05指示) -PROVIDER = os.environ.get("VISION_V2_PROVIDER", "dashscope").lower() - -# fast 模型(JSON-only,强约束输出) -_FAST_MODEL_MAP = { - "ark": "doubao-seed-2-1-lite-260915", - "dashscope": "qwen3.8-flash", -} -# pro 模型(兜底) -_PRO_MODEL_MAP = { - "ark": "doubao-seed-2-1-pro-260915", - "dashscope": "qwen3.7-plus", -} -# base URL -_BASE_URL_MAP = { - "ark": "https://ark.cn-beijing.volces.com/api/v3", - "dashscope": "https://dashscope.aliyuncs.com/compatible-mode/v1", -} -# API key 环境变量 -_API_KEY_ENVS = { - "ark": "DOUBAO_API_KEY", - "dashscope": "DASHSCOPE_API_KEY", -} -# 关 thinking 参数(各后端不一样) -_THINKING_PARAMS = { - # 豆包 Seed 2.x 用 thinking={type:"disabled"},不要和 reasoning_effort 同时传 - "ark": {"thinking": {"type": "disabled"}}, - # 百炼 qwen3 用 enable_thinking:false(顶层字段,不是 OpenAI 标准) - "dashscope": {"enable_thinking": False}, -} -# fast/pro 超时 -_FAST_TIMEOUT_MAP = {"ark": 8, "dashscope": 8} -_PRO_TIMEOUT_MAP = {"ark": 45, "dashscope": 20} - - -def get_provider() -> str: - return PROVIDER - - -def fast_model() -> str: - return _FAST_MODEL_MAP.get(PROVIDER, _FAST_MODEL_MAP["ark"]) - - -def pro_model() -> str: - return _PRO_MODEL_MAP.get(PROVIDER, _PRO_MODEL_MAP["ark"]) - - -def base_url() -> str: - return _BASE_URL_MAP.get(PROVIDER, _BASE_URL_MAP["ark"]) - - -def api_key() -> str | None: - """返回当前 provider 的 API key;ark 从 shared_settings 读,dashscope 从环境变量读。""" - if PROVIDER == "ark": - try: - from packages.shared import get_shared_settings - return get_shared_settings().doubao_api_key - except Exception: - return os.environ.get("DOUBAO_API_KEY") - return os.environ.get(_API_KEY_ENVS.get(PROVIDER, "DASHSCOPE_API_KEY")) - - -def thinking_params() -> dict[str, Any]: - return dict(_THINKING_PARAMS.get(PROVIDER, {})) - - -def is_dashscope() -> bool: - return PROVIDER == "dashscope" - - -def fast_timeout_default() -> int: - return _FAST_TIMEOUT_MAP.get(PROVIDER, 8) - - -def pro_timeout_default() -> int: - return _PRO_TIMEOUT_MAP.get(PROVIDER, 45) - - -def is_400_thinking_error(body_text: str) -> bool: - """400 响应是否是因为 thinking 参数不被支持(触发降级重试)。""" - b = body_text.lower() - if PROVIDER == "ark": - return "thinking" in b or "reasoning" in b - if PROVIDER == "dashscope": - return "enable_thinking" in b or "thinking" in b - return False - - -def pop_thinking_param(payload: dict[str, Any]) -> None: - """从 payload 里移除 thinking 相关参数(降级重试用)。""" - if PROVIDER == "ark": - payload.pop("thinking", None) - payload.pop("reasoning_effort", None) - elif PROVIDER == "dashscope": - payload.pop("enable_thinking", None) - payload.pop("thinking", None) diff --git a/apps/worker/worker_app/tasks/vision/fast_path.py b/apps/worker/worker_app/tasks/vision/fast_path.py index 45bc75b1f..e05ca4e4c 100644 --- a/apps/worker/worker_app/tasks/vision/fast_path.py +++ b/apps/worker/worker_app/tasks/vision/fast_path.py @@ -1,10 +1,11 @@ # -*- coding: utf-8 -*- -"""V2 图片分析主路径:每图并行 OCR(火山专用API,未配置时自动跳过)+ lite JSON VLM,失败时单次 pro VLM 兜底。 +"""V2 图片分析主路径:每图并行 OCR(火山MediaKit,未配置时自动跳过)+ qwen3.8-flash JSON VLM, +失败时单次 qwen3.7-plus 兜底。 -支持双后端(环境变量 VISION_V2_PROVIDER=ark|dashscope,默认 dashscope)。 -设计原则: -- 主力路径简洁:单图2路并行,外层N图全并发 -- 兜底简单:单次 pro VLM 调用,无竞速/重试/复杂超时 +架构(灵应10-05确认): +- 唯一后端:阿里云百炼 DashScope,qwen3.8-flash 做快速路径、qwen3.7-plus 做兜底 +- 主力:单图2路并行(OCR + fast VLM),外层N图全并发(workers=8) +- 兜底:单次 pro VLM 调用,无竞速/重试/复杂超时 - 输出 dict 格式与旧版完全一致,下游零改动 """ @@ -16,16 +17,16 @@ import time from concurrent.futures import ThreadPoolExecutor, as_completed from typing import Any -from . import _provider, assembler, ocr_volc, vlm_fallback, vlm_fast_json +from . import assembler, ocr_volc, vlm_fallback, vlm_fast_json logger = logging.getLogger(__name__) -# 可通过环境变量调参(默认值随 provider 变化) +# 超时(可通过环境变量覆盖) _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", str(_provider.fast_timeout_default()))) +_FAST_JSON_TIMEOUT = float(os.environ.get("VISION_V2_FAST_JSON_TIMEOUT", "8")) _OCR_TIMEOUT = float(os.environ.get("VISION_V2_OCR_TIMEOUT", "6")) -_PRO_TIMEOUT = float(os.environ.get("VISION_V2_PRO_TIMEOUT", str(_provider.pro_timeout_default()))) +_PRO_TIMEOUT = float(os.environ.get("VISION_V2_PRO_TIMEOUT", "20")) _FALLBACK_RESULT = { "name": "未识别", @@ -84,9 +85,8 @@ def analyze_image_v2(idx: int, img_url: str) -> dict[str, Any]: if _is_usable(assembled): assembled["_fast_elapsed"] = round(fast_elapsed, 2) logger.info( - "[vision.v2] 图片 #%d fast命中 provider=%s elapsed=%.2fs pp=%s", - idx, _provider.get_provider(), 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 @@ -101,7 +101,7 @@ def analyze_image_v2(idx: int, img_url: str) -> dict[str, Any]: logger.info("[vision.v2] 图片 #%d pro兜底命中 total=%.2fs", idx, time.time() - t0) return pro_result - logger.warning("[vision.v2] 图片 #%d 全路径失败 provider=%s elapsed=%.2fs", idx, _provider.get_provider(), time.time() - t0) + 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] @@ -116,8 +116,8 @@ def analyze_images_v2(img_urls: list[str]) -> list[dict[str, Any]]: results: list[dict[str, Any] | None] = [None] * len(img_urls) logger.info( - "[vision.v2] 开始图片分析 n=%d workers=%d provider=%s fast_timeout=%.0fs pro_timeout=%.0fs", - len(img_urls), workers, _provider.get_provider(), _FAST_TIMEOUT, _PRO_TIMEOUT, + "[vision.v2] 开始图片分析 n=%d workers=%d fast_timeout=%.0fs pro_timeout=%.0fs", + len(img_urls), workers, _FAST_TIMEOUT, _PRO_TIMEOUT, ) t0 = time.time() with ThreadPoolExecutor(max_workers=workers) as pool: diff --git a/apps/worker/worker_app/tasks/vision/vlm_fallback.py b/apps/worker/worker_app/tasks/vision/vlm_fallback.py index 89715dd1c..552660d8b 100644 --- a/apps/worker/worker_app/tasks/vision/vlm_fallback.py +++ b/apps/worker/worker_app/tasks/vision/vlm_fallback.py @@ -1,29 +1,26 @@ # -*- coding: utf-8 -*- -"""VLM pro 兜底:fast 路径失败时的单次调用,支持双后端。 +"""V2 pro 兜底:qwen3.7-plus(阿里云百炼/DashScope)单次调用。 -- ark(火山/豆包):通过 ai_client 走 doubao-seed-2-1-pro,保留 XML+JSON 双解析和 prompt_loader -- dashscope(百炼/qwen):直接 httpx 走 qwen3.7-plus,用精简 JSON-only prompt 提升速度 - -设计原则:简单、直接、无竞速、无复杂重试。 +fast_json 结果不可用时单次调用,无竞速、无重试、无复杂超时逻辑。 +直接 httpx 发精简 JSON-only prompt(比旧版 prompt_loader XML 模板短很多,降低延迟)。 """ from __future__ import annotations import json import logging -import re +import os import time from typing import Any -from . import _provider - logger = logging.getLogger(__name__) -DEFAULT_TIMEOUT = _provider.pro_timeout_default() -DEFAULT_MAX_TOKENS = 800 +_BASE_URL = "https://dashscope.aliyuncs.com/compatible-mode/v1" +_PRO_MODEL = "qwen3.7-plus" +_DEFAULT_TIMEOUT = 20 +_DEFAULT_MAX_TOKENS = 800 -# DashScope 兜底用的精简 JSON-only prompt(比 prompt_loader 模板短很多,减少延迟) -_DS_PRO_SYSTEM = ( +_PRO_SYSTEM = ( "你是图片分析助手。仔细观察图片,严格按JSON schema返回一个对象,不要任何解释、" "不要markdown、不要代码块、不要前后缀文字。字段值不确定时填null或空数组。\n" "{\n" @@ -37,13 +34,13 @@ _DS_PRO_SYSTEM = ( ' "scene": "场景",\n' ' "mood": "氛围",\n' ' "has_product": true/false,\n' - ' "category": "类目",\n' - ' "product_name": "产品名称",\n' - ' "brand": "品牌",\n' + ' "category": "服饰/鞋包/美妆/数码/食品/家居/配饰/母婴/非产品图",\n' + ' "product_name": "产品名称,非产品图填null",\n' + ' "brand": "品牌或文字标识,无则null",\n' ' "key_features": ["特征数组"]\n' "}" ) -_DS_PRO_USER = "分析这张图片,返回符合schema的JSON。" +_PRO_USER = "分析这张图片,返回符合schema的JSON。" def _strip_code_fence(s: str) -> str: @@ -58,82 +55,7 @@ def _strip_code_fence(s: str) -> str: return s -def _xml_text(tag: str, xml: str) -> str: - m = re.search(rf"<{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]: - """解析 ARK pro 返回的 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: - name = _xml_attr("product", "name", raw) or "未识别" - brand = _xml_attr("product", "brand", raw) or "无法判断" - category = _xml_attr("product", "category", raw) or "无法判断" - appearance = _xml_attr("product", "appearance", raw) or "无法判断" - packaging = _xml_attr("product", "packaging", raw) or "无法判断" - feat = _xml_attr("product", "features", raw) 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 "" - 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 f"{brand} {name}" - 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 _assemble_pp_from_json(obj: dict[str, Any]) -> str: - """从 DashScope pro 返回的 JSON 组装 portrait_prompt。""" +def _assemble_pp(obj: dict[str, Any]) -> str: if not obj.get("has_person", False): return "无人像" parts: list[str] = [] @@ -159,29 +81,36 @@ def _assemble_pp_from_json(obj: dict[str, Any]) -> str: return ",".join(parts) if parts else "无人像" -def _call_dashscope_pro(img_url: str, idx: int, timeout: int) -> dict[str, Any] | None: - """DashScope qwen3.7-plus 兜底,直接 httpx 发精简 JSON prompt。""" - import httpx +def call_pro_vlm( + img_url: str, + idx: int, + *, + timeout: int = _DEFAULT_TIMEOUT, +) -> dict[str, Any] | None: + """单次调用 qwen3.7-plus,解析后返回 product dict;失败返回 None。""" t0 = time.time() - api_key = _provider.api_key() + import httpx + + api_key = os.environ.get("DASHSCOPE_API_KEY") if not api_key: - logger.warning("[vision.vlm] dashscope api_key 未配置") + logger.warning("[vision.v2] DASHSCOPE_API_KEY 未配置,跳过 pro 兜底") return None - url = f"{_provider.base_url().rstrip('/')}/chat/completions" + + url = f"{_BASE_URL}/chat/completions" payload: dict[str, Any] = { - "model": _provider.pro_model(), + "model": _PRO_MODEL, "messages": [ - {"role": "system", "content": _DS_PRO_SYSTEM}, + {"role": "system", "content": _PRO_SYSTEM}, {"role": "user", "content": [ {"type": "image_url", "image_url": {"url": img_url}}, - {"type": "text", "text": _DS_PRO_USER}, + {"type": "text", "text": _PRO_USER}, ]}, ], "temperature": 0.3, - "max_tokens": DEFAULT_MAX_TOKENS, + "max_tokens": _DEFAULT_MAX_TOKENS, "stream": False, + "enable_thinking": False, } - payload.update(_provider.thinking_params()) try: r = httpx.post( url, @@ -191,127 +120,49 @@ def _call_dashscope_pro(img_url: str, idx: int, timeout: int) -> dict[str, Any] ) elapsed = time.time() - t0 if r.status_code != 200: - logger.warning("[vision.vlm] dashscope pro HTTP %d elapsed=%.1fs body=%s", r.status_code, elapsed, r.text[:200]) + logger.warning("[vision.v2] pro HTTP %d elapsed=%.1fs body=%s", r.status_code, elapsed, r.text[:200]) return None data = r.json() raw = (data.get("choices") or [{}])[0].get("message", {}).get("content") if not raw: - logger.warning("[vision.vlm] dashscope pro 返回空 elapsed=%.1fs", elapsed) + logger.warning("[vision.v2] pro 返回空 elapsed=%.1fs", elapsed) return None usage = data.get("usage") or {} logger.info( - "[vision.vlm] dashscope pro 完成 idx=%d elapsed=%.1fs in=%d out=%d", - idx, elapsed, usage.get("prompt_tokens", 0), usage.get("completion_tokens", 0), + "[vision.v2] pro 完成 idx=%d model=%s elapsed=%.1fs in=%d out=%d", + idx, _PRO_MODEL, elapsed, + usage.get("prompt_tokens", 0), usage.get("completion_tokens", 0), ) text = _strip_code_fence(raw) l, r_pos = text.find("{"), text.rfind("}") if l < 0 or r_pos <= l: - logger.warning("[vision.vlm] dashscope pro 无JSON elapsed=%.1fs head=%s", elapsed, raw[:200]) + logger.warning("[vision.v2] pro 无JSON elapsed=%.1fs head=%s", elapsed, raw[:200]) return None obj = json.loads(text[l:r_pos+1]) if not isinstance(obj, dict): return None scene = obj.get("scene") or "通用" mood = obj.get("mood") or "" - pp = _assemble_pp_from_json(obj) + pp = _assemble_pp(obj) + has_person = obj.get("has_person", False) + has_product = obj.get("has_product", False) name = obj.get("product_name") or "未识别" brand = obj.get("brand") or "无法判断" - category = obj.get("category") or ("非产品图" if obj.get("has_person") else "无法判断") + category = obj.get("category") or ("非产品图" if has_person and not has_product else "无法判断") return { - "name": name, "brand": brand, "category": category, + "name": name, + "brand": brand, + "category": category, "appearance": obj.get("outfit") or "无法判断", "packaging": "无法判断", "text_on_package": [], "key_features": obj.get("key_features") or ["无法判断"], - "scene": scene, "mood": mood, "portrait_prompt": pp, + "scene": scene, + "mood": mood, + "portrait_prompt": pp, "summary": f"{brand} {name}" if name != "未识别" else "未识别", - "_source": "vlm_pro_dashscope_json", + "_source": "vlm_pro", } except Exception as e: - logger.warning("[vision.vlm] dashscope pro 异常 idx=%d elapsed=%.1fs err=%s", idx, time.time()-t0, e) + logger.warning("[vision.v2] pro 异常 idx=%d elapsed=%.1fs err=%s", idx, time.time()-t0, e, exc_info=True) return None - - -def _call_ark_pro(img_url: str, idx: int, model: str | None, timeout: int) -> dict[str, Any] | None: - """ARK 豆包 pro 兜底,保留 ai_client + prompt_loader + XML/JSON 双解析。""" - 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 _provider.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}], - images=[img_url], - temperature=0.3, - max_tokens=DEFAULT_MAX_TOKENS, - timeout=timeout, - model=use_model, - ) - except Exception as e: - logger.warning("[vision.vlm] ark pro 调用失败 idx=%d 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: - logger.warning("[vision.vlm] ark pro 返回空 elapsed=%.1fs", elapsed) - return None - text = _strip_code_fence(raw) - l, r_pos = text.find("{"), text.rfind("}") - if l >= 0 and r_pos > l: - try: - obj = json.loads(text[l:r_pos+1]) - if isinstance(obj, dict): - logger.info("[vision.vlm] ark pro JSON 完成 idx=%d 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: - return _xml_to_product(text, idx) - except Exception as e: - logger.warning("[vision.vlm] ark pro 解析失败 idx=%d elapsed=%.1fs err=%s head=%s", idx, elapsed, e, raw[:200]) - return None - - -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。""" - if _provider.is_dashscope(): - return _call_dashscope_pro(img_url, idx, timeout) - return _call_ark_pro(img_url, idx, model, timeout) 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 47e885d38..5fb245356 100644 --- a/apps/worker/worker_app/tasks/vision/vlm_fast_json.py +++ b/apps/worker/worker_app/tasks/vision/vlm_fast_json.py @@ -1,30 +1,32 @@ # -*- coding: utf-8 -*- -"""V2 快速路径 JSON-only VLM 调用。 - -支持双后端切换(通过环境变量 VISION_V2_PROVIDER=ark|dashscope): -- ark(火山方舟/豆包):doubao-seed-2-1-lite-260915,thinking={type:"disabled"} -- dashscope(阿里云百炼/qwen,默认对比测试):qwen3.8-flash,enable_thinking=false +"""V2 快速路径:qwen3.8-flash(阿里云百炼/DashScope)强约束 JSON-only 调用。 +目标:替代"人体属性/商品检测/图像标签"三个火山不存在的专用云端 API。 设计要点: -- 直接用 httpx 发最小 payload,不走 ai_client 包装 -- 关闭 thinking/推理链(reasoning_tokens 是延迟主因) -- system prompt 极致精简,只给字段 schema 和强约束 +- 直接用 httpx 发最小 payload 到 DashScope OpenAI 兼容 endpoint,不走 ai_client 包装 +- enable_thinking=false 关闭推理链(reasoning 是延迟主因) +- system prompt 极致精简,只给字段 schema 和强约束(禁止自然语言、禁止 markdown) - max_tokens=350、temperature=0.1(稳定输出 JSON) -- 单次调用不重试(失败由外层走 pro 兜底) -- 期望返回纯 JSON object(无 ```json 包裹、无解释文字) +- timeout=8s(失败由外层走 pro 兜底) +- API Key 从环境变量 DASHSCOPE_API_KEY 读取 """ from __future__ import annotations import json import logging +import os import time from typing import Any -from . import _provider - logger = logging.getLogger(__name__) +# DashScope OpenAI 兼容 endpoint +_BASE_URL = "https://dashscope.aliyuncs.com/compatible-mode/v1" +_FAST_MODEL = "qwen3.8-flash" +_DEFAULT_TIMEOUT = 8 +_DEFAULT_MAX_TOKENS = 350 + # 极简 system prompt:只给字段定义 + 硬性输出要求 _FAST_SYSTEM = ( "你是图片结构化识别器。严格按下方 JSON schema 返回一个对象,不要任何解释、" @@ -57,8 +59,9 @@ _FAST_SYSTEM = ( _FAST_USER = "识别这张图片的人物穿搭与主体信息,只返回JSON对象。" -DEFAULT_TIMEOUT = _provider.fast_timeout_default() -DEFAULT_MAX_TOKENS = 350 + +def _api_key() -> str | None: + return os.environ.get("DASHSCOPE_API_KEY") def _strip_code_fence(s: str) -> str: @@ -76,42 +79,37 @@ def _strip_code_fence(s: str) -> str: def call_fast_json( img_url: str, *, - model: str | None = None, - timeout: int = DEFAULT_TIMEOUT, - max_tokens: int = DEFAULT_MAX_TOKENS, + timeout: int = _DEFAULT_TIMEOUT, + max_tokens: int = _DEFAULT_MAX_TOKENS, ) -> dict[str, Any] | None: - """调用 fast VLM 返回结构化 dict;失败/非 JSON 返回 None。""" + """调用 qwen3.8-flash 返回结构化 dict;失败/非 JSON 返回 None。""" t0 = time.time() import httpx + api_key = _api_key() + if not api_key: + logger.warning("[vision.v2] DASHSCOPE_API_KEY 未配置,跳过 fast_json") + return None + + url = f"{_BASE_URL}/chat/completions" + payload: dict[str, Any] = { + "model": _FAST_MODEL, + "messages": [ + {"role": "system", "content": _FAST_SYSTEM}, + { + "role": "user", + "content": [ + {"type": "image_url", "image_url": {"url": img_url}}, + {"type": "text", "text": _FAST_USER}, + ], + }, + ], + "temperature": 0.1, + "max_tokens": max_tokens, + "stream": False, + "enable_thinking": False, + } try: - api_key = _provider.api_key() - base_url = _provider.base_url().rstrip("/") - if not api_key: - logger.warning("[vision.v2] provider=%s api_key 未配置,跳过 fast_json", _provider.get_provider()) - return None - - use_model = model or _provider.fast_model() - url = f"{base_url}/chat/completions" - payload: dict[str, Any] = { - "model": use_model, - "messages": [ - {"role": "system", "content": _FAST_SYSTEM}, - { - "role": "user", - "content": [ - {"type": "image_url", "image_url": {"url": img_url}}, - {"type": "text", "text": _FAST_USER}, - ], - }, - ], - "temperature": 0.1, - "max_tokens": max_tokens, - "stream": False, - } - # 按 provider 设置关 thinking 参数 - payload.update(_provider.thinking_params()) - resp = httpx.post( url, headers={"Authorization": f"Bearer {api_key}", "Content-Type": "application/json"}, @@ -119,35 +117,24 @@ def call_fast_json( timeout=timeout, ) elapsed = time.time() - t0 - if resp.status_code == 400: - logger.warning( - "[vision.v2] fast_json HTTP 400 provider=%s elapsed=%.1fs body=%s", - _provider.get_provider(), elapsed, resp.text[:200], - ) - if _provider.is_400_thinking_error(resp.text[:500]): - _provider.pop_thinking_param(payload) - resp = httpx.post( - url, - headers={"Authorization": f"Bearer {api_key}", "Content-Type": "application/json"}, - json=payload, - timeout=timeout, - ) - elapsed = time.time() - t0 - if resp.status_code != 200: - logger.warning("[vision.v2] fast_json 降级重试 HTTP %d elapsed=%.1fs", resp.status_code, elapsed) - return None - else: - return None - elif resp.status_code != 200: - logger.warning( - "[vision.v2] fast_json HTTP %d provider=%s elapsed=%.1fs body=%s", - resp.status_code, _provider.get_provider(), elapsed, resp.text[:200], + if resp.status_code == 400 and "enable_thinking" in resp.text[:300].lower(): + # 极少数 endpoint 版本不识别 enable_thinking,重试一次不带 + logger.warning("[vision.v2] fast_json HTTP 400 thinking 参数不兼容,重试 elapsed=%.1fs", elapsed) + payload.pop("enable_thinking", None) + resp = httpx.post( + url, + headers={"Authorization": f"Bearer {api_key}", "Content-Type": "application/json"}, + json=payload, + timeout=timeout, ) + elapsed = time.time() - t0 + if resp.status_code != 200: + logger.warning("[vision.v2] fast_json HTTP %d elapsed=%.1fs body=%s", resp.status_code, elapsed, resp.text[:200]) return None data = resp.json() raw = (data.get("choices") or [{}])[0].get("message", {}).get("content") - if raw is None: - logger.warning("[vision.v2] fast_json 返回 None elapsed=%.1fs", elapsed) + if not raw: + logger.warning("[vision.v2] fast_json 返回空 elapsed=%.1fs", elapsed) return None usage = data.get("usage") or {} reasoning_tokens = usage.get("reasoning_tokens", 0) @@ -155,13 +142,12 @@ def call_fast_json( if not reasoning_tokens: reasoning_tokens = ctd.get("reasoning_tokens", 0) logger.info( - "[vision.v2] fast_json 完成 provider=%s model=%s elapsed=%.1fs in=%d out=%d reasoning=%d", - _provider.get_provider(), use_model, elapsed, + "[vision.v2] fast_json 完成 model=%s elapsed=%.1fs in=%d out=%d reasoning=%d", + _FAST_MODEL, elapsed, usage.get("prompt_tokens", 0), usage.get("completion_tokens", 0), reasoning_tokens, ) text = _strip_code_fence(raw) - l = text.find("{") - r = text.rfind("}") + l, r = text.find("{"), text.rfind("}") if l >= 0 and r > l: text = text[l : r + 1] try: @@ -173,12 +159,11 @@ def call_fast_json( logger.warning("[vision.v2] fast_json 非 dict: %s", type(obj)) return None logger.info( - "[vision.v2] fast_json 完成 provider=%s model=%s elapsed=%.1fs has_person=%s has_product=%s category=%s", - _provider.get_provider(), use_model, elapsed, - obj.get("has_person"), obj.get("has_product"), obj.get("category"), + "[vision.v2] fast_json 完成 elapsed=%.1fs has_person=%s has_product=%s category=%s", + elapsed, obj.get("has_person"), obj.get("has_product"), obj.get("category"), ) return obj except Exception as e: elapsed = time.time() - t0 - logger.warning("[vision.v2] fast_json 异常 provider=%s elapsed=%.1fs err=%s", _provider.get_provider(), elapsed, e, exc_info=True) + logger.warning("[vision.v2] fast_json 异常 elapsed=%.1fs err=%s", elapsed, e, exc_info=True) return None