feat(orchestrator): 视觉 fallback 降级 + 文本融合角色化
- run_visual_analysis 仅 vision 角色参与, fallback 顺序降级(Gemini→NVIDIA NIM)首个成功即采用 - run_text_fusion 固定用 role=text 的 Ollama(qwen2.5:7b) 融合, 支持 num_predict - config 改为多模型池(gemini/nvidia vision + ollama text)
This commit is contained in:
@@ -97,45 +97,87 @@ class AIOrchestrator:
|
||||
frame_paths: List[str],
|
||||
frame_timestamps: List[str],
|
||||
known_members_context: str) -> Dict[str, str]:
|
||||
"""并行调用所有健康模型进行视觉分析"""
|
||||
model_outputs = {}
|
||||
max_timeout = max((a.get_timeout() for a in adapters), default=240)
|
||||
"""视觉分析阶段:仅 role=vision 的适配器参与
|
||||
|
||||
with ThreadPoolExecutor(max_workers=len(adapters)) as pool:
|
||||
orchestrator.mode:
|
||||
- fallback (默认): 按 config 顺序依次尝试,首个成功即采用(单元素 dict)
|
||||
- ensemble: 并行所有健康 vision 模型,全部成功结果都保留(交叉验证)
|
||||
"""
|
||||
vision_adapters = [a for a in adapters if getattr(a, 'role', 'vision') == 'vision']
|
||||
if not vision_adapters:
|
||||
logger.error("没有 vision 角色的可用适配器")
|
||||
return {}
|
||||
|
||||
mode = self.config.get('orchestrator', {}).get('mode', 'fallback')
|
||||
|
||||
if mode == 'ensemble':
|
||||
return self._run_visual_ensemble(
|
||||
vision_adapters, frame_paths, frame_timestamps, known_members_context)
|
||||
|
||||
# fallback: 顺序降级,首个成功即采用
|
||||
model_outputs = {}
|
||||
for adapter in vision_adapters:
|
||||
if adapter.get_circuit_breaker().is_open():
|
||||
logger.warning(f"[{adapter.provider_name}] 熔断器 OPEN,跳过")
|
||||
continue
|
||||
start = time.time()
|
||||
try:
|
||||
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()
|
||||
log_task(logger, 0, f'model_{adapter.provider_name}',
|
||||
f'视觉分析成功', duration_ms=duration_ms)
|
||||
model_outputs[adapter.provider_name] = output
|
||||
logger.info(f"fallback 采用 [{adapter.provider_name}],停止降级")
|
||||
break
|
||||
else:
|
||||
adapter.get_circuit_breaker().record_failure()
|
||||
logger.warning(f"[{adapter.provider_name}] 视觉分析返回空,降级下一模型")
|
||||
except Exception as e:
|
||||
logger.error(f"[{adapter.provider_name}] 视觉分析异常: {e}")
|
||||
adapter.get_circuit_breaker().record_failure()
|
||||
return model_outputs
|
||||
|
||||
def _run_visual_ensemble(self, vision_adapters, frame_paths,
|
||||
frame_timestamps, known_members_context) -> Dict[str, str]:
|
||||
"""并行调用所有健康 vision 模型,保留全部成功结果(交叉验证)"""
|
||||
model_outputs = {}
|
||||
max_timeout = max((a.get_timeout() for a in vision_adapters), default=240)
|
||||
with ThreadPoolExecutor(max_workers=len(vision_adapters)) as pool:
|
||||
futures = {}
|
||||
for adapter in adapters:
|
||||
for adapter in vision_adapters:
|
||||
if adapter.get_circuit_breaker().is_open():
|
||||
logger.warning(f"[{adapter.provider_name}] 熔断器 OPEN,跳过")
|
||||
continue
|
||||
future = pool.submit(
|
||||
adapter.analyze_frames,
|
||||
frame_paths, frame_timestamps, known_members_context
|
||||
)
|
||||
frame_paths, frame_timestamps, known_members_context)
|
||||
futures[future] = adapter.provider_name
|
||||
|
||||
for future in as_completed(futures, timeout=max_timeout + 10):
|
||||
provider = futures[future]
|
||||
start = time.time()
|
||||
try:
|
||||
adapter = next(a for a in adapters if a.provider_name == provider)
|
||||
adapter = next(a for a in vision_adapters if a.provider_name == provider)
|
||||
output = future.result(timeout=adapter.get_timeout())
|
||||
duration_ms = int((time.time() - start) * 1000)
|
||||
if output:
|
||||
model_outputs[provider] = output
|
||||
adapter.get_circuit_breaker().record_success()
|
||||
log_task(logger, 0, f'model_{provider}', f'视觉分析成功,输出长度={len(output)}', duration_ms=duration_ms)
|
||||
log_task(logger, 0, f'model_{provider}',
|
||||
f'视觉分析成功,输出长度={len(output)}', duration_ms=duration_ms)
|
||||
else:
|
||||
adapter.get_circuit_breaker().record_failure()
|
||||
logger.warning(f"[{provider}] 视觉分析返回空")
|
||||
except FuturesTimeout:
|
||||
logger.warning(f"[{provider}] 视觉分析超时")
|
||||
adapter = next(a for a in adapters if a.provider_name == provider)
|
||||
adapter = next(a for a in vision_adapters if a.provider_name == provider)
|
||||
adapter.get_circuit_breaker().record_failure()
|
||||
except Exception as e:
|
||||
logger.error(f"[{provider}] 视觉分析异常: {e}")
|
||||
adapter = next(a for a in adapters if a.provider_name == provider)
|
||||
adapter = next(a for a in vision_adapters if a.provider_name == provider)
|
||||
adapter.get_circuit_breaker().record_failure()
|
||||
|
||||
return model_outputs
|
||||
|
||||
def run_text_fusion(self, model_outputs: Dict[str, str],
|
||||
@@ -152,17 +194,19 @@ class AIOrchestrator:
|
||||
known_members=known_members_context or '(暂无已知成员)'
|
||||
)
|
||||
|
||||
# 调用 Ollama 纯文本模式
|
||||
ollama_cfg = next(
|
||||
(cfg for cfg in self.config.get('models', []) if cfg.get('provider') == 'ollama'),
|
||||
None
|
||||
# 调用文本角色模型(role=text,默认 ollama / qwen2.5:7b)做融合
|
||||
text_cfg = next(
|
||||
(cfg for cfg in self.config.get('models', []) if cfg.get('role') == 'text'), None
|
||||
) or next(
|
||||
(cfg for cfg in self.config.get('models', []) if cfg.get('provider') == 'ollama'), None
|
||||
)
|
||||
if not ollama_cfg:
|
||||
raise VLMOutputInvalidError("没有 Ollama 配置,无法执行文本融合")
|
||||
if not text_cfg:
|
||||
raise VLMOutputInvalidError("没有文本角色模型配置,无法执行文本融合")
|
||||
|
||||
base_url = ollama_cfg.get('base_url', 'http://localhost:11434')
|
||||
model_name = ollama_cfg.get('model_name', 'llava-phi3')
|
||||
fusion_timeout = self.timeout_cfg.get('vlm_fusion', 120)
|
||||
base_url = text_cfg.get('base_url', 'http://localhost:11434')
|
||||
model_name = text_cfg.get('model_name', 'qwen2.5:7b')
|
||||
fusion_timeout = self.timeout_cfg.get('vlm_fusion', 300)
|
||||
num_predict = text_cfg.get('num_predict', 1024)
|
||||
|
||||
start = time.time()
|
||||
resp = requests.post(
|
||||
@@ -172,7 +216,7 @@ class AIOrchestrator:
|
||||
"prompt": prompt,
|
||||
"stream": False,
|
||||
"format": "json",
|
||||
"options": {"temperature": 0.0}
|
||||
"options": {"temperature": 0.0, "num_predict": num_predict}
|
||||
},
|
||||
timeout=fusion_timeout
|
||||
)
|
||||
|
||||
Reference in New Issue
Block a user