feat: 异步任务队列架构 - SQLite队列 + 速率限制 + NAS Poller

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'
This commit is contained in:
ericwyuan
2026-08-20 12:07:09 +08:00
parent a4b9178a59
commit 02da23ef42
14 changed files with 686 additions and 62 deletions

View File

@@ -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"

View File

@@ -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 获取 token2x 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 获取 token2x 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')

View File

@@ -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 响应返回

View File

@@ -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__':

View File

@@ -0,0 +1,5 @@
"""
SQLite 异步任务队列 (Oracle 端)
"""
from . import queue_manager
from .consumer import get_consumer

View File

@@ -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

View File

@@ -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()

View File

@@ -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)