feat(fam-edge): 视频生产-消费队列 - 新增 video_queue.py(生产者扫描新文件入队+重启恢复,独立消费者线程云端分析);模型消费超时按原配置×timeout_multiplier(2);failed 重试上限 max_retries(2);删除旧 watch_processor

This commit is contained in:
ericwyuan
2026-08-21 12:01:25 +08:00
parent ab5b706eeb
commit ce27e4ab0c
8 changed files with 242 additions and 108 deletions

View File

@@ -0,0 +1,199 @@
"""
video_queue - 视频生产-消费队列(新架构 v2.1
职责拆分:
生产者(Producer): 轮询 rclone 同步落地目录
- 新文件 -> 登记 videos 表(status=pending) -> 放入内存队列
- 启动时把 DB 中未处理完的(pending / failed 且未超重试上限)补入队
消费者(Consumer): 独立工作线程(max_concurrent 个)
- 从队列取 video_id -> 云端模型整视频分析(模型超时 = 原配置 × timeout_multiplier
- 成功 -> done全部模型失败 -> failedretry_count+1超上限不再重试
与旧 WatchProcessor 的区别: 生产/消费解耦,消费有独立线程与重试上限,
模型调用超时按各模型原配置 ×2应对大视频上传+分析耗时)。
"""
import os
import time
import queue
import threading
from typing import List, Optional
from .logger import setup_logger
from .config_loader import load_config
from . import oracle_db
from .video_processor import VideoProcessor
logger = setup_logger('fam-edge.video_queue')
VIDEO_EXTS = ('.mp4', '.mkv', '.avi', '.mov', '.ts')
class VideoQueue:
def __init__(self, db: oracle_db.OracleDB):
self.config = load_config()
self.db = db
self.processor = VideoProcessor(db)
self.local_dir = self.config.get('gdrive_sync', {}).get(
'local_dir', '/opt/fam-edge/gdrive_videos')
self.watch_interval = self.config.get('gdrive_sync', {}).get(
'watch_interval_sec', 30)
self.camera_name = self.config.get('gdrive_sync', {}).get(
'camera_name', '摄像头')
vp = self.config.get('video_processing', {})
self.max_concurrent = int(vp.get('max_concurrent', 1))
# 模型超时放大倍数:消费时在模型原配置 timeout 上 ×N
self.timeout_multiplier = float(vp.get('timeout_multiplier', 2.0))
# 单视频失败最大重试次数failed 且 retry_count>=max_retries 不再消费)
self.max_retries = int(vp.get('max_retries', 2))
self._queue: "queue.Queue[int]" = queue.Queue()
self._queued: set = set() # 已在队中的 video_id防重复入队
self._lock = threading.Lock()
self._producer_stop = threading.Event()
self._consumer_stop = threading.Event()
self._producer = None
self._consumers: List[threading.Thread] = []
self._stats = {"produced": 0, "consumed_ok": 0, "consumed_fail": 0}
# ------------------------------------------------------------------
# 生产者
# ------------------------------------------------------------------
def _scan_files(self) -> List[str]:
"""递归扫描监听目录,返回所有视频文件完整路径"""
if not os.path.isdir(self.local_dir):
logger.warning(f"监听目录不存在: {self.local_dir}")
return []
out = []
for root, _dirs, files in os.walk(self.local_dir):
for fn in sorted(files):
if fn.lower().endswith(VIDEO_EXTS):
out.append(os.path.join(root, fn))
return out
def _produce_once(self):
for path in self._scan_files():
fn = os.path.basename(path)
row = self.db.get_video_by_filename(fn)
if row is None:
vid = self.db.ensure_video(fn, path, camera_name=self.camera_name)
logger.info(f"登记新视频并入队: {fn} (id={vid})")
self._enqueue(vid)
elif row['status'] in ('pending', 'failed') and self._retry_allowed(row):
self._enqueue(row['id'])
def _enqueue_existing(self):
"""启动时把 DB 中未处理完的视频补入队(进程重启恢复)"""
rows = self.db.get_pending_videos(limit=10000)
for row in rows:
if self._retry_allowed(row):
self._enqueue(row['id'])
if rows:
logger.info(f"启动恢复入队 {len(rows)} 个待处理视频")
def _retry_allowed(self, row) -> bool:
if row['status'] == 'pending':
return True
# failed未超重试上限才允许再次消费
return int(row['retry_count'] or 0) < self.max_retries
def _enqueue(self, video_id: int):
with self._lock:
if video_id in self._queued:
return
self._queued.add(video_id)
self._queue.put(video_id)
self._stats["produced"] += 1
def _producer_loop(self):
logger.info(f"生产者启动,监听 {self.local_dir},间隔 {self.watch_interval}s")
while not self._producer_stop.is_set():
try:
self._produce_once()
except Exception as e:
logger.error(f"生产者扫描异常: {e}", exc_info=True)
for _ in range(self.watch_interval):
if self._producer_stop.is_set():
break
time.sleep(1)
# ------------------------------------------------------------------
# 消费者
# ------------------------------------------------------------------
def _consume_loop(self):
while not self._consumer_stop.is_set():
try:
video_id = self._queue.get(timeout=1)
except queue.Empty:
continue
try:
self._consume_one(video_id)
except Exception as e:
logger.error(f"消费 video_id={video_id} 异常: {e}", exc_info=True)
finally:
with self._lock:
self._queued.discard(video_id)
self._queue.task_done()
def _consume_one(self, video_id: int):
row = self.db.get_video_by_id(video_id)
if row is None:
logger.warning(f"消费到不存在的 video_id={video_id},跳过")
return
if row['status'] == 'done':
return
if not self._retry_allowed(row):
logger.warning(f"[video_id={video_id}] 已达重试上限({row['retry_count']}),放弃")
return
ok = self.processor.process_video(
video_id, row['filename'], row['local_path'],
timeout_multiplier=self.timeout_multiplier)
if ok:
self._stats["consumed_ok"] += 1
else:
self._stats["consumed_fail"] += 1
# ------------------------------------------------------------------
# 生命周期
# ------------------------------------------------------------------
def start(self):
if self._producer is not None and self._producer.is_alive():
return
# 启动恢复:先补入队 DB 中未处理完的
try:
self._enqueue_existing()
except Exception as e:
logger.error(f"启动恢复入队失败: {e}", exc_info=True)
self._producer = threading.Thread(
target=self._producer_loop, daemon=True, name='video-producer')
self._producer.start()
for i in range(max(1, self.max_concurrent)):
t = threading.Thread(
target=self._consume_loop, daemon=True,
name=f'video-consumer-{i}')
t.start()
self._consumers.append(t)
logger.info(f"VideoQueue 已启动: {max(1, self.max_concurrent)} 个消费者, "
f"超时倍数 ×{self.timeout_multiplier}, 重试上限 {self.max_retries}")
def stop(self):
self._producer_stop.set()
self._consumer_stop.set()
if self._producer:
self._producer.join(timeout=5)
for t in self._consumers:
t.join(timeout=5)
def stats(self) -> dict:
return {
"queue_size": self._queue.qsize(),
"consumers": len(self._consumers),
"timeout_multiplier": self.timeout_multiplier,
"max_retries": self.max_retries,
"produced": self._stats["produced"],
"consumed_ok": self._stats["consumed_ok"],
"consumed_fail": self._stats["consumed_fail"],
}
def is_alive(self) -> bool:
return (self._producer is not None and self._producer.is_alive()
and any(t.is_alive() for t in self._consumers))