Compare commits

..

2 Commits

Author SHA1 Message Date
CI Bot b4e84eb7b6 style: auto-format with black + isort + ruff + prettier [skip ci-format-check] 2026-10-05 09:12:45 +00:00
xiaoxia ed20fbad49 feat(vision): V2 快速图片分析路径 — OCR+lite JSON VLM并行,单图<3s目标
CI/CD Pipeline / Check push changed paths (pull_request) Has been skipped
CI/CD Pipeline / Dedup Check - skip PR tests when covered by push pipeline (pull_request) Successful in 1s
CI/CD Pipeline / Check if frontend-only change (pull_request) Successful in 1s
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 Web Image (pull_request) Has been skipped
CI/CD Pipeline / Build Staging API Image (pull_request) Has been skipped
CI/CD Pipeline / Build Staging Web Image (pull_request) Has been skipped
CI/CD Pipeline / Build Staging Worker Image (pull_request) Has been skipped
CI/CD Pipeline / Retag skipped Staging API Image (pull_request) Has been skipped
CI/CD Pipeline / Retag skipped Staging Web Image (pull_request) Has been skipped
CI/CD Pipeline / Retag skipped Staging Worker Image (pull_request) Has been skipped
CI/CD Pipeline / Deploy Staging (Watchtower auto-deploy) (pull_request) Has been skipped
CI/CD Pipeline / 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 / PR Build API Image (pull_request) Successful in 58s
CI/CD Pipeline / PR Build Worker Image (pull_request) Successful in 1m19s
Preview Deploy / Deploy Preview Environment (pull_request) Successful in 1m41s
ACR Cleanup / ACR Image Cleanup (pull_request_target) Waiting to run
Preview Cleanup / Cleanup Preview Environment (pull_request) Successful in 43s
PR Automation / Auto Approve on CI Green (pull_request) Successful in 2m59s
CI/CD Pipeline / Unit Tests (pull_request) Failing after 4m43s
CI/CD Pipeline / Validate - Python (mypy + alembic) (pull_request) Successful in 5m49s
CI/CD Pipeline / Integration Tests (pull_request) Successful in 5m50s
CI/CD Pipeline / Validate - Style (pull_request) Failing after 6m5s
AI Code Review / AI Code Review (pull_request) Successful in 6m59s
PR Automation / Auto Merge on CI Green + Approved (pull_request) Successful in 1m3s
CI/CD Pipeline / Validate - Security (pull_request) Successful in 12m21s
CI/CD Pipeline / Build Production API Image (pull_request) Has been skipped
CI/CD Pipeline / Build Production Web Image (pull_request) Has been skipped
CI/CD Pipeline / Build Production Worker Image (pull_request) Has been skipped
CI/CD Pipeline / CI Gate (pull_request) Failing after 6s
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
灵应直接指令(10-05):图片分析专用API组合方案直接干,不要等方案确认。

API现实说明:
- 火山引擎云端真实可用的视觉专用HTTP API:OCR(MediaKit tools-sync/ocr,Bearer鉴权)
- 人体属性/商品检测/图像标签:火山云端无公开HTTP API,仅有移动端SDK(智能美化特效,年费6-60万)
- 务实方案:OCR专用API + doubao-seed-2.1-lite强约束JSON-only prompt(替代3类缺失的专用API),
  pro VLM保留为终极兜底

架构:
- 新模块 apps/worker/worker_app/tasks/vision/:
  - vlm_fast_json.py:lite VLM极简JSON schema prompt,max_tokens=350,temp=0.1,timeout=8s
  - ocr_volc.py:MediaKit同步OCR封装,返回文本列表
  - assembler.py:字段映射+portrait_prompt模板拼接,输出格式与旧_normalize完全一致
  - fast_path.py:analyze_image_v2/analyze_images_v2,单图2路并行(OCR+lite JSON),
    外层8图全并发,置信度低/失败自动降级旧lite/pro竞速VLM
- _step_image_analysis增加VISION_V2_ENABLED环境变量开关:
  - true→走V2快速路径
  - false(默认)→走V1 #2198/#2199 lite/pro竞速路径(过渡期兜底)
- 下游信任链/t2i零改动:输出dict字段(name/brand/category/appearance/key_features/
  scene/mood/portrait_prompt/summary/_source)与旧格式完全兼容

性能目标:
- 单图fast路径目标<3s(OCR+lite JSON并行取最慢)
- 8图全并发<15s(较当前V1的~115s/3图提升10倍+)
- 兜底路径仍复用现有V1竞速,最坏情况不劣化
2026-10-05 17:06:05 +08:00
9 changed files with 862 additions and 752 deletions
@@ -1,182 +0,0 @@
/* DurationWheelPicker —— 弹层式滚轮选择器(样式与表单一致) */
/* 触发按钮:外观复用 .vv-select 风格 */
.dw-trigger {
display: flex;
align-items: center;
justify-content: space-between;
width: 100%;
height: 36px;
padding: 0 12px;
background: #fff;
border: 1px solid #e0e0e8;
border-radius: 8px;
font-size: 13px;
color: #1f2937;
cursor: pointer;
box-sizing: border-box;
transition: all 0.15s;
user-select: none;
}
.dw-trigger:hover {
border-color: #c0c0d0;
}
.dw-trigger-open,
.dw-trigger:focus-within {
border-color: #7c3aed !important;
box-shadow: 0 0 0 2px rgba(124, 58, 237, 0.12);
}
.dw-trigger-disabled {
opacity: 0.5;
pointer-events: none;
cursor: not-allowed;
}
.dw-trigger-val {
flex: 1;
overflow: hidden;
text-overflow: ellipsis;
white-space: nowrap;
}
.dw-trigger-placeholder {
color: #9ca3af;
}
.dw-trigger-arrow {
font-size: 10px;
color: #9ca3af;
margin-left: 8px;
transition: transform 0.2s;
}
.dw-trigger-arrow-up {
transform: rotate(180deg);
}
/* 弹层容器 */
.dw-popup {
padding: 8px;
min-width: 140px;
}
/* 滚轮 */
.dw-picker {
position: relative;
width: 100%;
overflow: hidden;
border-radius: 8px;
background: #fafafe;
border: 1px solid #e5e7eb;
}
.dw-picker-list {
margin: 0;
padding: 0;
list-style: none;
height: 100%;
overflow-y: scroll;
scroll-snap-type: y mandatory;
-webkit-overflow-scrolling: touch;
scrollbar-width: none;
}
.dw-picker-list::-webkit-scrollbar {
display: none;
}
.dw-picker-item {
display: flex;
align-items: baseline;
justify-content: center;
gap: 3px;
scroll-snap-align: center;
cursor: pointer;
font-size: 15px;
color: #9ca3af;
font-weight: 400;
transition:
color 0.15s,
transform 0.15s,
font-weight 0.15s;
}
.dw-picker-item-val {
font-variant-numeric: tabular-nums;
}
.dw-picker-item-unit {
font-size: 13px;
color: inherit;
}
.dw-picker-item-active {
color: #7c3aed;
font-weight: 600;
}
.dw-picker-item-active .dw-picker-item-val {
font-size: 18px;
}
.dw-picker-item-active .dw-picker-item-unit {
font-size: 14px;
}
/* 中心选中条 */
.dw-picker-mask {
position: absolute;
left: 6px;
right: 6px;
pointer-events: none;
background: #f5f0ff;
border-radius: 6px;
z-index: 1;
}
.dw-picker-mask::before,
.dw-picker-mask::after {
content: "";
position: absolute;
left: 0;
right: 0;
height: 1px;
background: #d8c4ff;
}
.dw-picker-mask::before {
top: 0;
}
.dw-picker-mask::after {
bottom: 0;
}
/* 上下渐变 */
.dw-picker-fade {
position: absolute;
left: 0;
right: 0;
height: 40%;
pointer-events: none;
z-index: 2;
}
.dw-picker-fade-top {
top: 0;
background: linear-gradient(to bottom, #fafafe 25%, rgba(250, 250, 254, 0));
}
.dw-picker-fade-bottom {
bottom: 0;
background: linear-gradient(to top, #fafafe 25%, rgba(250, 250, 254, 0));
}
/* 弹层按钮区 */
.dw-popup-actions {
display: flex;
gap: 8px;
justify-content: flex-end;
margin-top: 8px;
}
.dw-popup-actions .ant-btn {
border-radius: 6px;
}
.dw-popup-actions .ant-btn-primary {
background: #7c3aed;
}
.dw-popup-actions .ant-btn-primary:hover {
background: #6d28d9 !important;
}
/* 覆盖 antd Popover 默认内边距 */
.dw-popover .ant-popover-inner {
padding: 0 !important;
overflow: hidden;
}
.dw-popover .ant-popover-arrow {
display: none;
}
@@ -1,180 +0,0 @@
/**
* DurationWheelPicker —— 竖屏滚轮式时长选择器(弹层版)
*
* 设计:
* - 外观是和其他表单 Select 一致的输入框(白色底+1px灰边+紫色focus ring)
* - 点击输入框弹出 Popover,内部是滚轮 picker(原生 scroll-snap,零依赖)
* - 滚轮样式:白底容器,选中行 #7c3aed 紫字加粗+浅紫背景条
* - 支持触摸/鼠标滚轮/点击;松手吸附;底部"确认/取消"按钮
* - 默认范围 15–30 秒,步长 1 秒
*/
import React, { useEffect, useMemo, useRef, useState, useCallback } from "react"
import { Popover, Button } from "antd"
import { DownOutlined } from "@ant-design/icons"
import "./DurationWheelPicker.css"
export interface DurationWheelPickerProps {
value?: number
min?: number
max?: number
step?: number
unit?: string
onChange?: (value: number) => void
placeholder?: string
disabled?: boolean
/** 弹层宽度,默认 160px */
popupWidth?: number
/** 弹层内滚轮高度,默认 180px */
wheelHeight?: number
}
const ITEM_HEIGHT = 36
const DurationWheelPicker: React.FC<DurationWheelPickerProps> = ({
value = 20,
min = 15,
max = 30,
step = 1,
unit = "秒",
onChange,
placeholder = "请选择时长",
disabled = false,
popupWidth = 160,
wheelHeight = 180,
}) => {
const options = useMemo(() => {
const arr: number[] = []
for (let v = min; v <= max; v += step) arr.push(v)
return arr
}, [min, max, step])
const [open, setOpen] = useState(false)
// 弹层内暂存值,点确认才提交
const [draft, setDraft] = useState<number>(value)
const listRef = useRef<HTMLUListElement>(null)
const scrollTimerRef = useRef<ReturnType<typeof setTimeout> | null>(null)
useEffect(() => {
if (open) {
setDraft(value)
// 下一帧滚到当前值
requestAnimationFrame(() => scrollToValue(value, false))
}
// eslint-disable-next-line react-hooks/exhaustive-deps
}, [open])
const scrollToValue = useCallback(
(v: number, smooth = true) => {
const list = listRef.current
if (!list) return
const idx = options.indexOf(v)
if (idx < 0) return
list.scrollTo({ top: idx * ITEM_HEIGHT, behavior: smooth ? "smooth" : "auto" })
},
[options],
)
const handleScroll = () => {
if (scrollTimerRef.current) clearTimeout(scrollTimerRef.current)
scrollTimerRef.current = setTimeout(() => {
const list = listRef.current
if (!list) return
const idx = Math.round(list.scrollTop / ITEM_HEIGHT)
const clamped = Math.max(0, Math.min(options.length - 1, idx))
const targetTop = clamped * ITEM_HEIGHT
if (Math.abs(list.scrollTop - targetTop) > 1) {
list.scrollTo({ top: targetTop, behavior: "smooth" })
}
setDraft(options[clamped])
}, 100)
}
const handleConfirm = () => {
onChange?.(draft)
setOpen(false)
}
const handleCancel = () => {
setOpen(false)
}
const handleItemClick = (v: number) => {
setDraft(v)
scrollToValue(v, true)
}
const maskTop = wheelHeight / 2 - ITEM_HEIGHT / 2
const wheel = (
<div className="dw-popup">
<div
className="dw-picker"
style={{ height: wheelHeight, width: popupWidth - 24 /* padding */ }}
>
<div className="dw-picker-mask" style={{ top: maskTop, height: ITEM_HEIGHT }} aria-hidden />
<div className="dw-picker-fade dw-picker-fade-top" aria-hidden />
<div className="dw-picker-fade dw-picker-fade-bottom" aria-hidden />
<ul
ref={listRef}
className="dw-picker-list"
onScroll={handleScroll}
style={{
paddingTop: wheelHeight / 2 - ITEM_HEIGHT / 2,
paddingBottom: wheelHeight / 2 - ITEM_HEIGHT / 2,
}}
>
{options.map((v) => {
const isActive = v === draft
return (
<li
key={v}
className={`dw-picker-item${isActive ? " dw-picker-item-active" : ""}`}
style={{ height: ITEM_HEIGHT, lineHeight: `${ITEM_HEIGHT}px` }}
onClick={() => handleItemClick(v)}
aria-selected={isActive}
role="option"
>
<span className="dw-picker-item-val">{v}</span>
<span className="dw-picker-item-unit">{unit}</span>
</li>
)
})}
</ul>
</div>
<div className="dw-popup-actions">
<Button size="small" onClick={handleCancel}>
取消
</Button>
<Button size="small" type="primary" onClick={handleConfirm}>
确认
</Button>
</div>
</div>
)
return (
<Popover
open={!disabled && open}
onOpenChange={(v) => setOpen(v)}
content={wheel}
trigger="click"
placement="bottomLeft"
overlayClassName="dw-popover"
overlayStyle={{ padding: 0 }}
overlayInnerStyle={{ padding: 0, borderRadius: 10 }}
destroyTooltipOnHide
>
<div
className={`dw-trigger${disabled ? " dw-trigger-disabled" : ""}${open ? " dw-trigger-open" : ""}`}
style={{ height: 36 }}
>
<span className={`dw-trigger-val${value != null ? "" : " dw-trigger-placeholder"}`}>
{value != null ? `${value}${unit}` : placeholder}
</span>
<DownOutlined className={`dw-trigger-arrow${open ? " dw-trigger-arrow-up" : ""}`} />
</div>
</Popover>
)
}
export default DurationWheelPicker
@@ -218,12 +218,12 @@ const PURPOSES = [
"悬念短剧",
"情绪短片",
]
const DURATIONS = [15, 20, 30, 45, 60]
const RATIOS = [
{ v: "9:16", label: "9:16 竖屏(抖音/视频号)" },
{ v: "16:9", label: "16:9 横屏(B站/YouTube)" },
{ v: "1:1", label: "1:1 方形(小红书)" },
]
const DURATIONS = Array.from({ length: 16 }, (_, i) => 15 + i)
/** 兜底模型列表(接口未返回时使用,字段与 ViralVideoModel 对齐;后端返回后自动覆盖) */
const FALLBACK_VIDEO_MODELS: ViralVideoModel[] = [
{
@@ -437,7 +437,7 @@ const emptyTask = (id: string, title: string): TabTask => ({
language: "中文(普通话)",
viralStructure: STRUCTURES[0],
marketingPurpose: "",
duration: 20,
duration: 15,
persona: "",
videoRatio: "9:16",
videoModel: "seedance-2.5",
@@ -2431,7 +2431,7 @@ const ViralVideoPage: React.FC = () => {
options={PURPOSES.map((i) => ({ value: i, label: i }))}
/>
</div>
<div className="vv-form-row">
<div className="vv-form-row" style={{ gridColumn: "1 / -1" }}>
<label className="vv-label">文案视频时长</label>
<Select
className="vv-select"
+611 -40
View File
@@ -1,9 +1,7 @@
"""爆款视频 Celery 编排器 — ViralVideoOrchestrator.
"""爆款视频 Celery 编排器 — ViralVideoOrchestrator (v1.6 单次 Seedance 出片版).
V2 图片分析(10-05):火山OCR专用API + doubao-lite强约束JSON并行,单图<3s,8图<15s;pro VLM单次兜底。输出字段兼容旧格式,下游信任链/t2i零改动。
流水线步骤:
1. _step_image_analysis 图片分析(V2: OCR+qwen3.8-flash并行 + qwen3.7-plus兜底)
v1.6 重大简化(Seedance 2.5 单次最长 30 秒,直接出片):
1. _step_image_analysis 图片 VLM 分析(保留)
1.5 _step_video_analysis 参考视频风格分析(可选)
2. _step_intent_parsing 用户文案意图解析
3. _step_script_generation 编导分镜脚本生成(融合原 copy_fusion+storyboard+review,输出 copy_result 结构 + voiceover_script)
@@ -24,9 +22,11 @@ from __future__ import annotations
import json
import logging
import os
import re
import tempfile
import threading
import time
from concurrent.futures import ThreadPoolExecutor, as_completed
from pathlib import Path
from typing import Any
@@ -330,62 +330,633 @@ def _vision_fallback(idx: int, reason: str, extra: dict | None = None) -> dict:
return d
def _normalize_image_url(raw: str, idx: int) -> str:
"""将 job.images 中的 storage_key/相对路径/空值统一归一化为可公网访问 URL。
- http(s):// → 直接用
- 其他 → storage_key,通过 SharedStorageService.get_url() 转公网 URL
- 空值/非字符串 → 抛 ValueError
def _is_vision_result_usable(result: dict) -> bool:
"""判断 VLM 返回是否有效。
#2198b: 判定条件放宽——有人像(portrait_prompt非空/非'无人像')即视为usable(爆款视频核心
是要人物描述给信任链t2i用,product name/brand/features识别不准是次要的)。
非人像场景才要求name+summary+features有效。
"""
if not isinstance(result, dict):
return False
# 有人像描述(爆款视频最核心需求,portrait_prompt给信任链t2i做参考)就视为usable
pp = (result.get("portrait_prompt") or "").strip()
if pp and pp not in ("无人像", "无法判断", "未识别"):
return True
# 非人像场景:要求name+summary有效
name = (result.get("name") or "").strip()
if not name or name in ("未识别", "无法判断", "未知"):
return False
summary = (result.get("summary") or "").strip()
if len(summary) < 5 or summary in ("无法判断", "未识别"):
return False
category = (result.get("category") or "").strip()
if category == "非产品图":
return True
feats = result.get("key_features") or []
if not isinstance(feats, list) or len(feats) == 0 or feats == ["无法判断"]:
return False
return True
def _normalize_image_url(raw: str, idx: int) -> str:
"""#2188: 将 job.images 中的 storage_key/相对路径/空值统一归一化为可公网访问 URL。
- 以 http:// 或 https:// 开头 → 视为公网 URL
- 其他 → 视为 storage_key,用 SharedStorageService.get_url() 转公网 URL
- 空值/None/非字符串 → 抛 ValueError(上层 catch 后走 400 错误)
返回前做 HTTP 可达性检查(GET+Range:0-1024 避免 OSS 签名 URL 对 HEAD 返回 403 的假阴性)。
"""
import requests as _req
if not raw or not isinstance(raw, str):
raise ValueError(f"图片 #{idx} URL 为空或类型错误: {type(raw).__name__}={raw!r}")
url = raw.strip()
if not url:
raise ValueError(f"图片 #{idx} URL 为空白字符串")
if url.startswith("http://") or url.startswith("https://"):
return url
storage_key = url.lstrip("/")
try:
from packages.shared.storage import get_storage_service
# storage_key 判定:不以 http 开头
if not url.startswith("http://") and not url.startswith("https://"):
# 去掉可能的前导斜杠
storage_key = url.lstrip("/")
try:
from packages.shared.storage import get_storage_service
url = get_storage_service().get_url(storage_key)
_svc = get_storage_service()
url = _svc.get_url(storage_key)
except Exception as _e:
raise ValueError(f"图片 #{idx} storage_key={storage_key!r} 转公网URL失败: {_e}") from _e
logger.info("[爆款视频] 图片 #%d storage_key 已转公网 URL: %s", idx, url[:120])
# #2194: 用 GET+Range 代替 HEAD。
# Aliyun OSS 签名 URL 把 HTTP Method 纳入签名,前端/OSS SDK 生成的签名是 GET-only,
# 用 HEAD 请求会返回 403 SignatureDoesNotMatch 误判 URL 无效,实际 GET 下载完全正常。
# Range: bytes=0-1024 只取前1KB,开销极小。
try:
_r = _req.get(url, timeout=5, allow_redirects=True, stream=True, headers={"Range": "bytes=0-1024"})
if _r.status_code >= 400:
logger.warning("[爆款视频] 图片 #%d URL 可达性检查返回 %d: %s", idx, _r.status_code, url[:120])
_r.close()
except Exception as _e:
raise ValueError(f"图片 #{idx} storage_key={storage_key!r} 转公网URL失败: {_e}") from _e
logger.info("[爆款视频] 图片 #%d storage_key → 公网URL: %s", idx, url[:120])
logger.warning("[爆款视频] 图片 #%d URL 可达性检查异常: %s url=%s", idx, _e, url[:120])
return url
def _step_image_analysis(job: ViralVideoJob) -> dict:
"""步骤 1: 图片分析(V2 主路径)。
def _analyze_single_image(
idx: int,
img_url: str,
vision_model: str,
timeout: int,
*,
pro_fallback_model: str | None = None,
) -> dict:
"""单张图片 VLM 分析(#2040:改为从 prompt_loader 读模板 + XML 解析)。
架构:
- 主力:火山 MediaKit OCR(专用API,未配置时自动跳过)+ qwen3.8-flash 强约束 JSON,每图2路并行,目标<3s;
- 外层全并发(workers=8),目标8图<15s;
- 兜底: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 零改动。
lite 失败/不可用时用 pro 降级重试 1 次。失败/None 最终返回含默认字段的 dict。
"""
try:
from packages.application.viral_video import xml_parser as xp
from packages.application.viral_video.prompt_loader import (
get_template,
render_system_prompt,
render_user_prompt,
)
from packages.shared.ai_service import call_vision
except ImportError as e:
logger.warning("[爆款视频] prompt 模板/解析模块不可用: %s", e)
return _vision_fallback(idx, f"fallback_import_error:{e}")
if not img_url or not isinstance(img_url, str):
return _vision_fallback(idx, "invalid_url")
template = get_template("image_analysis")
system = render_system_prompt(template)
user = render_user_prompt(
template,
image_count=1,
industry="通用",
image_urls=f"第1张:{img_url}",
)
def _call(model: str, tmo: int, label: str):
# #2194/#2198: max_retries=0 由外层 _step_image_analysis 统一设置(阶段前置0、阶段后恢复),
# 子线程只读不改,避免嵌套并行竞速时多线程同时改 client.max_retries 产生竞态
_t0 = time.time()
try:
import json as _json
from packages.shared.ai_client import get_doubao_client as _gdc
_client = _gdc()
_messages = [
{"role": "system", "content": system},
{"role": "user", "content": user},
]
raw = _client.vision_completion(
messages=_messages,
images=[img_url],
temperature=0.3,
max_tokens=1200,
timeout=tmo,
model=model,
)
_elapsed = time.time() - _t0
logger.info(
"[爆款视频] 图片 #%d VLM(%s/%s) 完成 elapsed=%.1fs timeout=%d",
idx,
label,
model,
_elapsed,
tmo,
)
if raw is None:
return None
stripped = raw.strip()
if stripped.startswith("```"):
stripped = stripped.strip("`")
if stripped.startswith("json"):
stripped = stripped[4:].lstrip()
try:
return _json.loads(stripped)
except (_json.JSONDecodeError, TypeError):
return stripped
except Exception as e:
_elapsed = time.time() - _t0
logger.warning(
"[爆款视频] 图片 #%d call_vision(%s/%s) 异常 elapsed=%.1fs err=%s",
idx,
label,
model,
_elapsed,
e,
)
return None
def _xml_to_product(nodes: list, raw_text: str) -> dict:
product_nodes = [n for n in nodes if n["tag"] == "product"]
scene = xp.text_of(raw_text, "scene") or "通用"
mood = xp.text_of(raw_text, "mood") or ""
# #2184: #2177 XML 重构后人物信息放在顶层 <people has_person count gender age_range pose expression/>,
# 不再是 <product> 的 portrait_prompt 属性。需从顶层 people 标签提取并拼装 portrait_prompt。
portrait_prompt = "无人像"
try:
people_node = xp.find_first(raw_text, "people")
if people_node:
_pa = people_node.get("attrs") or {}
_has_person = xp.attr_bool(_pa.get("has_person"), False)
if _has_person:
_gender = _pa.get("gender", "无法判断") or "无法判断"
_age = _pa.get("age_range", "无法判断") or "无法判断"
_hair = _pa.get("hair", "无法判断") or "无法判断"
_skin = _pa.get("skin_tone", "无法判断") or "无法判断"
_face = _pa.get("face_shape", "无法判断") or "无法判断"
_outfit = _pa.get("outfit", "无法判断") or "无法判断"
_pose = _pa.get("pose", "无法判断") or "无法判断"
_expr = _pa.get("expression", "无法判断") or "无法判断"
_count = xp.attr_int(_pa.get("count"), 1)
# #2185: VLM有时对外貌属性输出"无法判断",用通用兜底值确保portrait_prompt始终有完整外貌描述
if _hair == "无法判断":
_hair = "自然发型"
if _skin == "无法判断":
_skin = "自然"
if _face == "无法判断":
_face = "标准"
if _outfit == "无法判断":
_outfit = "日常服装"
_parts = []
if _gender != "无法判断":
_g = _gender + ("性" if not _gender.endswith("性") else "")
_parts.append(_g)
else:
_parts.append("成年人")
if _age != "无法判断":
_parts.append(_age)
_parts.append("人物")
_parts.append(_hair)
_parts.append(f"{_skin}肤色")
_parts.append(f"{_face}脸型")
_parts.append(f"身着{_outfit}")
if _pose != "无法判断":
_parts.append(f"姿态{_pose}")
if _expr != "无法判断":
_parts.append(f"表情{_expr}")
else:
_parts.append("表情自然")
# #2186: 智能回填——VLM有时省略hair/outfit等外貌属性,但product.name/features/colors里已有相关信息
# 从product名字和features中提取服装关键词回填outfit
if _outfit in ("日常服装", "无法判断"):
for _ppn in product_nodes:
_pn = (_ppn.get("attrs") or {}).get("name", "") or ""
_pf = (_ppn.get("attrs") or {}).get("features", "") or ""
_ptxt = _pn + " " + _pf
# 服装关键词识别(常见上装/下装/裙装/套装)
_cloth_kws = [
# 衬衫/T恤类
"衬衫",
"T恤",
"POLO衫",
"polo衫",
"Polo衫",
"打底衫",
"雪纺衫",
"罩衫",
"针织衫",
# 毛衣/卫衣/针织类
"毛衣",
"卫衣",
"帽衫",
"针织",
"毛衫",
"开衫",
# 外套/西装/夹克/风衣类
"外套",
"西装",
"西服",
"夹克",
"皮衣",
"皮夹克",
"风衣",
"大衣",
"羽绒服",
"棉服",
"棉服",
"马甲",
"背心",
"开衫外套",
# 裙装
"连衣裙",
"半身裙",
"短裙",
"长裙",
"百褶裙",
"A字裙",
"旗袍",
"汉服",
"JK裙",
# 裤装
"牛仔裤",
"休闲裤",
"西裤",
"运动裤",
"短裤",
"阔腿裤",
"打底裤",
# 制服/套装
"制服",
"套装",
"职业装",
"工装",
# 通用上装/下装词(兜底)
"上衣",
"短袖",
"长袖",
"无袖",
"半袖",
"吊带",
"背心",
"网纱",
"雪纺",
"真丝",
"纯棉",
"亚麻",
]
for _ckw in _cloth_kws:
if _ckw in _ptxt:
_ci = _ptxt.find(_ckw)
# 向前找颜色/材质/款式形容词(白/黑/米/红/蓝/灰/棉/麻/长/短/厚/薄/长袖/短袖/翻领/圆领/V领/印花/条纹等)
_start = max(0, _ci - 12)
# 向后包含款式词(长袖/短袖/外套/套装/上衣等后续修饰)
_end = min(len(_ptxt), _ci + len(_ckw) + 8)
_outfit_extract = _ptxt[_start:_end].strip(" ,,。.、")
# 仅清理明确的品牌/产品类前缀(不清理颜色/款式/尺寸形容词)
_outfit_extract = re.sub(
r"^(\S{0,4}牌|\S{0,3}品牌|\S{0,3}款|产品|商品|的)", "", _outfit_extract
).strip()
# 尾部清理:去掉残留的品牌字/型号字(如"标""ml""g""装"等单字杂字)
_outfit_extract = re.sub(
r"(标[0-9a-zA-Z]*|\d+\s*(?:ml|g|L|斤|件|个|瓶|盒|包|袋|装)|\s+\d+\s*)$",
"",
_outfit_extract,
flags=re.IGNORECASE,
).strip()
if len(_outfit_extract) >= 2:
_outfit = _outfit_extract
break
if _outfit not in ("日常服装", "无法判断"):
break
# 从color标签中提取头发颜色回填hair
if _hair in ("自然发型", "无法判断"):
_hair_color = ""
_color_nodes = [n for n in nodes if n["tag"] == "color"]
_hair_kws_map = {
"黑": "黑色",
"棕": "棕色",
"金": "金色",
"栗": "栗色",
"红": "红色",
"白": "白色",
"灰": "灰色",
"蓝": "蓝色",
"黄": "黄色",
"紫": "紫色",
}
for _cn in _color_nodes:
_cname = (_cn.get("attrs") or {}).get("name", "") or ""
# 小占比颜色更可能是发色(非主色的小面积色),且名称含头发/黑/棕/金等
_ccov = 0.0
try:
_ccov = float((_cn.get("attrs") or {}).get("coverage", "0") or 0)
except Exception:
pass
for _hk, _hv in _hair_kws_map.items():
if _hk in _cname and _ccov < 0.3:
_hair_color = _hv
break
if _hair_color:
break
if _hair_color:
_hair = f"{_hair_color}头发"
else:
_hair = "自然发型"
# 重新拼装_parts(回填后)
_parts = []
if _gender != "无法判断":
_g = _gender + ("性" if not _gender.endswith("性") else "")
_parts.append(_g)
else:
_parts.append("成年人")
if _age != "无法判断":
_parts.append(_age)
_parts.append("人物")
_parts.append(_hair)
_parts.append(f"{_skin}肤色")
_parts.append(f"{_face}脸型")
_parts.append(f"身着{_outfit}")
if _pose != "无法判断":
_parts.append(f"姿态{_pose}")
if _expr != "无法判断":
_parts.append(f"表情{_expr}")
else:
_parts.append("表情自然")
portrait_prompt = ",".join(_parts)
logger.info(
"[爆款视频] 图片 #%d 解析<people>(回填后): count=%d gender=%s age=%s hair=%s skin=%s face=%s outfit=%s pose=%s expr=%s → %s",
idx,
_count,
_gender,
_age,
_hair,
_skin,
_face,
_outfit,
_pose,
_expr,
portrait_prompt,
)
except Exception as _pe:
logger.warning("[爆款视频] 图片 #%d 解析<people>标签异常: %s,回退无人像", idx, _pe)
for p in product_nodes:
a = p["attrs"]
text_on_pkg = a.get("text_on_package", "")
p_body = p.get("text", "") or ""
if not text_on_pkg and p_body:
text_on_pkg = xp.text_of(p_body, "text_on_package") or ""
text_list = [x.strip() for x in re.split(r"[,,;;]", text_on_pkg) if x.strip()] if text_on_pkg else []
features = a.get("features", "")
feat_list = [x.strip() for x in re.split(r"[,,;;]", features) if x.strip()] if features else []
name = a.get("name", "") or "未识别"
brand = a.get("brand", "") or "无法判断"
category = a.get("category", "") or "无法判断"
appearance = a.get("appearance", "") or "无法判断"
packaging = a.get("packaging", "") or "无法判断"
summary = a.get("summary", "") or f"{brand} {name}"
# 优先取 product 属性上的 portrait_prompt(兼容旧schema),否则用顶层 <people> 解析结果
_pp_from_attr = a.get("portrait_prompt", "")
if _pp_from_attr and _pp_from_attr != "无人像":
portrait_prompt = _pp_from_attr
return {
"name": name,
"brand": brand,
"category": category,
"appearance": appearance,
"packaging": packaging,
"text_on_package": text_list,
"key_features": feat_list or [features] if features else ["无法判断"],
"scene": scene,
"mood": mood,
"portrait_prompt": portrait_prompt,
"summary": summary,
"_source": "xml",
}
# 没有 product 标签但有 <people has_person="true"> 也要能取到人物描述(兜底)
if portrait_prompt != "无人像":
return {
"name": "未识别",
"brand": "无法判断",
"category": "无法判断",
"appearance": "无法判断",
"packaging": "无法判断",
"text_on_package": [],
"key_features": ["无法判断"],
"scene": scene,
"mood": mood,
"portrait_prompt": portrait_prompt,
"summary": "未识别",
"_source": "xml_no_product",
}
return _vision_fallback(idx, "no_product_tag")
def _normalize(raw, source: str) -> dict:
if raw is None:
return _vision_fallback(idx, f"{source}_none")
# VLM 偶尔直接返回 JSON 对象(不包裹```json),_call 里 json.loads 后已是 dict
if isinstance(raw, dict):
_prod = {
"name": raw.get("name") or "未识别",
"brand": raw.get("brand") or "无法判断",
"category": raw.get("category") or "无法判断",
"appearance": raw.get("appearance") or "无法判断",
"packaging": raw.get("packaging") or "无法判断",
"text_on_package": raw.get("text_on_package") or [],
"key_features": raw.get("key_features") or raw.get("features") or ["无法判断"],
"scene": raw.get("scene") or "通用",
"mood": raw.get("mood") or "",
"portrait_prompt": raw.get("portrait_prompt") or "无人像",
"summary": raw.get("summary") or f"{raw.get('brand','')} {raw.get('name','')}",
"_source": source,
}
return _prod
if not isinstance(raw, str):
return _vision_fallback(idx, f"{source}_badtype")
nodes = xp.parse_tags(raw)
if not nodes:
# 不是 XML 也不是 dict:尝试当作纯 JSON 字符串再解析一次
try:
import json as _j2
_jd = _j2.loads(raw)
if isinstance(_jd, dict):
return _normalize(_jd, source)
except Exception:
pass
logger.warning("[爆款视频] 图片 #%d XML/JSON 解析都失败 source=%s raw_head=%s", idx, source, raw[:200])
return _vision_fallback(idx, f"{source}_parse_fail", {"_raw": raw[:500]})
product = _xml_to_product(nodes, raw)
product.setdefault("_source", source)
product["raw"] = raw[:500]
return product
# #2198: lite/pro 并行竞速。同时发两个请求,先返回 usable 结果就用哪个,避免
# 串行 lite超时→再发pro 累计80-100s的惩罚。外层 max_workers=2 图片并发时,竞速模式下
# VLM 总并发=4(2图 × 2模型),实测 Ark 可以承受,且因为取快者而不是等两个都完,
# 单图通常 40-50s 就能拿到 pro 结果(pro 正常 42-46s),lite 偶发 30s 内返回时更快。
race_t0 = time.time()
lite_tag = vision_model.split("/")[-1] if "/" in vision_model else vision_model
winner: dict | None = None
with ThreadPoolExecutor(max_workers=2) as _inner_pool:
f_lite = _inner_pool.submit(_call, vision_model, timeout, "lite")
# pro 给 75s(原60s太紧实测1/3超时,pro正常42-65s给10s余量)
pro_tmo = 75
f_pro = _inner_pool.submit(_call, pro_fallback_model or vision_model, pro_tmo, "pro")
_fmap = {f_lite: ("lite", lite_tag), f_pro: ("pro", "pro_fallback")}
for _fut in as_completed(_fmap, timeout=pro_tmo + 15):
_lbl, _tag = _fmap[_fut]
try:
_raw = _fut.result()
except Exception as _e:
logger.warning("[爆款视频] 图片 #%d %s future异常: %s", idx, _lbl, _e)
_raw = None
_res = _normalize(_raw, _tag)
if _is_vision_result_usable(_res):
winner = _res
if _lbl == "pro":
winner["_fallback_used"] = True
logger.info(
"[爆款视频] 图片 #%d 竞速胜出=%s elapsed=%.1fs",
idx,
_lbl,
time.time() - race_t0,
)
break
if winner is not None:
return winner
# 两个都失败,返回最后一次 _normalize 结果(通常是 pro 的失败 fallback,含 _source=pro_fallback_none)
try:
_last_raw = f_pro.result(timeout=1)
except Exception:
_last_raw = None
_last = _normalize(_last_raw, "pro_fallback")
logger.warning(
"[爆款视频] 图片 #%d lite/pro 竞速均失败 elapsed=%.1fs",
idx,
time.time() - race_t0,
)
return _last
def _step_image_analysis(job: ViralVideoJob) -> dict:
"""步骤 1: 图片分析。
V2(VISION_V2_ENABLED=true,10-05 新方案):
- 每图并行 2 路:火山 MediaKit OCR(专用API)+ doubao-seed-2.1-lite 强约束 JSON
(火山云端无人体属性/商品检测/图像标签公开 HTTP API,用 lite JSON-only VLM 弥补),
目标单图 <3s;
- 外层 8 图全并发,目标 8 图 <15s;
- 置信度低/全失败时降级 doubao-seed-2.1-pro 完整 VLM 兜底(复用旧竞速逻辑);
- 输出 dict 格式与旧 _normalize() 完全一致,下游信任链/t2i 零改动。
V1(默认,#2198/#2199 lite/pro 并行竞速):
- 过渡版兜底,单图 lite(30s)/pro(75s) 竞速,外层 max_workers=2,典型 40-75s/图。
"""
try:
from packages.shared.ai_service import call_vision # noqa: F401
except ImportError:
logger.warning("[爆款视频] ai_service.call_vision 不可用,使用占位结果")
return {"products": [_vision_fallback(0, "fallback_import_error")]}
if not job.images:
logger.warning("[爆款视频] 任务无 images,跳过图片分析")
return {"products": []}
# URL 归一化(storage_key→公网URL;空值直接400)
# URL 归一化 — storage_key→公网URL + 空值报400
normalized_urls: list[str] = []
for idx, raw in enumerate(job.images):
normalized_urls.append(_normalize_image_url(raw, idx))
try:
from worker_app.tasks.vision import analyze_images_v2 as _aiv2
except ImportError:
try:
from tasks.vision import analyze_images_v2 as _aiv2 # type: ignore
except ImportError as e:
logger.error("[爆款视频] vision 模块导入失败: %s", e)
return {"products": [_vision_fallback(0, f"vision_import_error:{e}")]}
normalized_urls.append(_normalize_image_url(raw, idx))
except ValueError as _ve:
logger.error("[爆款视频] 图片 #%d URL 归一化失败: %s", idx, _ve)
raise
# V2 内部 httpx 直连 dashscope,单次调用无重试,无需调整全局 client
results = _aiv2(normalized_urls)
return {"products": list(results)}
# 模型配置(V1/V2 共用)
try:
_s = get_shared_settings()
lite_model = _s.doubao_vision_lite_model
pro_model = _s.doubao_vision_model
except Exception:
lite_model = "doubao-seed-2-1-lite-260915"
pro_model = "doubao-seed-2-1-pro-260915"
# ========== V2 路径(VISION_V2_ENABLED=true)==========
import os as _os_v2
_v2_enabled = _os_v2.environ.get("VISION_V2_ENABLED", "false").lower() in ("1", "true", "yes", "on")
if _v2_enabled:
try:
from worker_app.tasks.vision import analyze_images_v2 as _aiv2
except ImportError:
try:
from tasks.vision import analyze_images_v2 as _aiv2 # type: ignore
except ImportError:
logger.warning("[vision.v2] 模块导入失败,回退 V1 路径")
_aiv2 = None # type: ignore
if _aiv2 is not None:
from packages.shared.ai_client import get_doubao_client as _gdc_v2
_cli = _gdc_v2()
_orig_retries_v2 = _cli.max_retries
_cli.max_retries = 0
try:
results_v2 = _aiv2(normalized_urls, lite_model=lite_model, pro_model=pro_model)
finally:
_cli.max_retries = _orig_retries_v2
return {"products": list(results_v2)}
# 导入失败 → fallthrough 走 V1
# ========== V1 路径(默认,lite/pro 并行竞速)==========
vision_model = lite_model
vision_timeout = 30
from packages.shared.ai_client import get_doubao_client as _gdc_step
_step_client = _gdc_step()
_step_orig_retries = _step_client.max_retries
_step_client.max_retries = 0
results: list[dict] = [None] * len(normalized_urls) # type: ignore
max_workers = min(2, max(1, len(normalized_urls))) # 并发≤2 防方舟限流(竞速模式下总并发=4)
logger.info(
"[爆款视频] 开始并行竞速图片分析(V1) n=%d lite=%s(%ds) pro=%s(75s) img_workers=%d",
len(normalized_urls),
vision_model,
vision_timeout,
pro_model,
max_workers,
)
try:
with ThreadPoolExecutor(max_workers=max_workers) as pool:
future_to_idx = {
pool.submit(
_analyze_single_image, idx, url, vision_model, vision_timeout, pro_fallback_model=pro_model
): idx
for idx, url in enumerate(normalized_urls)
}
for fut in as_completed(future_to_idx):
idx = future_to_idx[fut]
try:
results[idx] = fut.result()
except Exception as e:
logger.warning("[爆款视频] 图片 #%d future 异常 err=%s", idx, e, exc_info=True)
results[idx] = _vision_fallback(idx, "future_exception", {"_error": str(e)[:200]})
finally:
_step_client.max_retries = _step_orig_retries
return {"products": results}
def _step_video_analysis(job: ViralVideoJob) -> dict | None:
@@ -1,4 +1,16 @@
# -*- coding: utf-8 -*-
"""V2 图片分析:火山OCR专用API + doubao-lite强约束JSON并行,单次pro VLM兜底。"""
"""V2 图片分析:专用 API 组合路径(OCR + lite JSON VLM 并行 + pro VLM 兜底)。
灵应指令(10-05):
- 优先火山引擎视觉智能 API:OCR 是真实专用云端 API;
- 人体属性/商品检测/图像标签:火山云端无公开 HTTP API(仅有移动端 SDK),
采用 doubao-seed-2.1-lite + 强约束 JSON-only prompt 作为"伪专用 API",
目标 1-3s 返回结构化字段;
- VLM(doubao-seed-2.1-pro)保留为终极兜底(置信度低/全失败时降级);
- 单图并行 2 路(OCR + lite JSON VLM),外层 8 图全并发,目标 8 图 <15s。
输出 dict 格式与 viral_video._normalize() 完全一致,下游信任链/t2i 零改动。
灰度开关:VISION_V2_ENABLED=true(默认 false,走旧 #2198/#2199 竞速逻辑)。
"""
from .fast_path import analyze_image_v2, analyze_images_v2 # noqa: F401
@@ -21,6 +21,8 @@ def _join_parts(*parts: str | None) -> str:
_AGE_PREFIX = {
"儿童": "小女孩" if None else "儿童",
"青少年": "少女" if None else "少年",
"青年": "年轻",
"中年": "中年",
"老年": "老年",
@@ -153,8 +155,7 @@ def _build_portrait_prompt(fj: dict[str, Any]) -> str:
if detail_parts:
pieces.append(",".join(detail_parts))
if style_parts:
# 风格词之间不用逗号,用空格紧凑
pieces.append("".join(style_parts) + "风格")
pieces.append(",".join(style_parts) + "风格")
else:
pieces.append("人像写真")
+180 -81
View File
@@ -1,12 +1,9 @@
# -*- coding: utf-8 -*-
"""V2 图片分析主路径:每图并行 OCR(火山MediaKit,未配置时自动跳过)+ qwen3.8-flash JSON VLM,
失败时单次 qwen3.7-plus 兜底。
"""V2 快速路径:每图并行 OCR + lite JSON VLM,失败降级 pro VLM。
架构(灵应10-05确认):
- 唯一后端:阿里云百炼 DashScope,qwen3.8-flash 做快速路径、qwen3.7-plus 做兜底
- 主力:单图2路并行(OCR + fast VLM),外层N图全并发(workers=8)
- 兜底:单次 pro VLM 调用,无竞速/重试/复杂超时
- 输出 dict 格式与旧版完全一致,下游零改动
单图并行 2 路(OCR + lite JSON VLM),目标 <3s。
外层 8 图全并发,目标 8 图 <15s。
终极兜底:复用旧 _analyze_single_image 完整 pro VLM 逻辑。
"""
from __future__ import annotations
@@ -17,128 +14,230 @@ 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 assembler, ocr_volc, vlm_fast_json
logger = logging.getLogger(__name__)
# 超时(可通过环境变量覆盖)
_IMG_WORKERS = int(os.environ.get("VISION_V2_IMG_WORKERS", "8"))
_FAST_TIMEOUT = float(os.environ.get("VISION_V2_FAST_TIMEOUT", "12"))
_FAST_JSON_TIMEOUT = float(os.environ.get("VISION_V2_FAST_JSON_TIMEOUT", "12"))
_OCR_TIMEOUT = float(os.environ.get("VISION_V2_OCR_TIMEOUT", "6"))
_PRO_TIMEOUT = float(os.environ.get("VISION_V2_PRO_TIMEOUT", "25"))
_FALLBACK_RESULT = {
"name": "未识别",
"brand": "无法判断",
"category": "非产品图",
"appearance": "无法判断",
"packaging": "无法判断",
"text_on_package": [],
"key_features": ["无法判断"],
"scene": "通用",
"mood": "",
"portrait_prompt": "无法判断",
"summary": "未识别",
}
# ---------- 配置项(可通过环境变量覆盖) ----------
VISION_V2_ENABLED = os.environ.get("VISION_V2_ENABLED", "false").lower() in ("1", "true", "yes", "on")
# 单图 fast 路径总超时(包含 OCR + fast_json 并行)
V2_FAST_TIMEOUT = float(os.environ.get("VISION_V2_FAST_TIMEOUT", "10"))
# 外层图片并发(默认 8,即全并行)
V2_IMG_WORKERS = int(os.environ.get("VISION_V2_IMG_WORKERS", "8"))
# fast_json 单次超时
V2_FAST_JSON_TIMEOUT = float(os.environ.get("VISION_V2_FAST_JSON_TIMEOUT", "8"))
# OCR 单次超时
V2_OCR_TIMEOUT = float(os.environ.get("VISION_V2_OCR_TIMEOUT", "8"))
# pro VLM 兜底超时(仅在 fast 路径完全失败时触发)
V2_PRO_TIMEOUT = float(os.environ.get("VISION_V2_PRO_TIMEOUT", "45"))
V2_LITE_TIMEOUT = float(os.environ.get("VISION_V2_LITE_TIMEOUT", "20"))
def _is_usable(r: dict[str, Any]) -> bool:
pp = (r.get("portrait_prompt") or "").strip()
def _is_result_usable(result: dict[str, Any]) -> bool:
"""与 viral_video._is_vision_result_usable 对齐的可用判定。"""
pp = (result.get("portrait_prompt") or "").strip()
if pp and pp not in ("无人像", "无法判断", "未识别"):
return True
name = (r.get("name") or "").strip()
name = result.get("name") or ""
if name and name not in ("未识别", "无法判断", "未知"):
return True
summary = result.get("summary") or ""
if len(summary) >= 5 and summary not in ("无法判断", "未识别"):
return True
cat = result.get("category") or ""
if cat == "非产品图" and pp != "无人像":
return True
kf = result.get("key_features") or []
if kf and kf != ["无法判断"]:
# 只要有非默认特征且非空
return True
return False
def analyze_image_v2(idx: int, img_url: str) -> dict[str, Any]:
t0 = time.time()
def _call_pro_fallback(img_url: str, idx: int, lite_model: str, pro_model: str) -> dict[str, Any] | None:
"""fast 路径失败时,调用旧的 lite/pro 竞速 VLM。
fj_result: dict[str, Any] | None = None
ocr_result: list[str] = []
with ThreadPoolExecutor(max_workers=2) as pool:
f_fj = pool.submit(vlm_fast_json.call_fast_json, img_url, timeout=_FAST_JSON_TIMEOUT)
f_ocr = pool.submit(ocr_volc.call_ocr, img_url, timeout=_OCR_TIMEOUT)
复用 viral_video._analyze_single_image 的实现,避免重复代码。
"""
try:
from worker_app.tasks.viral_video import _analyze_single_image
except ImportError:
try:
for fut in as_completed([f_fj, f_ocr], timeout=_FAST_TIMEOUT):
try:
res = fut.result(timeout=1)
except Exception as e:
logger.warning("[vision.v2] 图片 #%d 子任务异常: %s", idx, e)
continue
if fut is f_fj and isinstance(res, dict):
fj_result = res
elif fut is f_ocr and isinstance(res, list):
ocr_result = res
except TimeoutError:
for f in (f_fj, f_ocr):
if not f.done():
f.cancel()
logger.warning("[vision.v2] 图片 #%d fast路径超时(%.0fs),走pro兜底", idx, _FAST_TIMEOUT)
from tasks.viral_video import _analyze_single_image # type: ignore
except ImportError:
logger.warning("[vision.v2] 无法 import _analyze_single_image,跳过 pro 兜底")
return None
try:
return _analyze_single_image(
idx,
img_url,
vision_model=lite_model,
timeout=int(V2_LITE_TIMEOUT),
pro_fallback_model=pro_model,
)
except Exception as e:
logger.warning("[vision.v2] 图片 #%d pro 兜底异常 err=%s", idx, e, exc_info=True)
return None
def analyze_image_v2(
idx: int,
img_url: str,
*,
lite_model: str | None = None,
pro_model: str | None = None,
) -> dict[str, Any]:
"""单张图片 V2 分析:OCR + lite JSON VLM 并行,必要时降级 pro VLM。
返回的 dict 与 viral_video._normalize() 输出格式完全一致。
"""
t0 = time.time()
# ---- 第 1 层:fast 路径并行 ----
fast_json_result: dict[str, Any] | None = None
ocr_result: list[str] = []
with ThreadPoolExecutor(max_workers=2) as pool:
f_fj = pool.submit(
vlm_fast_json.call_fast_json,
img_url,
model=lite_model,
timeout=V2_FAST_JSON_TIMEOUT,
)
f_ocr = pool.submit(ocr_volc.call_ocr, img_url, timeout=V2_OCR_TIMEOUT)
# 等全部完成或超时
for fut in as_completed([f_fj, f_ocr], timeout=V2_FAST_TIMEOUT):
try:
res = fut.result(timeout=1)
except Exception as e:
logger.warning("[vision.v2] 图片 #%d fast 子任务异常: %s", idx, e)
continue
if fut is f_fj:
fast_json_result = res if isinstance(res, dict) else None
elif fut is f_ocr:
ocr_result = res if isinstance(res, list) else []
fast_elapsed = time.time() - t0
if fj_result:
assembled = assembler.assemble_result(idx, fj_result, ocr_result)
if _is_usable(assembled):
# ---- 组装 fast 结果 ----
assembled: dict[str, Any] | None = None
if fast_json_result:
assembled = assembler.assemble_result(idx, fast_json_result, ocr_result)
if _is_result_usable(assembled):
assembled["_fast_elapsed"] = round(fast_elapsed, 2)
logger.info(
"[vision.v2] 图片 #%d fast命中 elapsed=%.2fs pp=%s",
"[vision.v2] 图片 #%d fast 路径命中 elapsed=%.2fs portrait_prompt=%s",
idx,
fast_elapsed,
(assembled.get("portrait_prompt") or "")[:40],
)
return assembled
logger.info(
"[vision.v2] 图片 #%d fast 结果不可用 portrait_prompt=%s,走 pro 兜底",
idx,
(assembled.get("portrait_prompt") or "")[:40],
)
else:
logger.info("[vision.v2] 图片 #%d fast_json 返回空 elapsed=%.2fs,走 pro 兜底", idx, fast_elapsed)
# ---- 第 2 层:pro VLM 兜底(复用旧竞速逻辑)----
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):
use_lite = lite_model or vlm_fast_json.DEFAULT_LITE_MODEL
use_pro = pro_model or "doubao-seed-2-1-pro-260915"
pro_result = _call_pro_fallback(img_url, idx, use_lite, use_pro)
if pro_result and _is_result_usable(pro_result):
pro_result["_fallback_used"] = True
pro_result["_fast_elapsed"] = round(fast_elapsed, 2)
pro_result["_pro_elapsed"] = round(time.time() - pro_t0, 2)
if ocr_result and not pro_result.get("text_on_package"):
pro_result["text_on_package"] = ocr_result[:8]
logger.info("[vision.v2] 图片 #%d pro兜底命中 total=%.2fs", idx, time.time() - t0)
logger.info(
"[vision.v2] 图片 #%d pro 兜底命中 total_elapsed=%.2fs",
idx,
time.time() - t0,
)
return pro_result
logger.warning("[vision.v2] 图片 #%d 全路径失败 elapsed=%.2fs", idx, time.time() - t0)
out = dict(_FALLBACK_RESULT)
out["_source"] = "v2_all_failed"
out["text_on_package"] = ocr_result[:8]
out["_fast_elapsed"] = round(fast_elapsed, 2)
return out
# ---- 第 3 层:兜底失败,返回 assembled 或标准 fallback ----
if assembled:
assembled["_source"] = "v2_fast_json_degraded"
logger.warning(
"[vision.v2] 图片 #%d pro 兜底也失败,返回降级 fast 结果 elapsed=%.2fs",
idx,
time.time() - t0,
)
return assembled
# 最后的最后:返回最小可用结构
logger.warning("[vision.v2] 图片 #%d 所有路径均失败 elapsed=%.2fs", idx, time.time() - t0)
return {
"name": "未识别",
"brand": "无法判断",
"category": "非产品图",
"appearance": "无法判断",
"packaging": "无法判断",
"text_on_package": ocr_result[:8],
"key_features": ["无法判断"],
"scene": "通用",
"mood": "",
"portrait_prompt": "无法判断",
"summary": "未识别",
"_source": "v2_all_failed",
}
def analyze_images_v2(img_urls: list[str]) -> list[dict[str, Any]]:
def analyze_images_v2(
img_urls: list[str],
*,
lite_model: str | None = None,
pro_model: str | None = None,
max_workers: int | None = None,
) -> list[dict[str, Any]]:
"""批量图片 V2 分析(外层全并行)。"""
if not img_urls:
return []
workers = min(_IMG_WORKERS, len(img_urls), 16)
workers = max_workers if max_workers and max_workers > 0 else V2_IMG_WORKERS
workers = min(workers, len(img_urls), 16) # 安全上限 16
results: list[dict[str, Any] | None] = [None] * len(img_urls)
logger.info(
"[vision.v2] 开始图片分析 n=%d workers=%d fast_timeout=%.0fs pro_timeout=%.0fs",
"[vision.v2] 开始 V2 并行图片分析 n=%d workers=%d fast_timeout=%.0fs",
len(img_urls),
workers,
_FAST_TIMEOUT,
_PRO_TIMEOUT,
V2_FAST_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)}
future_to_idx = {
pool.submit(analyze_image_v2, idx, url, lite_model=lite_model, pro_model=pro_model): idx
for idx, url in enumerate(img_urls)
}
for fut in as_completed(future_to_idx):
idx = future_to_idx[fut]
try:
results[idx] = fut.result()
except Exception as e:
logger.warning("[vision.v2] 图片 #%d future异常: %s", idx, e, exc_info=True)
r = dict(_FALLBACK_RESULT)
r["_source"] = "v2_future_exception"
results[idx] = r
logger.warning("[vision.v2] 图片 #%d future 异常 err=%s", idx, e, exc_info=True)
results[idx] = {
"name": "未识别",
"brand": "无法判断",
"category": "非产品图",
"appearance": "无法判断",
"packaging": "无法判断",
"text_on_package": [],
"key_features": ["无法判断"],
"scene": "通用",
"mood": "",
"portrait_prompt": "无法判断",
"summary": "未识别",
"_source": "v2_future_exception",
}
elapsed = time.time() - t0
succ = sum(1 for r in results if r and _is_usable(r))
succ = sum(1 for r in results if r and _is_result_usable(r))
fb = sum(1 for r in results if r and r.get("_fallback_used"))
logger.info("[vision.v2] 完成 n=%d usable=%d pro_fallback=%d elapsed=%.2fs", len(img_urls), succ, fb, elapsed)
logger.info(
"[vision.v2] V2 图片分析完成 n=%d success=%d pro_fallback=%d elapsed=%.2fs",
len(img_urls),
succ,
fb,
elapsed,
)
return [r for r in results if r is not None]
@@ -1,175 +0,0 @@
# -*- coding: utf-8 -*-
"""V2 pro 兜底:qwen3.7-plus(阿里云百炼/DashScope)单次调用。
fast_json 结果不可用时单次调用,无竞速、无重试、无复杂超时逻辑。
直接 httpx 发精简 JSON-only prompt(比旧版 prompt_loader XML 模板短很多,降低延迟)。
"""
from __future__ import annotations
import json
import logging
import os
import time
from typing import Any
logger = logging.getLogger(__name__)
_BASE_URL = "https://dashscope.aliyuncs.com/compatible-mode/v1"
_PRO_MODEL = "qwen3.7-plus"
_DEFAULT_TIMEOUT = 25
_DEFAULT_MAX_TOKENS = 800
_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": "产品名称,非产品图填null",\n'
' "brand": "品牌或文字标识,无则null",\n'
' "key_features": ["特征数组"]\n'
"}"
)
_PRO_USER = "分析这张图片,返回符合schema的JSON。"
def _strip_code_fence(s: str) -> str:
s = s.strip()
if s.startswith("```"):
lines = s.split("\n")
if lines and lines[0].startswith("```"):
lines = lines[1:]
if lines and lines[-1].strip().startswith("```"):
lines = lines[:-1]
s = "\n".join(lines).strip()
return s
def _assemble_pp(obj: dict[str, Any]) -> str:
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_pro_vlm(
img_url: str,
idx: int,
*,
timeout: int = _DEFAULT_TIMEOUT,
) -> dict[str, Any] | None:
"""单次调用 qwen3.7-plus,解析后返回 product dict;失败返回 None。"""
t0 = time.time()
import httpx
api_key = os.environ.get("DASHSCOPE_API_KEY")
if not api_key:
logger.warning("[vision.v2] DASHSCOPE_API_KEY 未配置,跳过 pro 兜底")
return None
url = f"{_BASE_URL}/chat/completions"
payload: dict[str, Any] = {
"model": _PRO_MODEL,
"messages": [
{"role": "system", "content": _PRO_SYSTEM},
{
"role": "user",
"content": [
{"type": "image_url", "image_url": {"url": img_url}},
{"type": "text", "text": _PRO_USER},
],
},
],
"temperature": 0.3,
"max_tokens": _DEFAULT_MAX_TOKENS,
"stream": False,
"enable_thinking": False,
"response_format": {"type": "json_object"},
}
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.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.v2] pro 返回空 elapsed=%.1fs", elapsed)
return None
usage = data.get("usage") or {}
logger.info(
"[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.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(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 has_person and not has_product 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",
}
except Exception as e:
logger.warning("[vision.v2] pro 异常 idx=%d elapsed=%.1fs err=%s", idx, time.time() - t0, e, exc_info=True)
return None
@@ -1,43 +1,35 @@
# -*- coding: utf-8 -*-
"""V2 快速路径:qwen3.8-flash(阿里云百炼/DashScope)强约束 JSON-only 调用。
"""doubao-seed-2.1-lite 强约束 JSON-only 调用。
目标:替代"人体属性/商品检测/图像标签"三个火山不存在的专用云端 API。
设计要点:
- 直接用 httpx 发最小 payload 到 DashScope OpenAI 兼容 endpoint,不走 ai_client 包装
- enable_thinking=false 关闭推理链(reasoning 是延迟主因)
- system prompt 极致精简,只给字段 schema 和强约束(禁止自然语言、禁止 markdown)
- max_tokens=350、temperature=0.1(稳定输出 JSON)
- timeout=12s(失败由外层走 pro 兜底)
- API Key 从环境变量 DASHSCOPE_API_KEY 读取
- max_tokens=350(比旧 VLM 的 1200 小很多,降低延迟)
- temperature=0.1(极低,稳定输出 JSON)
- timeout=8s(够快,失败则由外层走 pro VLM 兜底)
- 期望返回纯 JSON object(无 ```json 包裹、无解释文字)
"""
from __future__ import annotations
import json
import logging
import os
import time
from typing import Any
logger = logging.getLogger(__name__)
# DashScope OpenAI 兼容 endpoint
_BASE_URL = "https://dashscope.aliyuncs.com/compatible-mode/v1"
_FAST_MODEL = "qwen3.8-flash"
_DEFAULT_TIMEOUT = 12
_DEFAULT_MAX_TOKENS = 350
# 极简 system prompt:只给字段定义 + 硬性输出要求
_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'
@@ -46,7 +38,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'
@@ -59,17 +51,21 @@ _FAST_SYSTEM = (
_FAST_USER = "识别这张图片的人物穿搭与主体信息,只返回JSON对象。"
def _api_key() -> str | None:
return os.environ.get("DASHSCOPE_API_KEY")
# 默认模型
DEFAULT_LITE_MODEL = "doubao-seed-2-1-lite-260915"
DEFAULT_TIMEOUT = 8
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,93 +75,61 @@ def _strip_code_fence(s: str) -> str:
def call_fast_json(
img_url: str,
*,
timeout: int = _DEFAULT_TIMEOUT,
max_tokens: int = _DEFAULT_MAX_TOKENS,
model: str | None = None,
timeout: int = DEFAULT_TIMEOUT,
max_tokens: int = DEFAULT_MAX_TOKENS,
) -> dict[str, Any] | None:
"""调用 qwen3.8-flash 返回结构化 dict;失败/非 JSON 返回 None。"""
"""调用 lite VLM 返回结构化 dict;失败/非 JSON 返回 None。
注意:不做重试(外层竞速/降级逻辑负责),max_retries=0 由外层统一设置。
"""
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,
"response_format": {"type": "json_object"},
}
try:
resp = httpx.post(
url,
headers={"Authorization": f"Bearer {api_key}", "Content-Type": "application/json"},
json=payload,
from packages.shared.ai_client import get_doubao_client
client = get_doubao_client()
if not client.is_available:
logger.warning("[vision.v2] doubao client 不可用,跳过 fast_json")
return None
use_model = model or DEFAULT_LITE_MODEL
raw = client.vision_completion(
messages=[
{"role": "system", "content": _FAST_SYSTEM},
{"role": "user", "content": _FAST_USER},
],
images=[img_url],
temperature=0.1,
max_tokens=max_tokens,
timeout=timeout,
model=use_model,
)
elapsed = time.time() - t0
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]
)
if raw is None:
logger.warning("[vision.v2] fast_json 返回 None elapsed=%.1fs model=%s", elapsed, use_model)
return None
data = resp.json()
raw = (data.get("choices") or [{}])[0].get("message", {}).get("content")
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)
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",
_FAST_MODEL,
elapsed,
usage.get("prompt_tokens", 0),
usage.get("completion_tokens", 0),
reasoning_tokens,
)
text = _strip_code_fence(raw)
l, r = text.find("{"), text.rfind("}")
# 截到第一个 { 和最后一个 } 之间,容忍前后偶发文字
l = text.find("{")
r = text.rfind("}")
if l >= 0 and r > l:
text = text[l : r + 1]
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 完成 elapsed=%.1fs has_person=%s has_product=%s category=%s",
"[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"),