From 4e85a989448d5d1b737052233dcd07dfb01edd85 Mon Sep 17 00:00:00 2001 From: ericwyuan Date: Fri, 21 Aug 2026 14:24:12 +0800 Subject: [PATCH] =?UTF-8?q?fix(=E9=94=81=E4=B8=8E=E6=9F=A5=E8=AF=A2?= =?UTF-8?q?=E7=BB=84):=20#12=20=E7=86=94=E6=96=AD=E5=99=A8=E7=8A=B6?= =?UTF-8?q?=E6=80=81=E8=BD=AC=E6=8D=A2=E5=8A=A0=E9=94=81=EF=BC=9B#10=20?= =?UTF-8?q?=E9=97=AE=E7=AD=94=E4=BA=BA=E5=90=8D=E5=8C=B9=E9=85=8D=E6=94=B9?= =?UTF-8?q?=20JSON=5FCONTAINS=20=E7=B2=BE=E7=A1=AE=E5=8C=B9=E9=85=8D(?= =?UTF-8?q?=E9=98=B2=20LIKE=20=E5=AD=90=E4=B8=B2/=E7=89=B9=E6=AE=8A?= =?UTF-8?q?=E5=AD=97=E7=AC=A6=E8=AF=AF=E5=8C=B9=E9=85=8D)=EF=BC=9B#4=20Ora?= =?UTF-8?q?cleDB=20=E5=A4=8D=E5=90=88=E5=86=99=E5=8A=A0=E9=94=81=EF=BC=9B#?= =?UTF-8?q?6=20=E6=AF=8F=E6=B6=88=E8=B4=B9=E8=80=85=E7=8B=AC=E7=AB=8B=20Vi?= =?UTF-8?q?deoProcessor(=E6=B6=88=E9=99=A4=20adapter.timeout=20=E5=85=B1?= =?UTF-8?q?=E4=BA=AB=E7=AB=9E=E4=BA=89)=EF=BC=9B#15=20done=20=E8=A7=86?= =?UTF-8?q?=E9=A2=91=E6=96=87=E4=BB=B6=E8=A2=AB=E8=A6=86=E7=9B=96(mtime>pr?= =?UTF-8?q?ocessed=5Fat)=E8=87=AA=E5=8A=A8=E9=87=8D=E7=BD=AE=E9=87=8D?= =?UTF-8?q?=E5=88=86=E6=9E=90?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- fam-core/src/fam_core/db_layer.py | 6 ++- .../model_adapters/circuit_breaker.py | 23 +++++++---- fam-edge/src/fam_edge/oracle_db.py | 41 ++++++++++--------- fam-edge/src/fam_edge/video_queue.py | 32 +++++++++++++-- 4 files changed, 68 insertions(+), 34 deletions(-) diff --git a/fam-core/src/fam_core/db_layer.py b/fam-core/src/fam_core/db_layer.py index fe90442..c73b7b0 100644 --- a/fam-core/src/fam_core/db_layer.py +++ b/fam-core/src/fam_core/db_layer.py @@ -212,9 +212,11 @@ def query_sync_events_for_person_date(person: str, date_str: str) -> List[Dict]: v.camera_name, v.filename, v.event_start_time, v.processed_at FROM sync_events e JOIN sync_videos v ON e.video_id = v.id - WHERE v.processed_at LIKE %s AND e.person_list_json LIKE %s + WHERE v.processed_at LIKE %s + AND e.person_list_json IS NOT NULL + AND JSON_CONTAINS(e.person_list_json, JSON_QUOTE(%s), '$') ORDER BY v.processed_at ASC, e.ts ASC""", - (f'{date_str}%', f'%{person}%')) + (f'{date_str}%', person)) return cur.fetchall() finally: conn.close() diff --git a/fam-edge/src/fam_edge/model_adapters/circuit_breaker.py b/fam-edge/src/fam_edge/model_adapters/circuit_breaker.py index 85fee74..b494bce 100644 --- a/fam-edge/src/fam_edge/model_adapters/circuit_breaker.py +++ b/fam-edge/src/fam_edge/model_adapters/circuit_breaker.py @@ -8,6 +8,7 @@ """ from collections import deque import time +import threading class CircuitBreaker: @@ -21,27 +22,31 @@ class CircuitBreaker: self.cooldown = cooldown self.state = 'CLOSED' self.last_failure = None + self._lock = threading.Lock() # 状态转换加锁,防多线程竞争 def record_failure(self): if not self.enabled: return - self.failures.append(time.time()) - if len(self.failures) >= self.threshold: - self.state = 'OPEN' - self.last_failure = time.time() + with self._lock: + self.failures.append(time.time()) + if len(self.failures) >= self.threshold: + self.state = 'OPEN' + self.last_failure = time.time() def record_success(self): if not self.enabled: return - self.failures.clear() - self.state = 'CLOSED' + with self._lock: + self.failures.clear() + self.state = 'CLOSED' def is_open(self): if not self.enabled: return False - if self.state == 'OPEN' and self.last_failure and time.time() - self.last_failure > self.cooldown: - self.state = 'HALF_OPEN' - return self.state == 'OPEN' + with self._lock: + if self.state == 'OPEN' and self.last_failure and time.time() - self.last_failure > self.cooldown: + self.state = 'HALF_OPEN' + return self.state == 'OPEN' def __repr__(self): return f"CircuitBreaker(state={self.state}, enabled={self.enabled})" diff --git a/fam-edge/src/fam_edge/oracle_db.py b/fam-edge/src/fam_edge/oracle_db.py index e99c75e..0de774a 100644 --- a/fam-edge/src/fam_edge/oracle_db.py +++ b/fam-edge/src/fam_edge/oracle_db.py @@ -17,6 +17,7 @@ Oracle 本地库(SQLite) - 视频摘要 / 事件 / 人物 存储 import os import json import sqlite3 +import threading from datetime import datetime, timezone, timedelta from typing import Dict, List, Optional @@ -35,6 +36,7 @@ class OracleDB: self._conn.row_factory = sqlite3.Row self._conn.execute("PRAGMA journal_mode=WAL") self._conn.execute("PRAGMA busy_timeout=10000") + self._write_lock = threading.Lock() # 复合写(如 DELETE+INSERT+commit)串行化 self._init_schema() # ------------------------------------------------------------------ @@ -182,25 +184,26 @@ class OracleDB: def mark_video_processed(self, video_id: int, summary: str, events: List[dict], people: List[str], compute_provider: str) -> List[int]: """落库视频结果;返回新插入事件的 id 列表(与 events 参数一一对应,供事件截图用)""" - now = _now_iso() - self._conn.execute( - "UPDATE videos SET status='done', summary_json=?, events_json=?, " - "people_json=?, compute_provider=?, updated_at=?, processed_at=? WHERE id=?", - (summary, json.dumps(events, ensure_ascii=False), json.dumps(people, ensure_ascii=False), - compute_provider, now, now, video_id)) - # 事件落独立表,便于 NAS 拉取 - self._conn.execute("DELETE FROM events WHERE video_id=?", (video_id,)) - event_ids: List[int] = [] - for ev in events: - cur = self._conn.execute( - "INSERT INTO events (video_id, ts, description, person_list_json, " - "is_attention_event) VALUES (?,?,?,?,?)", - (video_id, ev.get('timestamp', ''), ev.get('description', ''), - json.dumps(ev.get('people', []), ensure_ascii=False), - 1 if ev.get('is_attention_event') else 0)) - event_ids.append(cur.lastrowid) - self._conn.commit() - return event_ids + with self._write_lock: + now = _now_iso() + self._conn.execute( + "UPDATE videos SET status='done', summary_json=?, events_json=?, " + "people_json=?, compute_provider=?, updated_at=?, processed_at=? WHERE id=?", + (summary, json.dumps(events, ensure_ascii=False), json.dumps(people, ensure_ascii=False), + compute_provider, now, now, video_id)) + # 事件落独立表,便于 NAS 拉取 + self._conn.execute("DELETE FROM events WHERE video_id=?", (video_id,)) + event_ids: List[int] = [] + for ev in events: + cur = self._conn.execute( + "INSERT INTO events (video_id, ts, description, person_list_json, " + "is_attention_event) VALUES (?,?,?,?,?)", + (video_id, ev.get('timestamp', ''), ev.get('description', ''), + json.dumps(ev.get('people', []), ensure_ascii=False), + 1 if ev.get('is_attention_event') else 0)) + event_ids.append(cur.lastrowid) + self._conn.commit() + return event_ids def mark_video_failed(self, video_id: int, error: str = ''): now = _now_iso() diff --git a/fam-edge/src/fam_edge/video_queue.py b/fam-edge/src/fam_edge/video_queue.py index 88458ee..3fa432f 100644 --- a/fam-edge/src/fam_edge/video_queue.py +++ b/fam-edge/src/fam_edge/video_queue.py @@ -32,7 +32,6 @@ 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( @@ -105,6 +104,28 @@ class VideoQueue: self._enqueue(vid) elif row['status'] in ('pending', 'failed') and self._retry_allowed(row): self._enqueue(row['id']) + elif row['status'] == 'done' and self._file_changed(row, path): + # 文件被覆盖(rclone 重新同步/更新):重置 pending 重新分析 + logger.info(f"文件内容变更,重置重新分析: {fn} (id={row['id']})") + self.db._conn.execute( + "UPDATE videos SET status='pending', retry_count=0, summary_json=NULL, " + "events_json=NULL, people_json=NULL, compute_provider=NULL, " + "processed_at=NULL, file_valid=1 WHERE id=?", + (row['id'],)) + self.db._conn.commit() + self._enqueue(row['id']) + + def _file_changed(self, row, path: str) -> bool: + """判断视频文件在处理后被覆盖(mtime 晚于 processed_at)。""" + try: + proc = row['processed_at'] or '' + if not proc: + return False + from datetime import datetime + pt = datetime.strptime(proc[:19], '%Y-%m-%d %H:%M:%S').timestamp() + return os.path.getmtime(path) > pt + 5 # 5s 容差 + except Exception: + return False def _enqueue_existing(self): """启动时把 DB 中未处理完的视频补入队(进程重启恢复)""" @@ -161,13 +182,16 @@ class VideoQueue: # 消费者 # ------------------------------------------------------------------ def _consume_loop(self): + # 每消费者独立 VideoProcessor(各自 build_adapters), + # 避免多消费者共享 adapter 实例导致 timeout 等实例属性改写竞争 + processor = VideoProcessor(self.db) 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) + self._consume_one(video_id, processor) except Exception as e: logger.error(f"消费 video_id={video_id} 异常: {e}", exc_info=True) finally: @@ -175,7 +199,7 @@ class VideoQueue: self._queued.discard(video_id) self._queue.task_done() - def _consume_one(self, video_id: int): + def _consume_one(self, video_id: int, processor: VideoProcessor): row = self.db.get_video_by_id(video_id) if row is None: logger.warning(f"消费到不存在的 video_id={video_id},跳过") @@ -185,7 +209,7 @@ class VideoQueue: if not self._retry_allowed(row): logger.warning(f"[video_id={video_id}] 已达重试上限({row['retry_count']}),放弃") return - ok = self.processor.process_video( + ok = processor.process_video( video_id, row['filename'], row['local_path'], timeout_multiplier=self.timeout_multiplier) if ok: