From 02da23ef4291d1c8e3158d515636e4dcc1df6c2c Mon Sep 17 00:00:00 2001 From: ericwyuan Date: Thu, 20 Aug 2026 12:07:09 +0800 Subject: [PATCH] =?UTF-8?q?feat:=20=E5=BC=82=E6=AD=A5=E4=BB=BB=E5=8A=A1?= =?UTF-8?q?=E9=98=9F=E5=88=97=E6=9E=B6=E6=9E=84=20-=20SQLite=E9=98=9F?= =?UTF-8?q?=E5=88=97=20+=20=E9=80=9F=E7=8E=87=E9=99=90=E5=88=B6=20+=20NAS?= =?UTF-8?q?=20Poller?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Edge端: - 新增 SQLite 异步任务队列 (queue_manager + consumer) - 新增 TokenBucket 速率限制器 (Gemini 1000 RPM, NVIDIA 40 RPM, burst 2x) - 新增 /api/edge/video/enqueue + /api/edge/results 端点 - 消费者线程从队列消费任务,按速率限制调用AI模型 - orchestrator 集成 rate_limiter,Gemini优先→NVIDIA兜底 NAS端: - Dispatcher 重构为 enqueue 模式(上传后立即返回,不等结果) - 新增 Poller 线程(定期从Edge拉取结果写 MariaDB) - app.py 启动 Poller,config.yaml 新增 poller 配置 - db_layer 更新 valid_stages 添加 'process' --- fam-core/config/config.yaml | 15 +- fam-core/src/fam_core/app.py | 12 +- fam-core/src/fam_core/db_layer.py | 2 +- .../src/fam_core/dispatcher/dispatcher.py | 78 ++++---- fam-core/src/fam_core/poller/__init__.py | 0 fam-core/src/fam_core/poller/poller.py | 128 +++++++++++++ fam-edge/config/config.yaml | 18 +- .../fam_edge/ai_orchestrator/orchestrator.py | 31 +++- .../src/fam_edge/api_gateway/api_gateway.py | 94 +++++++++- fam-edge/src/fam_edge/app.py | 15 +- fam-edge/src/fam_edge/queue/__init__.py | 5 + fam-edge/src/fam_edge/queue/consumer.py | 128 +++++++++++++ fam-edge/src/fam_edge/queue/queue_manager.py | 174 ++++++++++++++++++ fam-edge/src/fam_edge/rate_limiter.py | 48 +++++ 14 files changed, 686 insertions(+), 62 deletions(-) create mode 100644 fam-core/src/fam_core/poller/__init__.py create mode 100644 fam-core/src/fam_core/poller/poller.py create mode 100644 fam-edge/src/fam_edge/queue/__init__.py create mode 100644 fam-edge/src/fam_edge/queue/consumer.py create mode 100644 fam-edge/src/fam_edge/queue/queue_manager.py create mode 100644 fam-edge/src/fam_edge/rate_limiter.py diff --git a/fam-core/config/config.yaml b/fam-core/config/config.yaml index 0cb0bed..7d6e7d4 100644 --- a/fam-core/config/config.yaml +++ b/fam-core/config/config.yaml @@ -1,7 +1,7 @@ # FAM-Core 配置文件 (NAS 端) - 实际部署配置 # 注: Tailscale 防火墙待修复,当前 edge_url 使用 Oracle 公网 IP -# 推送模式: NAS 上传视频到 Edge /api/edge/video/push,结果同步随响应返回, -# 无需 Oracle 反向访问 NAS(webhook/Video-Server 拉取均不再使用) +# 异步队列模式: NAS 上传视频到 Edge /api/edge/video/enqueue 入队 → Edge 消费者异步处理 +# → NAS Poller 定期从 /api/edge/results 拉取结果写库 # Ollama 未对外暴露,chat_handler 通过 FAM-Edge 代理 server: @@ -27,9 +27,16 @@ scheduler: dispatcher: poll_interval: 30 - edge_url: "http://129.146.203.203:5000/api/edge/video/push" + edge_url: "http://129.146.203.203:5000/api/edge/video/enqueue" max_retries: 3 - push_timeout: 1800 + upload_timeout: 300 # 仅视频上传时间(不含 AI 处理) + stale_timeout: 3600 # PROCESSING 超时回收(Edge 处理 + 队列等待) + +poller: + poll_interval: 30 + results_url: "http://129.146.203.203:5000/api/edge/results" + batch_size: 10 + timeout: 30 video_server: base_url: "http://127.0.0.1:8000/media" diff --git a/fam-core/src/fam_core/app.py b/fam-core/src/fam_core/app.py index f04f5d7..5428abb 100644 --- a/fam-core/src/fam_core/app.py +++ b/fam-core/src/fam_core/app.py @@ -1,7 +1,7 @@ """ FAM-Core 主应用 - Flask 单进程 -承载: Task-Scheduler / Dispatcher / Event-Receiver / Chat-Handler / Member-Manager / Video-Server +承载: Task-Scheduler / Dispatcher / Poller / Event-Receiver / Chat-Handler / Member-Manager / Video-Server """ import os import sys @@ -14,6 +14,7 @@ from .config_loader import load_config from .logger import setup_logger from .scheduler.scheduler import TaskScheduler from .dispatcher.dispatcher import Dispatcher +from .poller.poller import Poller from .event_receiver.event_receiver import event_bp from .chat_handler.chat_handler import chat_bp from .member_manager.member_manager import member_bp @@ -37,6 +38,7 @@ def health(): # 初始化后台线程 _scheduler = None _dispatcher = None +_poller = None try: _scheduler = TaskScheduler() @@ -52,6 +54,13 @@ try: except Exception as e: logger.error(f"Dispatcher 启动失败: {e}") +try: + _poller = Poller() + _poller.start() + logger.info("Poller 已启动") +except Exception as e: + logger.error(f"Poller 启动失败: {e}") + @app.route('/api/status', methods=['GET']) def status(): @@ -59,6 +68,7 @@ def status(): return jsonify({ "scheduler_running": _scheduler._running if _scheduler else False, "dispatcher_running": _dispatcher._running if _dispatcher else False, + "poller_running": _poller._running if _poller else False, }), 200 diff --git a/fam-core/src/fam_core/db_layer.py b/fam-core/src/fam_core/db_layer.py index 1c5f9d0..c8f9253 100644 --- a/fam-core/src/fam_core/db_layer.py +++ b/fam-core/src/fam_core/db_layer.py @@ -92,7 +92,7 @@ def get_tasks_by_status(status: str, limit=10) -> List[Dict]: def update_task_status(task_id: int, status: str, error_message: str = None, failure_stage: str = None): """更新任务状态""" - valid_stages = {'download', 'extract', 'vlm_visual', 'vlm_fusion', 'callback'} + valid_stages = {'download', 'extract', 'vlm_visual', 'vlm_fusion', 'callback', 'process'} if failure_stage and failure_stage not in valid_stages: failure_stage = 'callback' conn = get_conn() diff --git a/fam-core/src/fam_core/dispatcher/dispatcher.py b/fam-core/src/fam_core/dispatcher/dispatcher.py index 143dd59..a1d90f1 100644 --- a/fam-core/src/fam_core/dispatcher/dispatcher.py +++ b/fam-core/src/fam_core/dispatcher/dispatcher.py @@ -1,11 +1,12 @@ """ -Dispatcher - 30s 轮询 PENDING 任务,推送视频至 Edge(推送模式) +Dispatcher - 30s 轮询 PENDING 任务,上传视频至 Edge 异步队列 -流程(纯单向通信,NAS → Oracle,无需 Oracle 反向访问 NAS): +流程(异步队列模式): 1. 读取任务对应的本地视频文件 -2. multipart 上传至 Edge /api/edge/video/push(附已知成员清单等元数据) -3. Edge 同步抽帧+分析,结果直接随 HTTP 响应返回 -4. Dispatcher 收到响应后直接写 monitor_events/event_details,任务标记 SUCCESS +2. multipart 上传至 Edge /api/edge/video/enqueue(附已知成员清单等元数据) +3. Edge 保存视频 + 入 SQLite 队列,立即返回 202 +4. Dispatcher 标记任务为 PROCESSING(已派发,等待 Poller 拉取结果) +5. Poller 线程定期从 Edge /api/edge/results 拉取结果,写库后标记 SUCCESS 退避重试: min(60 * (retry_count + 1) * 2, 600) 秒 """ @@ -18,22 +19,23 @@ from datetime import datetime, timedelta from ..logger import setup_logger, log_task from ..config_loader import load_config from .. import db_layer -from ..event_receiver.event_receiver import apply_success_event logger = setup_logger('fam-core.dispatcher') class Dispatcher: - """任务下发器,30s 轮询(推送模式)""" + """任务下发器,30s 轮询(异步队列模式)""" def __init__(self): cfg = load_config() self.poll_interval = cfg.get('dispatcher', {}).get('poll_interval', 30) - self.edge_url = cfg.get('dispatcher', {}).get('edge_url', 'http://localhost:5000/api/edge/video/push') + self.edge_url = cfg.get('dispatcher', {}).get('edge_url', 'http://localhost:5000/api/edge/video/enqueue') self.max_retries = cfg.get('dispatcher', {}).get('max_retries', 3) - # 推送+分析全程同步,耗时较长(大视频上传 + ARM 多帧分析) - self.push_timeout = cfg.get('dispatcher', {}).get('push_timeout', 1800) + # enqueue 模式只需上传时间,不含 AI 处理时间 + self.upload_timeout = cfg.get('dispatcher', {}).get('upload_timeout', 300) self.camera_name = cfg.get('scheduler', {}).get('camera_name', '默认摄像头') + # PROCESSING 超时回收:Edge 处理 + 队列等待可能很长 + self.stale_timeout = cfg.get('dispatcher', {}).get('stale_timeout', 3600) self._running = False self._thread = None @@ -61,7 +63,6 @@ class Dispatcher: "event_start_time": "", "event_end_time": "", } - # 用文件 mtime 近似事件开始时间 try: mtime = os.path.getmtime(video_path) start_dt = datetime.fromtimestamp(mtime) @@ -71,11 +72,10 @@ class Dispatcher: return payload def _dispatch_one(self, task): - """推送单个任务:上传视频 → 同步等结果 → 直接写库""" + """上传视频至 Edge 异步队列,立即返回""" task_id = task['task_id'] video_path = task['video_path'] - # 本地视频必须存在,否则直接失败(重试无意义) if not video_path or not os.path.isfile(video_path): db_layer.update_task_status( task_id, 'FAILED', @@ -89,7 +89,7 @@ class Dispatcher: size_mb = os.path.getsize(video_path) / (1024 * 1024) db_layer.update_task_status(task_id, 'PROCESSING') log_task(logger, task_id, 'dispatcher', - f'推送视频至 Edge: {self.edge_url} ({size_mb:.1f}MB)') + f'上传视频至 Edge 队列: {self.edge_url} ({size_mb:.1f}MB)') try: with open(video_path, 'rb') as fh: @@ -97,41 +97,29 @@ class Dispatcher: self.edge_url, data=payload, files={'video': (os.path.basename(video_path), fh, 'video/mp4')}, - timeout=self.push_timeout + timeout=self.upload_timeout ) except requests.RequestException as e: - logger.error(f"[task_id={task_id}] 推送失败: {e}") + logger.error(f"[task_id={task_id}] 上传失败: {e}") self._schedule_retry(task) return + if resp.status_code == 202: + try: + data = resp.json() + queue_id = data.get('queue_id', '?') + logger.info(f"[task_id={task_id}] 已入 Edge 队列 (queue_id={queue_id}),等待 Poller 拉取结果") + except ValueError: + logger.info(f"[task_id={task_id}] 已入 Edge 队列,等待 Poller 拉取结果") + return + if resp.status_code == 429: - logger.warning(f"[task_id={task_id}] Edge 忙 (429),回到 PENDING 稍后重试") + logger.warning(f"[task_id={task_id}] Edge 队列满 (429),回到 PENDING 稍后重试") db_layer.update_task_status(task_id, 'PENDING') return - if resp.status_code != 200: - logger.error(f"[task_id={task_id}] Edge 返回异常状态码: {resp.status_code}") - self._schedule_retry(task) - return - - try: - result = resp.json() - except ValueError: - result = {} - if result.get('status') == 'success': - try: - event_id = apply_success_event(task_id, result) - log_task(logger, task_id, 'dispatcher', f'推送任务完成: event_id={event_id}') - except Exception as e: - logger.error(f"[task_id={task_id}] 结果落库失败: {e}", exc_info=True) - db_layer.update_task_status( - task_id, 'FAILED', error_message=str(e), failure_stage='callback') - else: - error_message = result.get('error_message', 'unknown') - failure_stage = result.get('failure_stage', 'edge') - logger.error(f"[task_id={task_id}] Edge 分析失败: stage={failure_stage}, error={error_message}") - db_layer.update_task_status( - task_id, 'FAILED', error_message=error_message, failure_stage=failure_stage) + logger.error(f"[task_id={task_id}] Edge 返回异常状态码: {resp.status_code}") + self._schedule_retry(task) def _schedule_retry(self, task): """调度重试""" @@ -152,13 +140,11 @@ class Dispatcher: def _poll_once(self): """执行一次轮询""" - # 先回收僵尸任务:PROCESSING 超过 push_timeout+缓冲 说明推送进程已丢失 - # (如 fam-core 重启、Edge 重启掐断 in-flight 连接),重置回 PENDING 走重试 - stale_timeout = self.push_timeout + 120 + # 回收僵尸任务:PROCESSING 超过 stale_timeout 说明 Edge 丢失了任务 try: - stale_ids = db_layer.reclaim_stale_processing(stale_timeout) + stale_ids = db_layer.reclaim_stale_processing(self.stale_timeout) for tid in stale_ids: - logger.warning(f"[task_id={tid}] PROCESSING 超时 {stale_timeout}s,回收为 PENDING 重试") + logger.warning(f"[task_id={tid}] PROCESSING 超时 {self.stale_timeout}s,回收为 PENDING 重试") except Exception as e: logger.error(f"僵尸任务回收失败: {e}", exc_info=True) @@ -172,7 +158,7 @@ class Dispatcher: def _run(self): """线程主循环""" - logger.info(f"Dispatcher 启动,轮询间隔 {self.poll_interval}s") + logger.info(f"Dispatcher 启动 (enqueue 模式),轮询间隔 {self.poll_interval}s") while self._running: try: self._poll_once() diff --git a/fam-core/src/fam_core/poller/__init__.py b/fam-core/src/fam_core/poller/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/fam-core/src/fam_core/poller/poller.py b/fam-core/src/fam_core/poller/poller.py new file mode 100644 index 0000000..818567a --- /dev/null +++ b/fam-core/src/fam_core/poller/poller.py @@ -0,0 +1,128 @@ +""" +Poller - 定期从 Edge 拉取已完成的任务结果,写入 MariaDB + +流程: +1. 每 N 秒请求 Edge /api/edge/results?limit=10 +2. 遍历结果列表,对每个 nas_task_id: + - success: 调用 apply_success_event 写入 monitor_events + event_details,标记 SUCCESS + - failed: 更新任务状态为 FAILED,记录 failure_stage 和 error_message +3. Edge 端自动标记已拉取的结果为 delivered +""" +import time +import threading +import requests + +from ..logger import setup_logger, log_task +from ..config_loader import load_config +from .. import db_layer +from ..event_receiver.event_receiver import apply_success_event + +logger = setup_logger('fam-core.poller') + + +class Poller: + """结果拉取器,定期从 Edge 拉取处理结果""" + + def __init__(self): + cfg = load_config() + poller_cfg = cfg.get('poller', {}) + self.poll_interval = poller_cfg.get('poll_interval', 30) + self.results_url = poller_cfg.get('results_url', 'http://localhost:5000/api/edge/results') + self.batch_size = poller_cfg.get('batch_size', 10) + self.timeout = poller_cfg.get('timeout', 30) + self._running = False + self._thread = None + + def _handle_result(self, item: dict): + """处理单个结果""" + nas_task_id = item.get('nas_task_id') + result = item.get('result') + + if not nas_task_id: + logger.warning(f"结果缺少 nas_task_id,跳过: {item}") + return + + if result is None: + logger.error(f"[task_id={nas_task_id}] Edge 返回空结果,标记 FAILED") + db_layer.update_task_status( + nas_task_id, 'FAILED', + error_message='Edge returned empty result', + failure_stage='callback' + ) + return + + status = result.get('status') + if status == 'success': + try: + event_id = apply_success_event(nas_task_id, result) + log_task(logger, nas_task_id, 'poller', f'结果落库成功: event_id={event_id}') + except Exception as e: + logger.error(f"[task_id={nas_task_id}] 结果落库失败: {e}", exc_info=True) + db_layer.update_task_status( + nas_task_id, 'FAILED', error_message=str(e), failure_stage='callback') + elif status == 'failed': + error_message = result.get('error_message', 'unknown') + failure_stage = result.get('failure_stage', '') + logger.error(f"[task_id={nas_task_id}] Edge 处理失败: stage={failure_stage}, error={error_message}") + db_layer.update_task_status( + nas_task_id, 'FAILED', error_message=error_message, failure_stage=failure_stage) + else: + logger.warning(f"[task_id={nas_task_id}] 未知状态: {status}") + + def _poll_once(self): + """执行一次拉取""" + try: + resp = requests.get( + self.results_url, + params={'limit': self.batch_size}, + timeout=self.timeout + ) + except requests.RequestException as e: + logger.error(f"拉取结果失败: {e}") + return + + if resp.status_code != 200: + logger.warning(f"Edge 返回 {resp.status_code}") + return + + try: + data = resp.json() + except ValueError: + logger.error("Edge 返回非 JSON 响应") + return + + results = data.get('results', []) + if not results: + return + + logger.info(f"拉取到 {len(results)} 条结果") + for item in results: + try: + self._handle_result(item) + except Exception as e: + task_id = item.get('nas_task_id', '?') + logger.error(f"[task_id={task_id}] 处理结果异常: {e}", exc_info=True) + + def _run(self): + """线程主循环""" + logger.info(f"Poller 启动,轮询间隔 {self.poll_interval}s,目标: {self.results_url}") + while self._running: + try: + self._poll_once() + 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='poller') + self._thread.start() + + def stop(self): + """停止拉取线程""" + self._running = False + if self._thread: + self._thread.join(timeout=5) diff --git a/fam-edge/config/config.yaml b/fam-edge/config/config.yaml index 79f0e68..a1e6eaa 100644 --- a/fam-edge/config/config.yaml +++ b/fam-edge/config/config.yaml @@ -1,7 +1,12 @@ # FAM-Edge 配置文件 (Oracle 端) - 多模型池配置 # Tailscale: Oracle=100.74.137.126, NAS=100.70.234.39 +# +# 异步队列模式: +# NAS 上传视频 → /api/edge/video/enqueue 入 SQLite 队列 → 消费者线程异步处理 +# → NAS Poller 从 /api/edge/results 拉取结果 +# 速率限制: Gemini 1000RPM x2 burst, NVIDIA 40RPM x2 burst -# NAS 端回调地址 +# NAS 端回调地址(旧 webhook 模式保留,异步模式不使用) nas: webhook_url: "http://100.70.234.39:8000/api/core/callback/event" media_base_url: "http://100.70.234.39:8000/media" @@ -13,6 +18,17 @@ server: port: 5000 max_concurrent_tasks: 1 +# 异步任务队列 +queue: + db_path: "/opt/fam-edge/data/fam_queue.db" + upload_dir: "/tmp/fam_uploads" + poll_interval: 10 # 消费者轮询间隔(秒) + # API 速率限制 (RPM),burst_factor=2 表示突发容量为 2 倍 RPM + rate_limit: + gemini_rpm: 1000 + nvidia_rpm: 40 + burst_factor: 2 + # 编排调度模式: fallback(顺序降级, 默认) | ensemble(并行交叉验证) orchestrator: mode: "fallback" diff --git a/fam-edge/src/fam_edge/ai_orchestrator/orchestrator.py b/fam-edge/src/fam_edge/ai_orchestrator/orchestrator.py index 66f4b5d..a7ce75f 100644 --- a/fam-edge/src/fam_edge/ai_orchestrator/orchestrator.py +++ b/fam-edge/src/fam_edge/ai_orchestrator/orchestrator.py @@ -52,12 +52,15 @@ class AIOrchestrator: def run_visual_analysis(self, adapters: List[BaseModelAdapter], frame_paths: List[str], frame_timestamps: List[str], - known_members_context: str) -> Dict[str, dict]: + known_members_context: str, + rate_limiter=None) -> Dict[str, dict]: """视觉分析阶段:仅 role=vision 的适配器参与 orchestrator.mode: - fallback (默认): 按 config 顺序依次尝试,首个成功即采用(单元素 dict) - ensemble: 并行所有健康 vision 模型,全部成功结果都保留(交叉验证) + + rate_limiter: 可选 RateLimiter 实例,按 provider 限速(2x burst) """ vision_adapters = [a for a in adapters if getattr(a, 'role', 'vision') == 'vision'] if not vision_adapters: @@ -68,7 +71,8 @@ class AIOrchestrator: if mode == 'ensemble': return self._run_visual_ensemble( - vision_adapters, frame_paths, frame_timestamps, known_members_context) + vision_adapters, frame_paths, frame_timestamps, + known_members_context, rate_limiter) # fallback: 顺序降级,首个成功即采用 model_outputs = {} @@ -76,6 +80,12 @@ class AIOrchestrator: if adapter.get_circuit_breaker().is_open(): logger.warning(f"[{adapter.provider_name}] 熔断器 OPEN,跳过") continue + # 速率限制:按 provider 获取 token(2x burst) + if rate_limiter: + acquired = rate_limiter.acquire(adapter.provider_name, timeout=300) + if not acquired: + logger.warning(f"[{adapter.provider_name}] 速率限制超时,跳过") + continue start = time.time() try: output = adapter.analyze_frames( @@ -97,7 +107,8 @@ class AIOrchestrator: return model_outputs def _run_visual_ensemble(self, vision_adapters, frame_paths, - frame_timestamps, known_members_context) -> Dict[str, dict]: + frame_timestamps, known_members_context, + rate_limiter=None) -> Dict[str, dict]: """并行调用所有健康 vision 模型,保留全部成功结果(交叉验证)""" model_outputs = {} max_timeout = max((a.get_timeout() for a in vision_adapters), default=240) @@ -107,6 +118,12 @@ class AIOrchestrator: if adapter.get_circuit_breaker().is_open(): logger.warning(f"[{adapter.provider_name}] 熔断器 OPEN,跳过") continue + # 速率限制:按 provider 获取 token(2x burst) + if rate_limiter: + acquired = rate_limiter.acquire(adapter.provider_name, timeout=300) + if not acquired: + logger.warning(f"[{adapter.provider_name}] 速率限制超时,跳过") + continue future = pool.submit( adapter.analyze_frames, frame_paths, frame_timestamps, known_members_context) @@ -377,9 +394,12 @@ class AIOrchestrator: return 200 def process_push_task(self, task_data: dict, video_path: str, - preprocessor: 'VideoPreprocessor') -> dict: + preprocessor: 'VideoPreprocessor', + rate_limiter=None) -> dict: """推送模式:同步处理上传的视频,结果直接返回(无 webhook 回调) + rate_limiter: 可选 RateLimiter 实例,按 provider 限速(2x burst) + 返回 payload 结构与原 webhook 回调一致: - 成功: {task_id, status: "success", event_start_time, ..., frame_details, ...} - 失败: {task_id, status: "failed", failure_stage, error_message} @@ -429,7 +449,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, rate_limiter ) if not model_outputs: raise Exception('All models failed in visual analysis') diff --git a/fam-edge/src/fam_edge/api_gateway/api_gateway.py b/fam-edge/src/fam_edge/api_gateway/api_gateway.py index ce6828c..779e959 100644 --- a/fam-edge/src/fam_edge/api_gateway/api_gateway.py +++ b/fam-edge/src/fam_edge/api_gateway/api_gateway.py @@ -1,8 +1,12 @@ """ API-Gateway - Flask 蓝图,接收任务 -同时只允许 1 个任务在处理;新任务到达时若当前有任务处理中,返回 429 +模式: +1. enqueue (异步): NAS 上传视频 → Edge 入队 → 立即返回 → 消费者异步处理 → NAS 轮询拉取结果 +2. push (同步, 兼容保留): NAS 上传 → Edge 同步处理 → 结果随响应返回 +3. analyze (旧拉取模式, 兼容保留) """ +import os import threading import requests from flask import Blueprint, request, jsonify @@ -10,12 +14,12 @@ from flask import Blueprint, request, jsonify from ..logger import setup_logger from ..ai_orchestrator.orchestrator import AIOrchestrator from ..video_preprocessor.preprocessor import VideoPreprocessor +from ..queue import queue_manager logger = setup_logger('fam-edge.api_gateway') api_bp = Blueprint('api_gateway', __name__) -# 并发控制:同时只允许 1 个任务 _current_task_lock = threading.Lock() _currently_processing = False @@ -72,6 +76,92 @@ def receive_task(): return jsonify({"status": "accepted", "task_id": task_id}), 202 +@api_bp.route('/api/edge/video/enqueue', methods=['POST']) +def enqueue_task(): + """异步模式:接收 multipart 视频上传,入队后立即返回 + + NAS 上传视频 → Edge 保存到磁盘 + 入 SQLite 队列 → 返回 task_id + 消费者线程异步处理,NAS 通过 /api/edge/results 拉取结果 + """ + task_id_raw = request.form.get('task_id') + file = request.files.get('video') + if not task_id_raw or not file: + return jsonify({"error": "缺少必填字段: task_id, video"}), 400 + + try: + task_id = int(task_id_raw) + except ValueError: + return jsonify({"error": "task_id 必须是整数"}), 400 + + camera_name = request.form.get('camera_name', '') + event_start_time = request.form.get('event_start_time', '') + known_members_context = request.form.get('known_members_context', '') + + upload_dir = os.environ.get('FAM_UPLOAD_DIR', '/tmp/fam_uploads') + os.makedirs(upload_dir, exist_ok=True) + video_filename = f"task_{task_id}_{file.filename}" + video_path = os.path.join(upload_dir, video_filename) + + try: + file.save(video_path) + size_mb = os.path.getsize(video_path) / 1024 / 1024 + logger.info(f"[task_id={task_id}] 入队: {file.filename} ({size_mb:.1f}MB)") + + queue_id = queue_manager.enqueue( + nas_task_id=task_id, + video_filename=file.filename, + video_path=video_path, + camera_name=camera_name, + event_start_time=event_start_time, + known_members_context=known_members_context, + ) + + return jsonify({ + "status": "queued", + "task_id": task_id, + "queue_id": queue_id, + }), 202 + + except Exception as e: + logger.error(f"[task_id={task_id}] 入队失败: {e}", exc_info=True) + if os.path.exists(video_path): + os.remove(video_path) + return jsonify({"error": str(e)}), 500 + + +@api_bp.route('/api/edge/results', methods=['GET']) +def get_results(): + """返回已完成但未拉取的结果,标记为已交付""" + limit = int(request.args.get('limit', 10)) + results = queue_manager.get_undelivered_results(limit=limit) + + import json + payload = [] + task_ids = [] + for r in results: + try: + result_json = json.loads(r['result_json']) if r['result_json'] else None + except json.JSONDecodeError: + result_json = None + payload.append({ + "nas_task_id": r['nas_task_id'], + "result": result_json, + }) + task_ids.append(r['id']) + + if task_ids: + queue_manager.mark_delivered(task_ids) + + return jsonify({"results": payload, "count": len(payload)}), 200 + + +@api_bp.route('/api/edge/queue/stats', methods=['GET']) +def queue_stats(): + """队列状态统计""" + stats = queue_manager.get_queue_stats() + return jsonify(stats), 200 + + @api_bp.route('/api/edge/video/push', methods=['POST']) def receive_push_task(): """推送模式:接收 multipart 视频上传,同步分析,结果随 HTTP 响应返回 diff --git a/fam-edge/src/fam_edge/app.py b/fam-edge/src/fam_edge/app.py index 031c131..5e62ccd 100644 --- a/fam-edge/src/fam_edge/app.py +++ b/fam-edge/src/fam_edge/app.py @@ -1,7 +1,7 @@ """ FAM-Edge 主应用 - Flask 单进程 -承载: API-Gateway / Video-Preprocessor / AI-Orchestrator / Storage-Cleaner +承载: API-Gateway / Video-Preprocessor / AI-Orchestrator / Storage-Cleaner / Queue-Consumer """ import os import sys @@ -12,6 +12,7 @@ sys.path.insert(0, os.path.dirname(os.path.abspath(__file__))) from .config_loader import load_config from .logger import setup_logger from .api_gateway.api_gateway import api_bp +from .queue.consumer import get_consumer logger = setup_logger('fam-edge.app') @@ -21,7 +22,17 @@ app.register_blueprint(api_bp) @app.route('/', methods=['GET']) def root(): - return jsonify({"service": "fam-edge", "version": "1.0"}), 200 + return jsonify({"service": "fam-edge", "version": "2.0"}), 200 + + +# 启动消费者线程(异步任务队列) +_consumer = None +try: + _consumer = get_consumer() + _consumer.start() + logger.info("Queue-Consumer 已启动") +except Exception as e: + logger.error(f"Queue-Consumer 启动失败: {e}") if __name__ == '__main__': diff --git a/fam-edge/src/fam_edge/queue/__init__.py b/fam-edge/src/fam_edge/queue/__init__.py new file mode 100644 index 0000000..9f104b8 --- /dev/null +++ b/fam-edge/src/fam_edge/queue/__init__.py @@ -0,0 +1,5 @@ +""" +SQLite 异步任务队列 (Oracle 端) +""" +from . import queue_manager +from .consumer import get_consumer diff --git a/fam-edge/src/fam_edge/queue/consumer.py b/fam-edge/src/fam_edge/queue/consumer.py new file mode 100644 index 0000000..8e19bc9 --- /dev/null +++ b/fam-edge/src/fam_edge/queue/consumer.py @@ -0,0 +1,128 @@ +""" +消费者线程 - 从 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 diff --git a/fam-edge/src/fam_edge/queue/queue_manager.py b/fam-edge/src/fam_edge/queue/queue_manager.py new file mode 100644 index 0000000..07e0beb --- /dev/null +++ b/fam-edge/src/fam_edge/queue/queue_manager.py @@ -0,0 +1,174 @@ +""" +SQLite 队列管理器 - Oracle 端异步任务队列 + +表结构: +- task_queue: 任务队列 (PENDING → PROCESSING → SUCCESS/FAILED) +- 元数据: delivered 标记 NAS 是否已拉取结果 +""" +import os +import sqlite3 +import json +import threading +from typing import Optional, List, Dict + +DB_PATH = os.environ.get('FAM_QUEUE_DB', '/opt/fam-edge/data/fam_queue.db') + +_init_lock = threading.Lock() +_initialized = False + + +def _get_conn() -> sqlite3.Connection: + global _initialized + if not _initialized: + with _init_lock: + if not _initialized: + _init_db() + _initialized = True + conn = sqlite3.connect(DB_PATH, timeout=30) + conn.row_factory = sqlite3.Row + conn.execute("PRAGMA journal_mode=WAL") + return conn + + +def _init_db(): + os.makedirs(os.path.dirname(DB_PATH), exist_ok=True) + conn = sqlite3.connect(DB_PATH) + conn.execute(""" + CREATE TABLE IF NOT EXISTS task_queue ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + nas_task_id INTEGER NOT NULL, + video_filename TEXT NOT NULL, + video_path TEXT NOT NULL, + camera_name TEXT DEFAULT '', + event_start_time TEXT DEFAULT '', + known_members_context TEXT DEFAULT '', + status TEXT DEFAULT 'PENDING', + result_json TEXT, + error_message TEXT, + failure_stage TEXT, + retry_count INTEGER DEFAULT 0, + created_at TEXT DEFAULT (datetime('now', 'localtime')), + updated_at TEXT DEFAULT (datetime('now', 'localtime')), + delivered INTEGER DEFAULT 0, + UNIQUE(nas_task_id) + ) + """) + conn.execute("CREATE INDEX IF NOT EXISTS idx_status ON task_queue(status)") + conn.execute("CREATE INDEX IF NOT EXISTS idx_delivered ON task_queue(delivered)") + conn.commit() + conn.close() + + +def enqueue(nas_task_id: int, video_filename: str, video_path: str, + camera_name: str, event_start_time: str, + known_members_context: str) -> int: + conn = _get_conn() + try: + cur = conn.execute( + "INSERT OR IGNORE INTO task_queue " + "(nas_task_id, video_filename, video_path, camera_name, event_start_time, known_members_context) " + "VALUES (?, ?, ?, ?, ?, ?)", + (nas_task_id, video_filename, video_path, camera_name, event_start_time, known_members_context) + ) + conn.commit() + if cur.rowcount == 0: + row = conn.execute( + "SELECT id FROM task_queue WHERE nas_task_id=?", (nas_task_id,) + ).fetchone() + return row['id'] if row else 0 + return cur.lastrowid + finally: + conn.close() + + +def claim_next() -> Optional[Dict]: + conn = _get_conn() + try: + conn.execute("BEGIN IMMEDIATE") + row = conn.execute( + "SELECT * FROM task_queue WHERE status='PENDING' ORDER BY id LIMIT 1" + ).fetchone() + if row: + conn.execute( + "UPDATE task_queue SET status='PROCESSING', updated_at=datetime('now','localtime') WHERE id=?", + (row['id'],) + ) + conn.commit() + return dict(row) + conn.rollback() + return None + except Exception: + conn.rollback() + return None + finally: + conn.close() + + +def mark_success(task_id: int, result_json: str): + conn = _get_conn() + try: + conn.execute( + "UPDATE task_queue SET status='SUCCESS', result_json=?, updated_at=datetime('now','localtime') WHERE id=?", + (result_json, task_id) + ) + conn.commit() + finally: + conn.close() + + +def mark_failed(task_id: int, error_message: str, failure_stage: str = ''): + conn = _get_conn() + try: + conn.execute( + "UPDATE task_queue SET status='FAILED', error_message=?, failure_stage=?, " + "updated_at=datetime('now','localtime') WHERE id=?", + (error_message, failure_stage, task_id) + ) + conn.commit() + finally: + conn.close() + + +def get_undelivered_results(limit: int = 10) -> List[Dict]: + conn = _get_conn() + try: + rows = conn.execute( + "SELECT * FROM task_queue WHERE status='SUCCESS' AND delivered=0 " + "ORDER BY id LIMIT ?", (limit,) + ).fetchall() + return [dict(r) for r in rows] + finally: + conn.close() + + +def mark_delivered(task_ids: List[int]): + if not task_ids: + return + conn = _get_conn() + try: + placeholders = ','.join('?' * len(task_ids)) + conn.execute( + f"UPDATE task_queue SET delivered=1, updated_at=datetime('now','localtime') " + f"WHERE id IN ({placeholders})", task_ids + ) + conn.commit() + finally: + conn.close() + + +def get_queue_stats() -> Dict: + conn = _get_conn() + try: + stats = {} + for status in ['PENDING', 'PROCESSING', 'SUCCESS', 'FAILED']: + row = conn.execute( + "SELECT COUNT(*) as cnt FROM task_queue WHERE status=?", (status,) + ).fetchone() + stats[status] = row['cnt'] + row = conn.execute( + "SELECT COUNT(*) as cnt FROM task_queue WHERE status='SUCCESS' AND delivered=0" + ).fetchone() + stats['UNDELIVERED'] = row['cnt'] + return stats + finally: + conn.close() diff --git a/fam-edge/src/fam_edge/rate_limiter.py b/fam-edge/src/fam_edge/rate_limiter.py new file mode 100644 index 0000000..24a9ec1 --- /dev/null +++ b/fam-edge/src/fam_edge/rate_limiter.py @@ -0,0 +1,48 @@ +""" +Rate Limiter - Token Bucket 算法 + +按 API 限制速度的 2 倍设置突发容量,按 API 限制速度持续补充。 +""" +import time +import threading + + +class TokenBucket: + def __init__(self, rpm: int, burst_factor: int = 2): + self.capacity = rpm * burst_factor + self.refill_rate = rpm / 60.0 + self.tokens = float(self.capacity) + self.last_refill = time.monotonic() + self._lock = threading.Lock() + + def acquire(self, tokens: int = 1, timeout: float = 300.0) -> bool: + deadline = time.monotonic() + timeout + while True: + with self._lock: + now = time.monotonic() + elapsed = now - self.last_refill + self.tokens = min(self.capacity, self.tokens + elapsed * self.refill_rate) + self.last_refill = now + if self.tokens >= tokens: + self.tokens -= tokens + return True + wait = (tokens - self.tokens) / self.refill_rate + if time.monotonic() + wait > deadline: + return False + time.sleep(min(wait, 1.0)) + + +class RateLimiter: + """多 API 速率限制管理""" + + def __init__(self): + self._buckets = {} + + def register(self, name: str, rpm: int, burst_factor: int = 2): + self._buckets[name] = TokenBucket(rpm, burst_factor) + + def acquire(self, name: str, tokens: int = 1, timeout: float = 300.0) -> bool: + bucket = self._buckets.get(name) + if bucket is None: + return True + return bucket.acquire(tokens, timeout)