refactor(vision): 唯一后端DashScope(qwen),删除ARK/doubao和provider切换
CI/CD Pipeline / Dedup Check - skip PR tests when covered by push pipeline (pull_request) Successful in 5s
CI/CD Pipeline / Check if frontend-only change (pull_request) Successful in 2s
CI/CD Pipeline / Check push changed paths (pull_request) Has been skipped
CI/CD Pipeline / Build Staging API Image (pull_request) Has been skipped
CI/CD Pipeline / Build Staging Web Image (pull_request) Has been skipped
PR Automation / Auto Approve on CI Green (pull_request) Successful in 2m59s
CI/CD Pipeline / Frontend Lint (pull_request) Has been skipped
CI/CD Pipeline / Frontend Unit Tests (pull_request) Has been skipped
CI/CD Pipeline / PR Build API Image (pull_request) Successful in 16s
CI/CD Pipeline / PR Build Web Image (pull_request) Has been skipped
Preview Deploy / Deploy Preview Environment (pull_request) Successful in 2m28s
CI/CD Pipeline / Build Staging Worker Image (pull_request) Has been skipped
CI/CD Pipeline / PR Build Worker Image (pull_request) Successful in 1m38s
AI Code Review / AI Code Review (pull_request) Successful in 7m13s
CI/CD Pipeline / Retag skipped Staging API Image (pull_request) Has been skipped
CI/CD Pipeline / Retag skipped Staging Web Image (pull_request) Has been skipped
CI/CD Pipeline / Retag skipped Staging Worker Image (pull_request) Has been skipped
PR Automation / Auto Merge on CI Green + Approved (pull_request) Successful in 3m54s
CI/CD Pipeline / Deploy Staging (Watchtower auto-deploy) (pull_request) Has been skipped
Preview Cleanup / Cleanup Preview Environment (pull_request) Successful in 1m30s
CI/CD Pipeline / Unit Tests (pull_request) Failing after 10m42s
CI/CD Pipeline / Integration Tests (pull_request) Successful in 10m38s
ACR Cleanup / ACR Image Cleanup (pull_request_target) Successful in 22s
CI/CD Pipeline / Validate - Style (pull_request) Failing after 12m20s
CI/CD Pipeline / Staging E2E Tests (pull_request) Has been skipped
CI/CD Pipeline / Staging API Integration Tests (pull_request) Has been skipped
CI/CD Pipeline / ACR Image Cleanup (pull_request) Has been skipped
CI/CD Pipeline / Validate - Python (mypy + alembic) (pull_request) Successful in 13m59s
CI/CD Pipeline / Validate - Security (pull_request) Successful in 38m41s
CI/CD Pipeline / Build Production Web Image (pull_request) Has been skipped
CI/CD Pipeline / Build Production Worker Image (pull_request) Has been skipped
CI/CD Pipeline / Build Production API Image (pull_request) Has been skipped
CI/CD Pipeline / Deploy Production (pull_request) Has been skipped
CI/CD Pipeline / Canary Release to Production (pull_request) Has been skipped
CI/CD Pipeline / Production Browser E2E (pull_request) Has been skipped
CI/CD Pipeline / CI Gate (pull_request) Failing after 2s

按灵应指示清理冗余代码:
- 删除_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)
This commit is contained in:
xiaoxia
2026-10-05 19:52:14 +08:00
parent aac785a9ed
commit 257921abf6
5 changed files with 134 additions and 428 deletions
+6 -24
View File
@@ -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)}
@@ -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)
@@ -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:
@@ -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}[^>]*>(.*?)</{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"<product[^>]*>(.*?)</product>", 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)
@@ -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