feat: NVIDIA 切换到原生视频输入 — nemotron-3-nano-omni + 集锦视频单次调用

调研结论: build.nvidia.com 免费托管 API 上 video-llama3-8b 与
qwen2.5-vl-72b 已下线(404),nvidia/nemotron-3-nano-omni-30b-a3b-reasoning
可用(200, 40RPM 免费额度内),原生支持 video_url 输入(MP4 base64)。

实现:
1. nvidia_adapter 新增 analyze_video: 按关键帧时间点截取 ±1.5s 片段
   (drawtext 叠加时间戳,连字符避免冒号转义)拼集锦视频,640 宽 CRF28,
   base64 后经 video_url 单次调用,输出全 schema JSON(frame_details +
   global_summary + entities_json)并归一化对齐时间戳
2. analyze_frames 保留为无视频文件时的降级路径; chat max_tokens
   512→2048(reasoning 模型 token 消耗大); timeout 20→120s
3. orchestrator.run_visual_analysis 增加 video_path 参数,fallback 循环
   对支持 analyze_video 的适配器优先走视频模式,失败自动降级逐帧

实测(360MB 测试视频, 3 关键帧): 集锦 107KB, 全程 37s, 动态动作识别准确
(走动→坐沙发→坐餐桌),跨片段综合摘要正常 — 显著优于旧逐帧静态识别。
This commit is contained in:
ericwyuan
2026-08-20 17:21:33 +08:00
parent 8a62c46194
commit 55633d3302
3 changed files with 228 additions and 17 deletions

View File

@@ -53,9 +53,13 @@ class AIOrchestrator:
frame_paths: List[str],
frame_timestamps: List[str],
known_members_context: str,
rate_limiter=None) -> Dict[str, dict]:
rate_limiter=None,
video_path: str = None) -> Dict[str, dict]:
"""视觉分析阶段:仅 role=vision 的适配器参与
支持 analyze_video 的适配器(如 NVIDIA Omni优先走原生视频输入
失败自动降级回逐帧图片模式。
orchestrator.mode:
- fallback (默认): 按 config 顺序依次尝试,首个成功即采用(单元素 dict
- ensemble: 并行所有健康 vision 模型,全部成功结果都保留(交叉验证)
@@ -88,8 +92,20 @@ class AIOrchestrator:
continue
start = time.time()
try:
output = adapter.analyze_frames(
frame_paths, frame_timestamps, known_members_context)
output = None
if video_path and hasattr(adapter, 'analyze_video'):
try:
logger.info(f"[{adapter.provider_name}] 尝试原生视频输入分析")
output = adapter.analyze_video(
video_path, frame_timestamps, known_members_context)
if not output:
logger.warning(f"[{adapter.provider_name}] 视频模式失败,降级逐帧模式")
except Exception as ve:
logger.warning(f"[{adapter.provider_name}] 视频模式异常: {ve},降级逐帧模式")
output = None
if not output:
output = adapter.analyze_frames(
frame_paths, frame_timestamps, known_members_context)
duration_ms = int((time.time() - start) * 1000)
if output:
adapter.get_circuit_breaker().record_success()
@@ -354,7 +370,8 @@ class AIOrchestrator:
# 3. 并行视觉分析
model_outputs = self.run_visual_analysis(
healthy_adapters, compressed_frames, frame_timestamps, known_members
healthy_adapters, compressed_frames, frame_timestamps,
known_members, video_path=video_path
)
if not model_outputs:
@@ -450,7 +467,7 @@ class AIOrchestrator:
# 3. 并行视觉分析
model_outputs = self.run_visual_analysis(
healthy_adapters, compressed_frames, frame_timestamps,
known_members, rate_limiter
known_members, rate_limiter, video_path=video_path
)
if not model_outputs:
raise Exception('All models failed in visual analysis')

View File

@@ -2,16 +2,23 @@
NvidiaVisionAdapter - NVIDIA NIM 云端 VLM 适配器
provider_name = "nvidia"
模型: meta/llama-3.2-11b-vision-instruct
模型: nvidia/nemotron-3-nano-omni-30b-a3b-reasoning (Omni, 原生视频输入)
角色: vision (视觉分析直出结构化 JSON) + 智能问答
SDK: openai (NIM 兼容 OpenAI API 规范)
限制: NIM 单次请求最多 1 张图 -> 适配器内部逐帧调用,再聚合成 frame_details
熔断器: 启用
视频模式 (analyze_video): 按关键帧时间点截取 ±1.5s 片段拼接集锦视频
(片段左上角叠加原始时间戳)base64 后经 video_url 单次调用 —
模型看到动态画面而非静态帧,动作/轨迹识别显著优于逐帧图片。
图片模式 (analyze_frames): 逐帧 image_url 调用(无视频文件时的降级路径)。
注意: nemotron-omni 是 reasoning 模型max_tokens 需给足reasoning 消耗 token
"""
import os
import base64
import json
import re
import subprocess
import tempfile
from typing import Dict, List, Optional
from .base_adapter import BaseModelAdapter
@@ -25,16 +32,20 @@ try:
except ImportError:
OpenAI = None
VIDEO_SEGMENT_PAD = 1.5 # 关键帧前后各截取秒数
HIGHLIGHT_WIDTH = 640 # 集锦视频宽度(保持宽高比)
class NvidiaVisionAdapter(BaseModelAdapter):
"""NVIDIA NIM 云端 VLM 适配器 (逐帧结构化 + 聚合; 文本问答)"""
"""NVIDIA NIM 云端 VLM 适配器 (视频集锦单次调用; 逐帧降级; 文本问答)"""
def __init__(self, config: dict):
super().__init__("nvidia", config)
self.model_name = config.get('model_name', 'meta/llama-3.2-11b-vision-instruct')
self.model_name = config.get(
'model_name', 'nvidia/nemotron-3-nano-omni-30b-a3b-reasoning')
self.api_key = self._resolve_key(config.get('api_key', ''))
self.base_url = config.get('base_url', 'https://integrate.api.nvidia.com/v1')
self.timeout = config.get('timeout', 20)
self.timeout = config.get('timeout', 120)
cb_cfg = config.get('circuit_breaker', {})
self._cb = CircuitBreaker(
threshold=cb_cfg.get('threshold', 3),
@@ -67,7 +78,187 @@ class NvidiaVisionAdapter(BaseModelAdapter):
return False
# ------------------------------------------------------------------
# 视觉分析逐帧调用NIM 限 1 图/请求),聚合为 frame_details
# 视频模式:集锦视频 + video_url 单次调用(主路径)
# ------------------------------------------------------------------
@staticmethod
def _ts_to_seconds(ts: str) -> float:
"""'HH:MM:SS''HH:MM:SS.mmm' -> 秒"""
parts = str(ts).strip().split(':')
try:
if len(parts) == 3:
return int(parts[0]) * 3600 + int(parts[1]) * 60 + float(parts[2])
if len(parts) == 2:
return int(parts[0]) * 60 + float(parts[1])
return float(ts)
except ValueError:
return -1.0
def _build_highlight_video(self, video_path: str,
frame_timestamps: List[str]) -> Optional[str]:
"""按关键帧时间点截取 ±pad 秒片段,叠加时间戳后拼接集锦视频
时间戳以 00-05-00 形式叠加(连字符避免 ffmpeg drawtext 冒号转义)。
"""
clips = []
for ts in frame_timestamps:
sec = self._ts_to_seconds(ts)
if sec < 0:
continue
start = max(0.0, sec - VIDEO_SEGMENT_PAD)
label = str(ts).replace(':', '-')
clips.append((start, label))
if not clips:
return None
out_path = os.path.join(
tempfile.mkdtemp(prefix='nim_highlight_'), 'highlight.mp4')
cmd = ['ffmpeg', '-y', '-loglevel', 'error']
for start, _ in clips:
cmd += ['-ss', f'{start:.2f}', '-t', f'{VIDEO_SEGMENT_PAD * 2}', '-i', video_path]
parts = []
for i, (_, label) in enumerate(clips):
parts.append(
f"[{i}:v]scale={HIGHLIGHT_WIDTH}:-2,"
f"drawtext=text='ts {label}':x=8:y=8:fontsize=22:"
f"fontcolor=white:box=1:boxcolor=black@0.6[v{i}]")
concat_in = ''.join(f'[v{i}]' for i in range(len(clips)))
parts.append(f'{concat_in}concat=n={len(clips)}:v=1:a=0[out]')
cmd += ['-filter_complex', ';'.join(parts), '-map', '[out]',
'-c:v', 'libx264', '-preset', 'veryfast', '-crf', '28',
'-an', out_path]
try:
subprocess.run(cmd, check=True, capture_output=True, timeout=120)
except subprocess.TimeoutExpired:
logger.warning("NVIDIA 集锦视频生成超时")
return None
except subprocess.CalledProcessError as e:
logger.warning(f"NVIDIA 集锦视频生成失败: {e.stderr.decode()[:200] if e.stderr else e}")
return None
size = os.path.getsize(out_path)
logger.info(f"NVIDIA 集锦视频生成: {len(clips)} 片段, {size // 1024}KB")
if size < 1024:
return None
return out_path
def analyze_video(self, video_path: str,
frame_timestamps: List[str],
known_members_context: str) -> Optional[Dict]:
"""原生视频输入分析: 集锦片段 -> video_url 单次调用"""
if self._cb.is_open():
logger.warning("NVIDIA 熔断器 OPEN跳过视频分析")
return None
if self._client is None:
logger.warning("NVIDIA 客户端未初始化,跳过视频分析")
return None
highlight = self._build_highlight_video(video_path, frame_timestamps)
if not highlight:
logger.warning("NVIDIA 集锦视频不可用,降级逐帧模式")
return None
try:
with open(highlight, 'rb') as f:
b64 = base64.b64encode(f.read()).decode('utf-8')
except Exception as e:
logger.warning(f"NVIDIA 读取集锦视频失败: {e}")
return None
finally:
try:
os.remove(highlight)
os.rmdir(os.path.dirname(highlight))
except OSError:
pass
prompt = self._build_video_prompt(frame_timestamps, known_members_context)
try:
resp = self._client.chat.completions.create(
model=self.model_name,
messages=[{"role": "user", "content": [
{"type": "text", "text": prompt},
{"type": "video_url", "video_url": {
"url": f"data:video/mp4;base64,{b64}"}}
]}],
temperature=0.2,
max_tokens=3072,
timeout=self.timeout
)
content = resp.choices[0].message.content
if not content:
logger.warning("NVIDIA 视频分析返回空 content")
return None
data = self._parse_single_frame_json(content)
if not data or 'frame_details' not in data:
logger.warning(f"NVIDIA 视频 JSON 解析失败: {content[:150]}")
return None
frame_details = self._normalize_frame_details(data, frame_timestamps)
if not frame_details:
return None
self._cb.record_success()
logger.info(f"NVIDIA 视频分析完成frame_details={len(frame_details)}")
result = {"frame_details": frame_details}
if data.get('global_summary'):
result['global_summary'] = str(data['global_summary'])
if data.get('entities_json'):
result['entities_json'] = data['entities_json']
return result
except Exception as e:
self._cb.record_failure()
logger.warning(f"NVIDIA 视频分析异常: {e}")
return None
def _normalize_frame_details(self, data: dict,
frame_timestamps: List[str]) -> List[Dict]:
"""归一化模型输出的 frame_details按已知时间戳对齐"""
details = []
for i, item in enumerate(data.get('frame_details', []), 1):
if not isinstance(item, dict):
continue
ts = str(item.get('frame_timestamp',
frame_timestamps[i - 1] if i <= len(frame_timestamps) else ''))
details.append({
"frame_index": i,
"frame_timestamp": ts,
"person": str(item.get('person', '无人')),
"action": str(item.get('action', '')),
"clothing": str(item.get('clothing', '')),
"is_attention_event": bool(item.get('is_attention_event', False)),
"source_providers": ["nvidia"],
})
return details
def _build_video_prompt(self, frame_timestamps: List[str],
known_members: str) -> str:
ts_list = '\n'.join(f' 片段{i}: 原始时间 {ts}' for i, ts in enumerate(frame_timestamps, 1))
return f"""你是家庭监控视频分析助手。下面的视频是由一段长时间监控录像中抽取的片段集锦,
{len(frame_timestamps)} 个片段(每个约 3 秒),按顺序拼接。每个片段左上角叠加了
原始时间戳ts 后的 00-05-00 表示 00:05:00
片段时间对照:
{ts_list}
只输出合法 JSON不要 markdown、不要解释结构如下
{{
"frame_details": [
{{
"frame_timestamp": "<片段原始时间>",
"person": "片段中的人物或'无人'",
"action": "片段中人物的动作(动态观察,如走动/跑动/坐下)",
"clothing": "衣着(颜色+类型)",
"is_attention_event": false
}}
],
"global_summary": "整段录像的综合摘要",
"entities_json": [{{"person": "人物名或人物X", "action": "行为概括", "clothing": "衣着"}}]
}}
规则:
1. 只描述客观画面,不猜测。
2. 已知家庭成员(按特征匹配,匹配到用 real_name否则用"人物X"
{known_members or '(暂无已知成员)'}
3. is_attention_event跌倒、危险、异常哭闹等需关注事件没有则为 false
4. 无人出现的片段 person 填"无人"action 填"""""
# ------------------------------------------------------------------
# 图片模式:逐帧调用(降级路径,无视频文件时使用)
# ------------------------------------------------------------------
def analyze_frames(self, frame_paths: List[str],
frame_timestamps: List[str],
@@ -184,9 +375,9 @@ class NvidiaVisionAdapter(BaseModelAdapter):
4. 没有人物出现的帧 person 填"无人"action 填"""""
# ------------------------------------------------------------------
# 智能问答:纯文本
# 智能问答:纯文本reasoning 模型max_tokens 需给足)
# ------------------------------------------------------------------
def chat(self, prompt: str, max_tokens: int = 512) -> Optional[str]:
def chat(self, prompt: str, max_tokens: int = 2048) -> Optional[str]:
if self._client is None:
logger.warning("NVIDIA 客户端未初始化,跳过问答")
return None