From cdd343131ede9048a0efe5c9f14eab628590b67e Mon Sep 17 00:00:00 2001 From: xiaoxia Date: Mon, 5 Oct 2026 19:06:02 +0800 Subject: [PATCH 1/3] =?UTF-8?q?fix(vision):=20#2205=20thinking=E5=8F=82?= =?UTF-8?q?=E6=95=B0=E4=BA=92=E6=96=A5=E4=BF=AE=E5=A4=8D=E2=80=94=E2=80=94?= =?UTF-8?q?=E5=8F=AA=E4=BC=A0thinking=3Ddisabled=EF=BC=8C=E5=8E=BB?= =?UTF-8?q?=E6=8E=89reasoning=5Feffort?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 根因:上版同时传 thinking={type:disabled} + reasoning_effort=low,方舟返回400 Invalid combination of reasoning_effort and thinking type。降级重试分支 把两个参数都pop掉,模型回到默认thinking开启→响应9-12s超时,fast全败。 修复: 1. 只传 thinking={"type":"disabled"},去掉互斥的 reasoning_effort 2. 400降级只pop thinking,保留精简payload重试(不pop多个) 3. 先单独验证 lite 单次HTTP 200且reasoning_tokens=0再跑E2E,省一轮部署 --- .../worker_app/tasks/vision/vlm_fast_json.py | 29 ++++++++----------- 1 file changed, 12 insertions(+), 17 deletions(-) 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 036ec8285..b67409a02 100644 --- a/apps/worker/worker_app/tasks/vision/vlm_fast_json.py +++ b/apps/worker/worker_app/tasks/vision/vlm_fast_json.py @@ -117,14 +117,9 @@ def call_fast_json( "max_tokens": max_tokens, "stream": False, } - # 关键:关闭 thinking(避免产生 reasoning_tokens 拖慢响应) - # 方舟/豆包 2.x 模型支持 thinking.type=disabled - try: - payload["thinking"] = {"type": "disabled"} - except Exception: - pass - # 部分模型用 reasoning_effort 控制思考深度 - payload["reasoning_effort"] = "low" + # 关键:关闭 thinking(reasoning_tokens 是延迟主因,单次要10-12s) + # 方舟/豆包 Seed 2.x 支持 thinking={type:"disabled"},且不要和 reasoning_effort 同时传(两者互斥会400) + payload["thinking"] = {"type": "disabled"} resp = httpx.post( url, @@ -133,15 +128,12 @@ def call_fast_json( 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] - ) - # 如果400说明不支持thinking参数,降级重试一次 - if resp.status_code == 400 and "thinking" in resp.text.lower(): + # 400 说明模型不支持 thinking 参数(极少数旧模型),重试一次不带 thinking + if resp.status_code == 400: + body_preview = resp.text[:300].lower() + logger.warning("[vision.v2] fast_json HTTP 400 elapsed=%.1fs body=%s", elapsed, resp.text[:200]) + if "thinking" in body_preview or "reasoning" in body_preview: payload.pop("thinking", None) - payload.pop("reasoning_effort", None) - time.time() resp = httpx.post( url, headers={"Authorization": f"Bearer {api_key}", "Content-Type": "application/json"}, @@ -150,10 +142,13 @@ def call_fast_json( ) elapsed = time.time() - t0 if resp.status_code != 200: - logger.warning("[vision.v2] fast_json 降级后 HTTP %d elapsed=%.1fs", resp.status_code, elapsed) + 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 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: -- 2.54.0 From aac785a9ed3d8ca59ee5aa5bbb84e860345a2767 Mon Sep 17 00:00:00 2001 From: xiaoxia Date: Mon, 5 Oct 2026 19:42:11 +0800 Subject: [PATCH 2/3] =?UTF-8?q?feat(vision):=20=E6=96=B0=E5=A2=9EDashScope?= =?UTF-8?q?(=E9=98=BF=E9=87=8C=E4=BA=91=E7=99=BE=E7=82=BC/qwen)=20provider?= =?UTF-8?q?=E6=94=AF=E6=8C=81?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - 新增_provider.py统一后端切换:VISION_V2_PROVIDER=ark|dashscope(默认dashscope对比测试) - vlm_fast_json.py:双后端httpx直连,ark用thinking={type:disabled}、dashscope用enable_thinking=false,400降级逻辑按provider处理 - vlm_fallback.py:dashscope用qwen3.7-plus+精简JSON prompt(<20s),ark保留原有ai_client+XML/JSON双解析 - fast_path.py:超时默认值从provider读取(dashscope pro 20s、ark pro 45s) - 通过环境变量DASHSCOPE_API_KEY配置key(敏感信息不入代码) - 本地curl验证:qwen3.8-flash关thinking单图2.6s、3图并发7.3s、8图并发10.3s(7/8可用) --- .../worker_app/tasks/vision/_provider.py | 112 ++++++++ .../worker_app/tasks/vision/fast_path.py | 33 +-- .../worker_app/tasks/vision/vlm_fallback.py | 269 ++++++++++++------ .../worker_app/tasks/vision/vlm_fast_json.py | 107 +++---- 4 files changed, 351 insertions(+), 170 deletions(-) create mode 100644 apps/worker/worker_app/tasks/vision/_provider.py diff --git a/apps/worker/worker_app/tasks/vision/_provider.py b/apps/worker/worker_app/tasks/vision/_provider.py new file mode 100644 index 000000000..c3df59132 --- /dev/null +++ b/apps/worker/worker_app/tasks/vision/_provider.py @@ -0,0 +1,112 @@ +# -*- 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 323ee44ea..45bc75b1f 100644 --- a/apps/worker/worker_app/tasks/vision/fast_path.py +++ b/apps/worker/worker_app/tasks/vision/fast_path.py @@ -1,7 +1,8 @@ # -*- coding: utf-8 -*- -"""V2 图片分析主路径:每图并行 OCR(火山专用API)+ lite JSON VLM,失败时单次 pro VLM 兜底。 +"""V2 图片分析主路径:每图并行 OCR(火山专用API,未配置时自动跳过)+ lite JSON VLM,失败时单次 pro VLM 兜底。 -设计原则(灵应10-05要求): +支持双后端(环境变量 VISION_V2_PROVIDER=ark|dashscope,默认 dashscope)。 +设计原则: - 主力路径简洁:单图2路并行,外层N图全并发 - 兜底简单:单次 pro VLM 调用,无竞速/重试/复杂超时 - 输出 dict 格式与旧版完全一致,下游零改动 @@ -15,16 +16,16 @@ import time from concurrent.futures import ThreadPoolExecutor, as_completed from typing import Any -from . import assembler, ocr_volc, vlm_fallback, vlm_fast_json +from . import _provider, 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", "8")) +_FAST_JSON_TIMEOUT = float(os.environ.get("VISION_V2_FAST_JSON_TIMEOUT", str(_provider.fast_timeout_default()))) _OCR_TIMEOUT = float(os.environ.get("VISION_V2_OCR_TIMEOUT", "6")) -_PRO_TIMEOUT = float(os.environ.get("VISION_V2_PRO_TIMEOUT", "45")) +_PRO_TIMEOUT = float(os.environ.get("VISION_V2_PRO_TIMEOUT", str(_provider.pro_timeout_default()))) _FALLBACK_RESULT = { "name": "未识别", @@ -42,7 +43,6 @@ _FALLBACK_RESULT = { 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 @@ -53,10 +53,8 @@ def _is_usable(r: dict[str, Any]) -> bool: def analyze_image_v2(idx: int, img_url: str) -> dict[str, Any]: - """单张图片 V2 分析。""" t0 = time.time() - # 第1层:OCR + lite JSON VLM 并行 fj_result: dict[str, Any] | None = None ocr_result: list[str] = [] with ThreadPoolExecutor(max_workers=2) as pool: @@ -74,7 +72,6 @@ def analyze_image_v2(idx: int, img_url: str) -> dict[str, Any]: 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() @@ -82,20 +79,17 @@ def analyze_image_v2(idx: int, img_url: str) -> dict[str, Any]: fast_elapsed = time.time() - t0 - # 组装 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 pp=%s", - idx, - fast_elapsed, + "[vision.v2] 图片 #%d fast命中 provider=%s elapsed=%.2fs pp=%s", + idx, _provider.get_provider(), fast_elapsed, (assembled.get("portrait_prompt") or "")[:40], ) return assembled - # 第2层:pro VLM 单次兜底 pro_t0 = time.time() pro_result = vlm_fallback.call_pro_vlm(img_url, idx, timeout=_PRO_TIMEOUT) if pro_result and _is_usable(pro_result): @@ -107,8 +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 全路径失败 elapsed=%.2fs", idx, time.time() - t0) + logger.warning("[vision.v2] 图片 #%d 全路径失败 provider=%s elapsed=%.2fs", idx, _provider.get_provider(), time.time() - t0) out = dict(_FALLBACK_RESULT) out["_source"] = "v2_all_failed" out["text_on_package"] = ocr_result[:8] @@ -117,13 +110,15 @@ def analyze_image_v2(idx: int, img_url: str) -> dict[str, Any]: def analyze_images_v2(img_urls: list[str]) -> list[dict[str, Any]]: - """批量图片 V2 分析,外层全并发。""" if not img_urls: return [] workers = min(_IMG_WORKERS, len(img_urls), 16) results: list[dict[str, Any] | None] = [None] * len(img_urls) - logger.info("[vision.v2] 开始图片分析 n=%d workers=%d fast_timeout=%.0fs", len(img_urls), workers, _FAST_TIMEOUT) + 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, + ) t0 = time.time() with ThreadPoolExecutor(max_workers=workers) as pool: future_to_idx = {pool.submit(analyze_image_v2, idx, url): idx for idx, url in enumerate(img_urls)} diff --git a/apps/worker/worker_app/tasks/vision/vlm_fallback.py b/apps/worker/worker_app/tasks/vision/vlm_fallback.py index 558500c17..89715dd1c 100644 --- a/apps/worker/worker_app/tasks/vision/vlm_fallback.py +++ b/apps/worker/worker_app/tasks/vision/vlm_fallback.py @@ -1,7 +1,10 @@ # -*- coding: utf-8 -*- -"""VLM 兜底:专用API路径失败时的最后一道防线,单次调用 doubao-seed-2.1-pro。 +"""VLM pro 兜底:fast 路径失败时的单次调用,支持双后端。 -设计原则:简单、直接、无竞速、无复杂超时逻辑。只在 fast_json 结果不可用时调用。 +- ark(火山/豆包):通过 ai_client 走 doubao-seed-2-1-pro,保留 XML+JSON 双解析和 prompt_loader +- dashscope(百炼/qwen):直接 httpx 走 qwen3.7-plus,用精简 JSON-only prompt 提升速度 + +设计原则:简单、直接、无竞速、无复杂重试。 """ from __future__ import annotations @@ -12,12 +15,36 @@ import re import time from typing import Any +from . import _provider + logger = logging.getLogger(__name__) -DEFAULT_PRO_MODEL = "doubao-seed-2-1-pro-260915" -DEFAULT_TIMEOUT = 45 +DEFAULT_TIMEOUT = _provider.pro_timeout_default() DEFAULT_MAX_TOKENS = 800 +# DashScope 兜底用的精简 JSON-only prompt(比 prompt_loader 模板短很多,减少延迟) +_DS_PRO_SYSTEM = ( + "你是图片分析助手。仔细观察图片,严格按JSON schema返回一个对象,不要任何解释、" + "不要markdown、不要代码块、不要前后缀文字。字段值不确定时填null或空数组。\n" + "{\n" + ' "has_person": true/false,\n' + ' "gender": "男"/"女"/null,\n' + ' "age_range": "儿童"/"青少年"/"青年"/"中年"/"老年"/null,\n' + ' "outfit": "人物穿搭描述,60字以内(例:白色T恤+牛仔裤)",\n' + ' "hair": "发型",\n' + ' "pose": "姿态",\n' + ' "expression": "表情",\n' + ' "scene": "场景",\n' + ' "mood": "氛围",\n' + ' "has_product": true/false,\n' + ' "category": "类目",\n' + ' "product_name": "产品名称",\n' + ' "brand": "品牌",\n' + ' "key_features": ["特征数组"]\n' + "}" +) +_DS_PRO_USER = "分析这张图片,返回符合schema的JSON。" + def _strip_code_fence(s: str) -> str: s = s.strip() @@ -42,7 +69,7 @@ def _xml_attr(tag: str, attr: str, xml: str) -> str: def _xml_to_product(raw: str, idx: int) -> dict[str, Any]: - """解析 VLM 输出的 XML 格式(简化版)。""" + """解析 ARK pro 返回的 XML 格式。""" scene = _xml_text("scene", raw) or "通用" mood = _xml_text("mood", raw) or "" @@ -71,86 +98,151 @@ def _xml_to_product(raw: str, idx: int) -> dict[str, Any]: 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 "" + 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 _xml_text("text_on_package", pbody) or "" + 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 _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 + 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", + "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", + "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", + "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。""" +def _assemble_pp_from_json(obj: dict[str, Any]) -> str: + """从 DashScope pro 返回的 JSON 组装 portrait_prompt。""" + if not obj.get("has_person", False): + return "无人像" + parts: list[str] = [] + gender = obj.get("gender") + age = obj.get("age_range") + if gender: + parts.append(gender + ("性" if not gender.endswith("性") else "")) + if age: + parts.append(age) + parts.append("人物") + hair = obj.get("hair") + if hair: + parts.append(hair) + outfit = obj.get("outfit") + if outfit: + parts.append(f"身着{outfit}") + pose = obj.get("pose") + if pose: + parts.append(f"姿态{pose}") + expr = obj.get("expression") + if expr: + parts.append(f"表情{expr}") + 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 + t0 = time.time() + api_key = _provider.api_key() + if not api_key: + logger.warning("[vision.vlm] dashscope api_key 未配置") + return None + url = f"{_provider.base_url().rstrip('/')}/chat/completions" + payload: dict[str, Any] = { + "model": _provider.pro_model(), + "messages": [ + {"role": "system", "content": _DS_PRO_SYSTEM}, + {"role": "user", "content": [ + {"type": "image_url", "image_url": {"url": img_url}}, + {"type": "text", "text": _DS_PRO_USER}, + ]}, + ], + "temperature": 0.3, + "max_tokens": DEFAULT_MAX_TOKENS, + "stream": False, + } + payload.update(_provider.thinking_params()) + try: + r = httpx.post( + url, + headers={"Authorization": f"Bearer {api_key}", "Content-Type": "application/json"}, + json=payload, + timeout=timeout, + ) + 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]) + 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) + 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), + ) + 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]) + 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) + name = obj.get("product_name") or "未识别" + brand = obj.get("brand") or "无法判断" + category = obj.get("category") or ("非产品图" if obj.get("has_person") else "无法判断") + return { + "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, + "summary": f"{brand} {name}" if name != "未识别" else "未识别", + "_source": "vlm_pro_dashscope_json", + } + except Exception as e: + logger.warning("[vision.vlm] dashscope pro 异常 idx=%d elapsed=%.1fs err=%s", idx, time.time()-t0, e) + 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, + 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) @@ -158,12 +250,10 @@ def call_pro_vlm( 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 + use_model = model or _provider.pro_model() _orig_retries = client.max_retries client.max_retries = 0 try: @@ -176,23 +266,21 @@ def call_pro_vlm( model=use_model, ) except Exception as e: - logger.warning("[vision.vlm] 图片 #%d pro VLM 调用失败 elapsed=%.1fs err=%s", idx, time.time() - t0, 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] 图片 #%d pro VLM 返回空 elapsed=%.1fs", idx, elapsed) + logger.warning("[vision.vlm] ark pro 返回空 elapsed=%.1fs", elapsed) return None - text = _strip_code_fence(raw) - l, r = text.find("{"), text.rfind("}") - if l >= 0 and r > l: + l, r_pos = text.find("{"), text.rfind("}") + if l >= 0 and r_pos > l: try: - obj = json.loads(text[l : r + 1]) + obj = json.loads(text[l:r_pos+1]) if isinstance(obj, dict): - logger.info("[vision.vlm] 图片 #%d pro VLM JSON 完成 elapsed=%.1fs", idx, elapsed) + logger.info("[vision.vlm] ark pro JSON 完成 idx=%d elapsed=%.1fs", idx, elapsed) return { "name": obj.get("name") or "未识别", "brand": obj.get("brand") or "无法判断", @@ -209,18 +297,21 @@ def call_pro_vlm( } 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 + return _xml_to_product(text, idx) except Exception as e: - logger.warning("[vision.vlm] 图片 #%d 解析失败 elapsed=%.1fs err=%s head=%s", idx, elapsed, e, raw[:200]) + 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 b67409a02..47e885d38 100644 --- a/apps/worker/worker_app/tasks/vision/vlm_fast_json.py +++ b/apps/worker/worker_app/tasks/vision/vlm_fast_json.py @@ -1,12 +1,16 @@ # -*- coding: utf-8 -*- -"""doubao-seed-2.1-lite 强约束 JSON-only 调用。 +"""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 -目标:替代"人体属性/商品检测/图像标签"三个火山不存在的专用云端 API。 设计要点: -- system prompt 极致精简,只给字段 schema 和强约束(禁止自然语言、禁止 markdown) -- max_tokens=350(比旧 VLM 的 1200 小很多,降低延迟) -- temperature=0.1(极低,稳定输出 JSON) -- timeout=8s(够快,失败则由外层走 pro VLM 兜底) +- 直接用 httpx 发最小 payload,不走 ai_client 包装 +- 关闭 thinking/推理链(reasoning_tokens 是延迟主因) +- system prompt 极致精简,只给字段 schema 和强约束 +- max_tokens=350、temperature=0.1(稳定输出 JSON) +- 单次调用不重试(失败由外层走 pro 兜底) - 期望返回纯 JSON object(无 ```json 包裹、无解释文字) """ @@ -17,6 +21,8 @@ import logging import time from typing import Any +from . import _provider + logger = logging.getLogger(__name__) # 极简 system prompt:只给字段定义 + 硬性输出要求 @@ -24,12 +30,12 @@ _FAST_SYSTEM = ( "你是图片结构化识别器。严格按下方 JSON schema 返回一个对象,不要任何解释、" "不要markdown、不要代码块、不要前后缀文字。字段值不确定时填 null 或空数组。\n" "{\n" - ' "has_person": true/false, // 图中是否有人\n' + ' "has_person": true/false,\n' ' "gender": "男"/"女"/null,\n' ' "age_range": "儿童"/"青少年"/"青年"/"中年"/"老年"/null,\n' ' "upper_wear": "上装款式,如T恤/衬衫/卫衣/毛衣/西装/夹克/连衣裙/吊带/背心/外套等",\n' ' "upper_color": "上装主色",\n' - ' "lower_wear": "下装款式,如牛仔裤/休闲裤/短裙/长裙/短裤/西裤/运动裤等;穿连衣裙时填null",\n' + ' "lower_wear": "下装款式;穿连衣裙时填null",\n' ' "lower_color": "下装主色",\n' ' "dress_color": "连衣裙主色(穿连衣裙时填)",\n' ' "accessories": ["眼镜"/"帽子"/"项链"/"耳环"/"背包"/"手表"等数组],\n' @@ -38,7 +44,7 @@ _FAST_SYSTEM = ( ' "pose": "姿势,如站立/坐姿/侧身/行走等",\n' ' "scene": "场景,如室内/街拍/户外/办公室/家居/海边/雪景/森林等",\n' ' "style": "风格,如休闲/商务/运动/复古/潮流/甜美/酷飒/优雅/街头/法式等",\n' - ' "has_product": true/false, // 是否有明确商品展示\n' + ' "has_product": true/false,\n' ' "category": "产品类目:服饰/鞋包/美妆/数码/食品/家居/配饰/母婴/非产品图",\n' ' "product_name": "产品名称,非产品图填null",\n' ' "brand": "品牌或文字标识,无则null",\n' @@ -51,21 +57,16 @@ _FAST_SYSTEM = ( _FAST_USER = "识别这张图片的人物穿搭与主体信息,只返回JSON对象。" -# 默认模型 -DEFAULT_LITE_MODEL = "doubao-seed-2-1-lite-260915" -DEFAULT_TIMEOUT = 8 +DEFAULT_TIMEOUT = _provider.fast_timeout_default() DEFAULT_MAX_TOKENS = 350 def _strip_code_fence(s: str) -> str: - """剥离 ```json ... ``` 包裹(即使要求纯 JSON,模型偶尔仍会包代码块)。""" s = s.strip() if s.startswith("```"): lines = s.split("\n") - # 去掉首行 ```json if lines and lines[0].startswith("```"): lines = lines[1:] - # 去掉尾行 ``` if lines and lines[-1].strip().startswith("```"): lines = lines[:-1] s = "\n".join(lines).strip() @@ -79,27 +80,18 @@ def call_fast_json( timeout: int = DEFAULT_TIMEOUT, max_tokens: int = DEFAULT_MAX_TOKENS, ) -> dict[str, Any] | None: - """调用 lite VLM 返回结构化 dict;失败/非 JSON 返回 None。 - - 直接用 httpx 发最小 payload(关闭 thinking),不走 ai_client 包装: - - 关闭 thinking/推理链(reasoning_tokens 是延迟主因,单次要10-12s) - - 单次调用不重试(失败由外层走 pro 兜底) - - 温度=0.1 稳定输出 JSON - """ + """调用 fast VLM 返回结构化 dict;失败/非 JSON 返回 None。""" t0 = time.time() import httpx try: - from packages.shared import get_shared_settings - - settings = get_shared_settings() - api_key = settings.doubao_api_key - base_url = (settings.doubao_base_url or "https://ark.cn-beijing.volces.com/api/v3").rstrip("/") + api_key = _provider.api_key() + base_url = _provider.base_url().rstrip("/") if not api_key: - logger.warning("[vision.v2] doubao api_key 未配置,跳过 fast_json") + logger.warning("[vision.v2] provider=%s api_key 未配置,跳过 fast_json", _provider.get_provider()) return None - use_model = model or DEFAULT_LITE_MODEL + use_model = model or _provider.fast_model() url = f"{base_url}/chat/completions" payload: dict[str, Any] = { "model": use_model, @@ -117,9 +109,8 @@ def call_fast_json( "max_tokens": max_tokens, "stream": False, } - # 关键:关闭 thinking(reasoning_tokens 是延迟主因,单次要10-12s) - # 方舟/豆包 Seed 2.x 支持 thinking={type:"disabled"},且不要和 reasoning_effort 同时传(两者互斥会400) - payload["thinking"] = {"type": "disabled"} + # 按 provider 设置关 thinking 参数 + payload.update(_provider.thinking_params()) resp = httpx.post( url, @@ -128,12 +119,13 @@ def call_fast_json( timeout=timeout, ) elapsed = time.time() - t0 - # 400 说明模型不支持 thinking 参数(极少数旧模型),重试一次不带 thinking if resp.status_code == 400: - body_preview = resp.text[:300].lower() - logger.warning("[vision.v2] fast_json HTTP 400 elapsed=%.1fs body=%s", elapsed, resp.text[:200]) - if "thinking" in body_preview or "reasoning" in body_preview: - payload.pop("thinking", None) + 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"}, @@ -147,7 +139,10 @@ def call_fast_json( else: return None elif resp.status_code != 200: - logger.warning("[vision.v2] fast_json HTTP %d elapsed=%.1fs body=%s", resp.status_code, elapsed, resp.text[: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], + ) return None data = resp.json() raw = (data.get("choices") or [{}])[0].get("message", {}).get("content") @@ -155,21 +150,16 @@ def call_fast_json( logger.warning("[vision.v2] fast_json 返回 None elapsed=%.1fs", elapsed) return None usage = data.get("usage") or {} + reasoning_tokens = usage.get("reasoning_tokens", 0) + ctd = usage.get("completion_tokens_details") or {} + if not reasoning_tokens: + reasoning_tokens = ctd.get("reasoning_tokens", 0) logger.info( - "[vision.v2] fast_json 直连完成 model=%s elapsed=%.1fs in=%d out=%d reasoning=%d", - use_model, - elapsed, - usage.get("prompt_tokens", 0), - usage.get("completion_tokens", 0), - usage.get("reasoning_tokens", 0), + "[vision.v2] fast_json 完成 provider=%s model=%s elapsed=%.1fs in=%d out=%d reasoning=%d", + _provider.get_provider(), use_model, elapsed, + usage.get("prompt_tokens", 0), usage.get("completion_tokens", 0), reasoning_tokens, ) - elapsed = time.time() - t0 - if raw is None: - logger.warning("[vision.v2] fast_json 返回 None elapsed=%.1fs model=%s", elapsed, use_model) - return None - text = _strip_code_fence(raw) - # 截到第一个 { 和最后一个 } 之间,容忍前后偶发文字 l = text.find("{") r = text.rfind("}") if l >= 0 and r > l: @@ -177,25 +167,18 @@ def call_fast_json( try: obj = json.loads(text) except json.JSONDecodeError: - logger.warning( - "[vision.v2] fast_json JSON 解析失败 elapsed=%.1fs head=%s", - elapsed, - raw[:200], - ) + logger.warning("[vision.v2] fast_json JSON 解析失败 elapsed=%.1fs head=%s", elapsed, raw[:200]) return None if not isinstance(obj, dict): logger.warning("[vision.v2] fast_json 非 dict: %s", type(obj)) return None logger.info( - "[vision.v2] fast_json 完成 model=%s elapsed=%.1fs has_person=%s has_product=%s category=%s", - use_model, - elapsed, - obj.get("has_person"), - obj.get("has_product"), - obj.get("category"), + "[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"), ) return obj except Exception as e: elapsed = time.time() - t0 - logger.warning("[vision.v2] fast_json 异常 elapsed=%.1fs err=%s", elapsed, e, exc_info=True) + logger.warning("[vision.v2] fast_json 异常 provider=%s elapsed=%.1fs err=%s", _provider.get_provider(), elapsed, e, exc_info=True) return None -- 2.54.0 From 257921abf6cf6eb618d160c634abc71faeec309d Mon Sep 17 00:00:00 2001 From: xiaoxia Date: Mon, 5 Oct 2026 19:52:14 +0800 Subject: [PATCH 3/3] =?UTF-8?q?refactor(vision):=20=E5=94=AF=E4=B8=80?= =?UTF-8?q?=E5=90=8E=E7=AB=AFDashScope(qwen)=EF=BC=8C=E5=88=A0=E9=99=A4ARK?= =?UTF-8?q?/doubao=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 -- 2.54.0