diff --git a/fam-edge/config/config.yaml b/fam-edge/config/config.yaml index 0b82fae..5f0409d 100644 --- a/fam-edge/config/config.yaml +++ b/fam-edge/config/config.yaml @@ -1,4 +1,4 @@ -# FAM-Edge 配置文件 (Oracle 端) - 实际部署配置 +# FAM-Edge 配置文件 (Oracle 端) - 多模型池配置 # Tailscale: Oracle=100.74.137.126, NAS=100.70.234.39 # NAS 端回调地址 @@ -13,6 +13,11 @@ server: port: 5000 max_concurrent_tasks: 1 +# 编排调度模式: fallback(顺序降级, 默认) | ensemble(并行交叉验证) +orchestrator: + mode: "fallback" + overall_timeout: 600 + # 关键帧筛选参数(自适应:帧数随视频时长动态计算) video: candidate_per_minute: 2 # 每分钟粗抽候选帧数 @@ -27,8 +32,6 @@ video: max_long_edge: 1024 # 超时(秒) -# vlm_visual 实测: 1024px 帧视觉编码 ~36s/帧 + 生成 ~12s/60token (Oracle ARM CPU) -# 30min 视频 12 帧 × ~50s ≈ 600s,超时需覆盖最坏情况 timeout: download: 60 vlm_visual: 600 @@ -36,26 +39,39 @@ timeout: callback: 30 overall: 1800 -# 模型清单 +# 多模型池配置 +# 视觉分析: Gemini(主) -> NVIDIA NIM(备) 顺序降级; 全失败 -> 任务 FAILED 走重试 +# 文本融合/对话: 本地 Ollama qwen2.5:7b 专职 (不参与视觉) models: - - provider: "ollama" - enabled: true - model_name: "llava-phi3" - base_url: "http://localhost:11434" - timeout: 600 - # num_predict 必须小: ARM CPU ~5 tok/s,500 会单帧跑数分钟触发超时 - num_predict: 60 - circuit_breaker: - enabled: false - threshold: 5 - cooldown: 900 - - provider: "gemini" - enabled: false - model_name: "gemini-1.5-flash" - api_key: "" - timeout: 8 + role: "vision" + enabled: true + model_name: "gemini-flash-latest" # v1beta 下 gemini-1.5-flash 会 404 + api_key: "${GEMINI_API_KEY}" + timeout: 15 circuit_breaker: enabled: true - threshold: 5 - cooldown: 900 + threshold: 3 + cooldown: 600 + + - provider: "nvidia" + role: "vision" + enabled: true + model_name: "meta/llama-3.2-11b-vision-instruct" + base_url: "https://integrate.api.nvidia.com/v1" + api_key: "${NVIDIA_API_KEY}" + timeout: 20 + circuit_breaker: + enabled: true + threshold: 3 + cooldown: 600 + + - provider: "ollama" + role: "text" + enabled: true + model_name: "qwen2.5:7b" + base_url: "http://localhost:11434" + timeout: 300 + num_predict: 1024 + circuit_breaker: + enabled: false diff --git a/fam-edge/src/fam_edge/ai_orchestrator/orchestrator.py b/fam-edge/src/fam_edge/ai_orchestrator/orchestrator.py index df4ad3c..d76e750 100644 --- a/fam-edge/src/fam_edge/ai_orchestrator/orchestrator.py +++ b/fam-edge/src/fam_edge/ai_orchestrator/orchestrator.py @@ -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 )