feat(fam-edge): 队列健壮性 - ①失败重试间隔 retry_interval_sec=1h(last_fail_at 记录,配额类故障等恢复再试);②max_retries 2→10;③入队前 OpenCV 文件校验(可解码/大小/元数据,失败标 invalid 不入队)+ mtime 稳定窗口防 rclone 半成品 + 处理前二次确认
This commit is contained in:
@@ -41,7 +41,10 @@ video_processing:
|
|||||||
max_concurrent: 1 # 消费者线程数(串行处理,避免云端并发超额)
|
max_concurrent: 1 # 消费者线程数(串行处理,避免云端并发超额)
|
||||||
timeout: 900 # 兜底单视频分析超时
|
timeout: 900 # 兜底单视频分析超时
|
||||||
timeout_multiplier: 2 # 模型消费超时倍数:在 models[i].timeout 原值上 ×2(大视频上传+分析耗时)
|
timeout_multiplier: 2 # 模型消费超时倍数:在 models[i].timeout 原值上 ×2(大视频上传+分析耗时)
|
||||||
max_retries: 2 # 单视频失败最大重试次数(failed 且 retry_count 超限不再消费)
|
max_retries: 10 # 单视频失败最大重试次数(配额/过载等瞬时故障给足重试机会)
|
||||||
|
retry_interval_sec: 3600 # 失败重试最小间隔:距上次失败 ≥1h 才重新入队,等配额恢复
|
||||||
|
file_validate: true # 登记入队前用 OpenCV 校验文件可解码;失败标记 invalid 不入队
|
||||||
|
stable_window_sec: 60 # 文件 mtime 稳定窗口:写入中(rclone 同步未完成)的文件跳过本轮
|
||||||
# 降级顺序:先 gemini 整视频,失败再 nvidia 整视频;两者都失败 -> 标记 failed
|
# 降级顺序:先 gemini 整视频,失败再 nvidia 整视频;两者都失败 -> 标记 failed
|
||||||
vision_order: ["gemini", "nvidia"]
|
vision_order: ["gemini", "nvidia"]
|
||||||
|
|
||||||
|
|||||||
@@ -51,6 +51,10 @@ class OracleDB:
|
|||||||
event_start_time TEXT,
|
event_start_time TEXT,
|
||||||
status TEXT DEFAULT 'pending',
|
status TEXT DEFAULT 'pending',
|
||||||
retry_count INTEGER DEFAULT 0,
|
retry_count INTEGER DEFAULT 0,
|
||||||
|
file_valid INTEGER DEFAULT 1,
|
||||||
|
file_error TEXT,
|
||||||
|
media_meta_json TEXT,
|
||||||
|
last_fail_at TEXT,
|
||||||
summary_json TEXT,
|
summary_json TEXT,
|
||||||
events_json TEXT,
|
events_json TEXT,
|
||||||
people_json TEXT,
|
people_json TEXT,
|
||||||
@@ -97,10 +101,17 @@ class OracleDB:
|
|||||||
CREATE INDEX IF NOT EXISTS idx_events_video ON events(video_id);
|
CREATE INDEX IF NOT EXISTS idx_events_video ON events(video_id);
|
||||||
CREATE INDEX IF NOT EXISTS idx_model_calls_created ON model_calls(created_at);
|
CREATE INDEX IF NOT EXISTS idx_model_calls_created ON model_calls(created_at);
|
||||||
""")
|
""")
|
||||||
# 兼容旧库:补 retry_count 列(生产-消费队列重试上限用)
|
# 兼容旧库:补 retry_count / file_valid / media 等列(生产-消费队列用)
|
||||||
cols = [r[1] for r in c.execute("PRAGMA table_info(videos)").fetchall()]
|
cols = [r[1] for r in c.execute("PRAGMA table_info(videos)").fetchall()]
|
||||||
if 'retry_count' not in cols:
|
for col, ddl in [
|
||||||
c.execute("ALTER TABLE videos ADD COLUMN retry_count INTEGER DEFAULT 0")
|
('retry_count', "ALTER TABLE videos ADD COLUMN retry_count INTEGER DEFAULT 0"),
|
||||||
|
('file_valid', "ALTER TABLE videos ADD COLUMN file_valid INTEGER DEFAULT 1"),
|
||||||
|
('file_error', "ALTER TABLE videos ADD COLUMN file_error TEXT"),
|
||||||
|
('media_meta_json', "ALTER TABLE videos ADD COLUMN media_meta_json TEXT"),
|
||||||
|
('last_fail_at', "ALTER TABLE videos ADD COLUMN last_fail_at TEXT"),
|
||||||
|
]:
|
||||||
|
if col not in cols:
|
||||||
|
c.execute(ddl)
|
||||||
self._conn.commit()
|
self._conn.commit()
|
||||||
|
|
||||||
# ------------------------------------------------------------------
|
# ------------------------------------------------------------------
|
||||||
@@ -128,6 +139,18 @@ class OracleDB:
|
|||||||
duration_sec, 1 if success else 0, error or '', now))
|
duration_sec, 1 if success else 0, error or '', now))
|
||||||
self._conn.commit()
|
self._conn.commit()
|
||||||
|
|
||||||
|
def set_video_file_status(self, video_id: int, valid: bool,
|
||||||
|
error: str = '', media_meta: dict = None):
|
||||||
|
"""登记/更新文件校验结果:valid / file_error / media_meta_json"""
|
||||||
|
now = _now_iso()
|
||||||
|
self._conn.execute(
|
||||||
|
"UPDATE videos SET file_valid=?, file_error=?, media_meta_json=?, "
|
||||||
|
"updated_at=? WHERE id=?",
|
||||||
|
(1 if valid else 0, error or '',
|
||||||
|
json.dumps(media_meta, ensure_ascii=False) if media_meta else None,
|
||||||
|
now, video_id))
|
||||||
|
self._conn.commit()
|
||||||
|
|
||||||
def ensure_video(self, filename: str, local_path: str,
|
def ensure_video(self, filename: str, local_path: str,
|
||||||
camera_name: str = '', event_start_time: str = '',
|
camera_name: str = '', event_start_time: str = '',
|
||||||
duration_sec: float = 0.0, drive_file_id: str = '') -> int:
|
duration_sec: float = 0.0, drive_file_id: str = '') -> int:
|
||||||
@@ -183,8 +206,8 @@ class OracleDB:
|
|||||||
now = _now_iso()
|
now = _now_iso()
|
||||||
self._conn.execute(
|
self._conn.execute(
|
||||||
"UPDATE videos SET status='failed', retry_count=retry_count+1, "
|
"UPDATE videos SET status='failed', retry_count=retry_count+1, "
|
||||||
"summary_json=?, updated_at=? WHERE id=?",
|
"summary_json=?, updated_at=?, last_fail_at=? WHERE id=?",
|
||||||
(error, now, video_id))
|
(error, now, now, video_id))
|
||||||
self._conn.commit()
|
self._conn.commit()
|
||||||
|
|
||||||
def get_all_videos(self) -> List[sqlite3.Row]:
|
def get_all_videos(self) -> List[sqlite3.Row]:
|
||||||
|
|||||||
@@ -86,6 +86,46 @@ def _parse_event_ts(ts: str, start_dt):
|
|||||||
return ts, 0.0
|
return ts, 0.0
|
||||||
|
|
||||||
|
|
||||||
|
def validate_video(path: str) -> tuple:
|
||||||
|
"""校验视频文件是否为正常可解码视频。
|
||||||
|
|
||||||
|
返回 (ok: bool, error: str, meta: dict|None)
|
||||||
|
- meta: {fps, frames, duration_sec, width, height}
|
||||||
|
用 OpenCV 打开并读取至少 1 帧(不校验会导致空/半成品文件浪费云端配额)。
|
||||||
|
"""
|
||||||
|
meta = None
|
||||||
|
try:
|
||||||
|
if not path or not os.path.isfile(path):
|
||||||
|
return False, "file_missing", None
|
||||||
|
if os.path.getsize(path) == 0:
|
||||||
|
return False, "file_empty", None
|
||||||
|
try:
|
||||||
|
import cv2
|
||||||
|
except ImportError:
|
||||||
|
return True, "", None # 无 cv2 时跳过深度校验(仅大小检查)
|
||||||
|
cap = cv2.VideoCapture(path)
|
||||||
|
try:
|
||||||
|
if not cap.isOpened():
|
||||||
|
return False, "cannot_open", None
|
||||||
|
ok, frame = cap.read()
|
||||||
|
if not ok or frame is None:
|
||||||
|
return False, "no_decodable_frame", None
|
||||||
|
fps = float(cap.get(cv2.CAP_PROP_FPS) or 0)
|
||||||
|
frames = int(cap.get(cv2.CAP_PROP_FRAME_COUNT) or 0)
|
||||||
|
meta = {
|
||||||
|
"fps": round(fps, 2),
|
||||||
|
"frames": frames,
|
||||||
|
"duration_sec": round(frames / max(fps, 0.01), 1),
|
||||||
|
"width": int(cap.get(cv2.CAP_PROP_FRAME_WIDTH) or 0),
|
||||||
|
"height": int(cap.get(cv2.CAP_PROP_FRAME_HEIGHT) or 0),
|
||||||
|
}
|
||||||
|
finally:
|
||||||
|
cap.release()
|
||||||
|
return True, "", meta
|
||||||
|
except Exception as e:
|
||||||
|
return False, f"validate_exc: {e}", None
|
||||||
|
|
||||||
|
|
||||||
class VideoProcessor:
|
class VideoProcessor:
|
||||||
def __init__(self, db: oracle_db.OracleDB):
|
def __init__(self, db: oracle_db.OracleDB):
|
||||||
self.config = load_config()
|
self.config = load_config()
|
||||||
@@ -93,6 +133,8 @@ class VideoProcessor:
|
|||||||
self.vision_order = self.config.get('video_processing', {}).get(
|
self.vision_order = self.config.get('video_processing', {}).get(
|
||||||
'vision_order', ['gemini', 'nvidia'])
|
'vision_order', ['gemini', 'nvidia'])
|
||||||
self.vision_timeout = self.config.get('video_processing', {}).get('timeout', 900)
|
self.vision_timeout = self.config.get('video_processing', {}).get('timeout', 900)
|
||||||
|
self.file_validate = bool(self.config.get('video_processing', {}).get(
|
||||||
|
'file_validate', True))
|
||||||
self.parse_start = self.config.get('gdrive_sync', {}).get(
|
self.parse_start = self.config.get('gdrive_sync', {}).get(
|
||||||
'parse_start_from_filename', True)
|
'parse_start_from_filename', True)
|
||||||
adapters = build_adapters(self.config.get('models', []))
|
adapters = build_adapters(self.config.get('models', []))
|
||||||
@@ -122,6 +164,17 @@ class VideoProcessor:
|
|||||||
self.db.mark_video_failed(video_id, "file_missing")
|
self.db.mark_video_failed(video_id, "file_missing")
|
||||||
return False
|
return False
|
||||||
|
|
||||||
|
# 处理前二次确认文件有效性(防止登记后文件被破坏/截断;校验结果落库)
|
||||||
|
if self.file_validate:
|
||||||
|
ok, verr, vmeta = validate_video(local_path)
|
||||||
|
if not ok:
|
||||||
|
logger.error(f"[video_id={video_id}] 文件校验失败({verr}),标记 failed: {local_path}")
|
||||||
|
self.db.set_video_file_status(video_id, False, verr)
|
||||||
|
self.db.mark_video_failed(video_id, f"invalid_file:{verr}")
|
||||||
|
return False
|
||||||
|
if vmeta:
|
||||||
|
self.db.set_video_file_status(video_id, True, '', vmeta)
|
||||||
|
|
||||||
camera_name = self.db.get_video_by_filename(filename)['camera_name'] or ''
|
camera_name = self.db.get_video_by_filename(filename)['camera_name'] or ''
|
||||||
event_start = ''
|
event_start = ''
|
||||||
if self.parse_start:
|
if self.parse_start:
|
||||||
|
|||||||
@@ -21,7 +21,7 @@ from typing import List, Optional
|
|||||||
from .logger import setup_logger
|
from .logger import setup_logger
|
||||||
from .config_loader import load_config
|
from .config_loader import load_config
|
||||||
from . import oracle_db
|
from . import oracle_db
|
||||||
from .video_processor import VideoProcessor
|
from .video_processor import VideoProcessor, validate_video
|
||||||
|
|
||||||
logger = setup_logger('fam-edge.video_queue')
|
logger = setup_logger('fam-edge.video_queue')
|
||||||
|
|
||||||
@@ -44,7 +44,12 @@ class VideoQueue:
|
|||||||
# 模型超时放大倍数:消费时在模型原配置 timeout 上 ×N
|
# 模型超时放大倍数:消费时在模型原配置 timeout 上 ×N
|
||||||
self.timeout_multiplier = float(vp.get('timeout_multiplier', 2.0))
|
self.timeout_multiplier = float(vp.get('timeout_multiplier', 2.0))
|
||||||
# 单视频失败最大重试次数(failed 且 retry_count>=max_retries 不再消费)
|
# 单视频失败最大重试次数(failed 且 retry_count>=max_retries 不再消费)
|
||||||
self.max_retries = int(vp.get('max_retries', 2))
|
self.max_retries = int(vp.get('max_retries', 10))
|
||||||
|
# 失败重试最小间隔(秒):配额类瞬时故障(如 429)需等其恢复再重试,避免短时间重复打爆
|
||||||
|
self.retry_interval_sec = int(vp.get('retry_interval_sec', 3600))
|
||||||
|
# 文件校验:入队前用 OpenCV 确认可解码;mtime 稳定窗口防 rclone 写入中的半成品
|
||||||
|
self.file_validate = bool(vp.get('file_validate', True))
|
||||||
|
self.stable_window_sec = int(vp.get('stable_window_sec', 60))
|
||||||
|
|
||||||
self._queue: "queue.Queue[int]" = queue.Queue()
|
self._queue: "queue.Queue[int]" = queue.Queue()
|
||||||
self._queued: set = set() # 已在队中的 video_id(防重复入队)
|
self._queued: set = set() # 已在队中的 video_id(防重复入队)
|
||||||
@@ -75,6 +80,26 @@ class VideoQueue:
|
|||||||
fn = os.path.basename(path)
|
fn = os.path.basename(path)
|
||||||
row = self.db.get_video_by_filename(fn)
|
row = self.db.get_video_by_filename(fn)
|
||||||
if row is None:
|
if row is None:
|
||||||
|
# 新文件:先过 mtime 稳定窗口 + 可解码校验,通过才登记入队;失败标记 invalid
|
||||||
|
if self.file_validate:
|
||||||
|
if time.time() - os.path.getmtime(path) < self.stable_window_sec:
|
||||||
|
logger.info(f"文件仍在写入(mtime 未稳定),跳过本轮: {fn}")
|
||||||
|
continue
|
||||||
|
ok, verr, vmeta = validate_video(path)
|
||||||
|
if not ok:
|
||||||
|
vid = self.db.ensure_video(fn, path, camera_name=self.camera_name)
|
||||||
|
self.db.set_video_file_status(vid, False, verr)
|
||||||
|
self.db._conn.execute(
|
||||||
|
"UPDATE videos SET status='invalid' WHERE id=?", (vid,))
|
||||||
|
self.db._conn.commit()
|
||||||
|
logger.warning(f"文件校验失败,标记 invalid 不入队: {fn} ({verr})")
|
||||||
|
continue
|
||||||
|
if vmeta:
|
||||||
|
vid = self.db.ensure_video(fn, path, camera_name=self.camera_name)
|
||||||
|
self.db.set_video_file_status(vid, True, '', vmeta)
|
||||||
|
logger.info(f"登记新视频并入队: {fn} (id={vid}, meta={vmeta})")
|
||||||
|
self._enqueue(vid)
|
||||||
|
continue
|
||||||
vid = self.db.ensure_video(fn, path, camera_name=self.camera_name)
|
vid = self.db.ensure_video(fn, path, camera_name=self.camera_name)
|
||||||
logger.info(f"登记新视频并入队: {fn} (id={vid})")
|
logger.info(f"登记新视频并入队: {fn} (id={vid})")
|
||||||
self._enqueue(vid)
|
self._enqueue(vid)
|
||||||
@@ -91,10 +116,26 @@ class VideoQueue:
|
|||||||
logger.info(f"启动恢复入队 {len(rows)} 个待处理视频")
|
logger.info(f"启动恢复入队 {len(rows)} 个待处理视频")
|
||||||
|
|
||||||
def _retry_allowed(self, row) -> bool:
|
def _retry_allowed(self, row) -> bool:
|
||||||
if row['status'] == 'pending':
|
"""是否允许消费/重试该视频:pending 直接放行;failed 需未超重试上限且距上次失败够久;invalid 永不放行"""
|
||||||
|
status = row['status']
|
||||||
|
if status == 'pending':
|
||||||
return True
|
return True
|
||||||
# failed:未超重试上限才允许再次消费
|
if status == 'invalid':
|
||||||
return int(row['retry_count'] or 0) < self.max_retries
|
return False
|
||||||
|
# failed:未超重试上限,且距上次失败 >= retry_interval_sec(等配额类故障恢复)
|
||||||
|
if int(row['retry_count'] or 0) >= self.max_retries:
|
||||||
|
return False
|
||||||
|
last_fail = row['last_fail_at'] or row['updated_at'] or ''
|
||||||
|
if last_fail:
|
||||||
|
try:
|
||||||
|
from datetime import datetime
|
||||||
|
lt = datetime.strptime(last_fail[:19], '%Y-%m-%d %H:%M:%S')
|
||||||
|
elapsed = (datetime.now() - lt).total_seconds()
|
||||||
|
if elapsed < self.retry_interval_sec:
|
||||||
|
return False
|
||||||
|
except ValueError:
|
||||||
|
pass
|
||||||
|
return True
|
||||||
|
|
||||||
def _enqueue(self, video_id: int):
|
def _enqueue(self, video_id: int):
|
||||||
with self._lock:
|
with self._lock:
|
||||||
|
|||||||
Reference in New Issue
Block a user