""" 消费者线程 - 从 SQLite 队列消费任务,限速处理 策略:Gemini 优先 → NVIDIA 兜底(与现有 orchestrator 一致) 速率:按 API 限制速度的 2 倍设置突发容量 """ import os import time import threading from typing import Optional from ..logger import setup_logger from ..config_loader import load_config from ..rate_limiter import RateLimiter from ..ai_orchestrator.orchestrator import AIOrchestrator from ..video_preprocessor.preprocessor import VideoPreprocessor from . import queue_manager logger = setup_logger('fam-edge.consumer') # 从配置加载速率限制参数 _cfg = load_config() _queue_cfg = _cfg.get('queue', {}) _rate_cfg = _queue_cfg.get('rate_limit', {}) GEMINI_RPM = _rate_cfg.get('gemini_rpm', 1000) NVIDIA_RPM = _rate_cfg.get('nvidia_rpm', 40) BURST_FACTOR = _rate_cfg.get('burst_factor', 2) POLL_INTERVAL = _queue_cfg.get('poll_interval', 10) # 同步 SQLite DB 路径到环境变量(供 queue_manager 读取) os.environ.setdefault('FAM_QUEUE_DB', _queue_cfg.get('db_path', '/opt/fam-edge/data/fam_queue.db')) os.environ.setdefault('FAM_UPLOAD_DIR', _queue_cfg.get('upload_dir', '/tmp/fam_uploads')) class Consumer: def __init__(self): self._running = False self._thread = None self._orchestrator = AIOrchestrator() self._rate_limiter = RateLimiter() self._rate_limiter.register('gemini', GEMINI_RPM, burst_factor=BURST_FACTOR) self._rate_limiter.register('nvidia', NVIDIA_RPM, burst_factor=BURST_FACTOR) self._poll_interval = POLL_INTERVAL def _process_one(self, task: dict) -> bool: task_id = task['id'] nas_task_id = task['nas_task_id'] video_path = task['video_path'] logger.info(f"[nas_task={nas_task_id}] 消费者开始处理") preprocessor = None try: preprocessor = VideoPreprocessor(nas_task_id) task_data = { "task_id": nas_task_id, "camera_name": task.get('camera_name', ''), "event_start_time": task.get('event_start_time', ''), "event_end_time": "", "known_members_context": task.get('known_members_context', ''), } result = self._orchestrator.process_push_task( task_data, video_path, preprocessor, self._rate_limiter ) if result.get('status') == 'success': import json queue_manager.mark_success(task_id, json.dumps(result, ensure_ascii=False)) logger.info(f"[nas_task={nas_task_id}] 消费者处理成功") return True else: error = result.get('error_message', 'unknown') stage = result.get('failure_stage', '') queue_manager.mark_failed(task_id, error, stage) logger.error(f"[nas_task={nas_task_id}] 消费者处理失败: {error}") return False except Exception as e: logger.error(f"[nas_task={nas_task_id}] 消费者异常: {e}", exc_info=True) queue_manager.mark_failed(task_id, str(e), 'process') return False finally: if preprocessor is not None: preprocessor.cleanup() try: if video_path and __import__('os').path.exists(video_path): __import__('os').remove(video_path) logger.info(f"[nas_task={nas_task_id}] 清理视频文件: {video_path}") except Exception: pass def _run(self): logger.info(f"消费者线程启动,轮询间隔 {self._poll_interval}s") logger.info(f"速率限制: Gemini {GEMINI_RPM}RPM x2 burst, NVIDIA {NVIDIA_RPM}RPM x2 burst") while self._running: try: task = queue_manager.claim_next() if task is None: time.sleep(self._poll_interval) continue self._process_one(task) except Exception as e: logger.error(f"消费者循环异常: {e}", exc_info=True) time.sleep(self._poll_interval) def start(self): if self._running: return self._running = True self._thread = threading.Thread(target=self._run, daemon=True, name='consumer') self._thread.start() logger.info("消费者线程已启动") def stop(self): self._running = False logger.info("消费者线程已停止") _consumer: Optional[Consumer] = None def get_consumer() -> Consumer: global _consumer if _consumer is None: _consumer = Consumer() return _consumer