From 02380301280136e0bc4a7d53c98f3f82ceb7ad4e Mon Sep 17 00:00:00 2001 From: ericwyuan Date: Thu, 3 Sep 2026 09:00:50 +0800 Subject: [PATCH] =?UTF-8?q?fix(sync):=20=E4=BF=AE=E5=A4=8D=20Oracle=20id?= =?UTF-8?q?=20=E9=87=8D=E6=8E=92=E5=AF=BC=E8=87=B4=E7=9A=84=20sync=5Fvideo?= =?UTF-8?q?s/sync=5Fpeople=201062?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Oracle 库重建/重排后 videos/people 的 id 跳变(2303 段→3408/4600+ 段), 旧 upsert 按 Oracle id 作主键插入时与已存在同 filename/label 行在 UNIQUE 键 上二次冲突报 1062,整批中止、游标不推进。 - upsert_sync_videos: 改为以 filename 业务键去重,命中就地 UPDATE 保留 NAS 原 id(防 sync_events/model_calls/identity_map 的 video_id 外键失效),Oracle id 落 oracle_id 列溯源;未命中插入优先用 Oracle id 对齐子表引用,主键冲突回退自增 - upsert_sync_people: 同模式,以 label 业务键去重(people.id 无外键引用) - ddl.sql: sync_videos/sync_people 的 id 改为 NAS 本地自增主键 + 新增 oracle_id 列 - NAS 现网已迁移(加列/回填/改自增)并重启验证:单批补拉 videos+348 events+1017 people+9 model_calls+500 identity_map+238,游标 09-02 13:55:21 → 09-03 08:51:11 --- fam-core/src/fam_core/db_layer.py | 155 ++++++++++++++++++++---------- scripts/ddl.sql | 12 ++- 2 files changed, 115 insertions(+), 52 deletions(-) diff --git a/fam-core/src/fam_core/db_layer.py b/fam-core/src/fam_core/db_layer.py index 45825d7..e4168b7 100644 --- a/fam-core/src/fam_core/db_layer.py +++ b/fam-core/src/fam_core/db_layer.py @@ -56,7 +56,16 @@ def get_conn(): # ============================================================ def upsert_sync_videos(rows: List[Dict]) -> int: - """批量 upsert Oracle 传来的 videos 增量。rows 为 Oracle 端 dict 列表。""" + """批量 upsert Oracle 传来的 videos 增量。rows 为 Oracle 端 dict 列表。 + + v2 修复 (2026-09-03):不再以 Oracle 的 videos.id 作为 NAS 主键。Oracle 库重建/ + 重排后 id 会复用/跳变(实测 2303 段跳到 3408/4600+ 段),旧写法按 id 插入时与 + 已存在的同 filename 行在 UNIQUE filename 上二次冲突报 1062,整批中止、游标不推进。 + 现改为:以 filename 为业务唯一键去重,命中则就地 UPDATE(保留 NAS 原 id,防止 + sync_events/model_calls/identity_map 的 video_id 外键失效),Oracle id 仅落 + oracle_id 列溯源;未命中则插入(优先用 Oracle id 作主键以对齐子表引用,主键冲突时 + 回退本地自增,避免 1062)。 + """ if not rows: return 0 conn = get_conn() @@ -64,35 +73,57 @@ def upsert_sync_videos(rows: List[Dict]) -> int: try: cur = conn.cursor() for r in rows: - cur.execute( - """INSERT INTO sync_videos - (id, drive_file_id, filename, camera_name, duration_sec, - event_start_time, status, summary_json, events_json, - people_json, compute_provider, created_at, updated_at, - processed_at, synced_at) - VALUES (%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s, NOW()) - ON DUPLICATE KEY UPDATE - drive_file_id=VALUES(drive_file_id), - filename=VALUES(filename), - camera_name=VALUES(camera_name), - duration_sec=VALUES(duration_sec), - event_start_time=VALUES(event_start_time), - status=VALUES(status), - summary_json=VALUES(summary_json), - events_json=VALUES(events_json), - people_json=VALUES(people_json), - compute_provider=VALUES(compute_provider), - created_at=VALUES(created_at), - updated_at=VALUES(updated_at), - processed_at=VALUES(processed_at), - synced_at=NOW()""", - (r.get('id'), r.get('drive_file_id'), r.get('filename'), - r.get('camera_name'), r.get('duration_sec') or 0, - r.get('event_start_time'), r.get('status'), - r.get('summary_json'), r.get('events_json'), r.get('people_json'), - r.get('compute_provider'), r.get('created_at'), - r.get('updated_at'), r.get('processed_at')) - ) + fn = r.get('filename') + oracle_id = r.get('id') + cur.execute("SELECT id FROM sync_videos WHERE filename=%s", (fn,)) + row = cur.fetchone() + if row: + # 命中业务键:就地更新,保留原 NAS id(外键不失效) + cur.execute( + """UPDATE sync_videos SET + oracle_id=%s, drive_file_id=%s, camera_name=%s, + duration_sec=%s, event_start_time=%s, status=%s, + summary_json=%s, events_json=%s, people_json=%s, + compute_provider=%s, created_at=%s, updated_at=%s, + processed_at=%s, synced_at=NOW() + WHERE id=%s""", + (oracle_id, r.get('drive_file_id'), r.get('camera_name'), + r.get('duration_sec') or 0, r.get('event_start_time'), + r.get('status'), r.get('summary_json'), r.get('events_json'), + r.get('people_json'), r.get('compute_provider'), + r.get('created_at'), r.get('updated_at'), + r.get('processed_at'), row[0])) + else: + # 未命中:新文件。优先用 Oracle id 作主键(对齐子表 video_id 引用) + try: + cur.execute( + """INSERT INTO sync_videos + (id, oracle_id, filename, camera_name, duration_sec, + event_start_time, status, summary_json, events_json, + people_json, compute_provider, created_at, updated_at, + processed_at, synced_at) + VALUES (%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s, NOW())""", + (oracle_id, oracle_id, fn, r.get('camera_name'), + r.get('duration_sec') or 0, r.get('event_start_time'), + r.get('status'), r.get('summary_json'), r.get('events_json'), + r.get('people_json'), r.get('compute_provider'), + r.get('created_at'), r.get('updated_at'), + r.get('processed_at'))) + except Exception: + # 主键 id 被复用(极端情况):回退本地自增,避免 1062 + cur.execute( + """INSERT INTO sync_videos + (oracle_id, filename, camera_name, duration_sec, + event_start_time, status, summary_json, events_json, + people_json, compute_provider, created_at, updated_at, + processed_at, synced_at) + VALUES (%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s, NOW())""", + (oracle_id, fn, r.get('camera_name'), + r.get('duration_sec') or 0, r.get('event_start_time'), + r.get('status'), r.get('summary_json'), r.get('events_json'), + r.get('people_json'), r.get('compute_provider'), + r.get('created_at'), r.get('updated_at'), + r.get('processed_at'))) n += 1 conn.commit() return n @@ -304,7 +335,12 @@ def query_sync_events_for_person_date(person: str, date_str: str) -> List[Dict]: # ============================================================ def upsert_sync_people(rows: List[Dict]) -> int: - """批量 upsert Oracle 传来的 people 增量。""" + """批量 upsert Oracle 传来的 people 增量。 + + v2 修复 (2026-09-03):同 videos,以 label 为业务唯一键去重,命中就地 UPDATE + 保留 NAS 原 id(people.id 无外键引用,仅作展示排序),Oracle id 落 oracle_id 列 + 溯源;未命中插入(优先 Oracle id,主键冲突回退自增,避免 1062)。 + """ if not rows: return 0 conn = get_conn() @@ -312,25 +348,44 @@ def upsert_sync_people(rows: List[Dict]) -> int: try: cur = conn.cursor() for r in rows: - cur.execute( - """INSERT INTO sync_people - (id, label, canonical_name, first_seen, appearances, - source, features_json, display_uid, updated_at, synced_at) - VALUES (%s,%s,%s,%s,%s,%s,%s,%s,%s, NOW()) - ON DUPLICATE KEY UPDATE - label=VALUES(label), - canonical_name=VALUES(canonical_name), - first_seen=VALUES(first_seen), - appearances=VALUES(appearances), - source=VALUES(source), - features_json=VALUES(features_json), - display_uid=VALUES(display_uid), - updated_at=VALUES(updated_at), - synced_at=NOW()""", - (r.get('id'), r.get('label'), r.get('canonical_name'), - r.get('first_seen'), r.get('appearances') or 0, - r.get('source'), r.get('features_json'), - r.get('display_uid'), r.get('updated_at'))) + label = r.get('label') + oracle_id = r.get('id') + cur.execute("SELECT id FROM sync_people WHERE label=%s", (label,)) + row = cur.fetchone() + if row: + cur.execute( + """UPDATE sync_people SET + oracle_id=%s, canonical_name=%s, first_seen=%s, + appearances=%s, source=%s, features_json=%s, + display_uid=%s, updated_at=%s, synced_at=NOW() + WHERE id=%s""", + (oracle_id, r.get('canonical_name'), r.get('first_seen'), + r.get('appearances') or 0, r.get('source'), + r.get('features_json'), r.get('display_uid'), + r.get('updated_at'), row[0])) + else: + try: + cur.execute( + """INSERT INTO sync_people + (id, oracle_id, label, canonical_name, first_seen, + appearances, source, features_json, display_uid, + updated_at, synced_at) + VALUES (%s,%s,%s,%s,%s,%s,%s,%s,%s,%s, NOW())""", + (oracle_id, oracle_id, label, r.get('canonical_name'), + r.get('first_seen'), r.get('appearances') or 0, + r.get('source'), r.get('features_json'), + r.get('display_uid'), r.get('updated_at'))) + except Exception: + cur.execute( + """INSERT INTO sync_people + (oracle_id, label, canonical_name, first_seen, + appearances, source, features_json, display_uid, + updated_at, synced_at) + VALUES (%s,%s,%s,%s,%s,%s,%s,%s,%s, NOW())""", + (oracle_id, label, r.get('canonical_name'), + r.get('first_seen'), r.get('appearances') or 0, + r.get('source'), r.get('features_json'), + r.get('display_uid'), r.get('updated_at'))) n += 1 conn.commit() return n diff --git a/scripts/ddl.sql b/scripts/ddl.sql index 422643f..ea036c2 100644 --- a/scripts/ddl.sql +++ b/scripts/ddl.sql @@ -128,8 +128,13 @@ CREATE TABLE IF NOT EXISTS family_members ( -- ============================================================ -- 7.1 视频会话表(Oracle videos 镜像) +-- v2 (2026-09-03):id 改为 NAS 本地自增主键(不再直接等于 Oracle id,因 Oracle 库 +-- 重建/重排会复用 id,导致 filename UNIQUE 二次冲突 1062);Oracle 的 videos.id 仅 +-- 落 oracle_id 列溯源。业务唯一键仍为 filename。子表 video_id 引用仍用 Oracle id +-- (新视频插入时 id=oracle_id;旧视频原地更新保留原 id,外键不失效)。 CREATE TABLE IF NOT EXISTS sync_videos ( - id INT PRIMARY KEY COMMENT 'Oracle videos.id', + id INT AUTO_INCREMENT PRIMARY KEY COMMENT 'NAS 本地自增主键', + oracle_id INT COMMENT 'Oracle videos.id(仅溯源参考)', drive_file_id VARCHAR(255) COMMENT 'Google 硬盘文件 ID', filename VARCHAR(500) NOT NULL UNIQUE COMMENT '视频文件名(唯一)', camera_name VARCHAR(50) COMMENT '摄像头名称/位置', @@ -164,8 +169,11 @@ CREATE TABLE IF NOT EXISTS sync_events ( ) ENGINE=InnoDB COMMENT='甲骨文事件镜像表'; -- 7.3 人物规范表(Oracle people 镜像) +-- v2 (2026-09-03):同 7.1,id 改为 NAS 本地自增主键,Oracle people.id 落 oracle_id +-- 列溯源;业务唯一键仍为 label。people.id 无外键引用,仅作展示排序。 CREATE TABLE IF NOT EXISTS sync_people ( - id INT PRIMARY KEY COMMENT 'Oracle people.id', + id INT AUTO_INCREMENT PRIMARY KEY COMMENT 'NAS 本地自增主键', + oracle_id INT COMMENT 'Oracle people.id(仅溯源参考)', label VARCHAR(100) NOT NULL UNIQUE COMMENT '抽象标识,如"人物A"', canonical_name VARCHAR(100) COMMENT '规范名(用户命名或 LLM 合并),NULL 表示未命名', first_seen VARCHAR(32) COMMENT '首次出现时间',